Flink 作业连续运行三个月,Checkpoint 从未失败,Doris 导入任务也全部成功。月底对账却发现部分订单被重复累计,另一些订单停在旧状态。团队最容易得出的结论是 Doris Connector 没有做到 Exactly-Once。
这个判断通常太早。Checkpoint、Doris 事务和业务结果分别解决三个不同问题。
Exactly-Once 只能证明一批数据不会因为失败恢复被重复提交,不能证明事件顺序正确、表模型正确,更不能替代业务对账。
一条链路里其实有三种正确
Flink 状态正确:失败后从同一 Checkpoint 恢复 ↓ Doris 提交正确:同一批数据只成功提交一次 ↓ 业务状态正确:乱序、删除、主键和字段合并符合业务规则第一层由 Flink Checkpoint 管,第二层依赖 Connector、Stream Load 2PC 和 Label,第三层依赖 Key 模型、Sequence、删除语义和对账。上层成功不能自动推出下层成功。
四种方案看起来都能重试,语义完全不同
下面做同条件对比:Kafka 至少一次投递,Sink 写入后在提交前崩溃,任务恢复并重放同一批数据。
| 写入方案 | 重试时发生什么 | 能保证什么 | 仍然解决不了什么 |
|---|---|---|---|
| 普通 Stream Load,每次随机 Label | 重放生成新事务 | 每次请求原子提交 | 重复数据 |
| 固定业务 Label | 相同 Label 被识别 | 单批幂等 | 跨批边界和 Label 过期 |
| Flink Connector 2PC | Checkpoint 与 Doris PRECOMMIT/COMMIT 协调 | 故障恢复后的提交唯一性 | 乱序、错误主键、漏同步 Delete |
| Batch Mode + Unique Key | 可能重复写,主键覆盖 | 最终一主键一行 | Duplicate/Aggregate 表重复,旧事件覆盖新状态 |
Doris 事务官方文档 明确区分两层:Label 保证单个事务不重复,2PC 用于跨系统协调。Label 还会按时间和数量清理,默认保留边界不能被当成永久幂等账本。
2PC 真正绑定的是 Checkpoint 与 Doris 事务
Flink Doris Connector 官方文档 中,流式写入默认依赖 Checkpoint,sink.enable-2pc默认开启。核心链路可以压成五步:
数据写入 Doris 临时事务 → Flink 触发 Checkpoint → Sink 预提交当前事务 → Checkpoint 全局完成 → Sink 提交 Doris 事务并开启下一事务故障发生在不同位置,恢复动作不同:
- PRECOMMIT 前失败,本批状态随 Checkpoint 一起回滚并重放;
- PRECOMMIT 后、Checkpoint 完成前失败,未完成事务不能冒充成功批次;
- Checkpoint 已完成但客户端未收到 COMMIT 响应,恢复后通过事务标识判断,而不是盲目再写一批。
sink.label-prefix必须全局唯一。两个作业共用前缀,不是增强幂等,而是在争用事务身份。
最小故障实验比看 SUCCESS 更有价值
建立 Unique Key 订单状态表,Sequence 使用单调业务版本op_version。准备三个事件:创建、支付、迟到取消。Connector 开启 2PC,Checkpoint 间隔设置为测试可观察的 30 秒。
CREATETABLEorder_state_eos(order_idBIGINT,op_versionBIGINT,statusVARCHAR(16),amountDECIMAL(12,2))UNIQUEKEY(order_id)DISTRIBUTEDBYHASH(order_id)BUCKETS4PROPERTIES("enable_unique_key_merge_on_write"="true","function_column.sequence_col"="op_version");测试不要只杀一次进程,而要覆盖三个切点:写入中、预提交后、Checkpoint 完成附近。每次恢复后检查四项证据:
SELECTorder_id,COUNT(*)FROMorder_state_eosGROUPBYorder_idHAVINGCOUNT(*)<>1;SELECTorder_id,MAX(op_version),MAX_BY(status,op_version)FROMorder_event_auditGROUPBYorder_id;第一条检查服务表的一主键一行,第二条从审计事件重建期望状态。只查 Doris 最终行无法发现源端事件是否漏掉,所以生产上最好保留 Duplicate Key 审计表作为旁路证据。
三个经典假成功
Batch Mode 把 EOS 悄悄关掉
Connector 开启sink.enable.batch-mode后,提交由时间和数据量触发,不再跟随 Checkpoint。官方文档明确说明此模式不保证 Exactly-Once。Unique Key 可以消除同主键重复,却不能保护 Duplicate Key 明细和 Aggregate Key 累加值。
旧事件被精确提交一次
迟到的取消事件只提交一次,完全满足传输 EOS;没有 Sequence 的 Unique Key 仍会让它覆盖支付状态。一次且仅一次地写错,仍然是错。
删除事件根本没有进入 Sink
Connector、上游 CDC 或表模型未正确处理 Delete,Checkpoint 仍会成功。状态表保留幽灵记录,技术指标全部绿色。
源码该追的是状态边界
发布文章前应将 Connector 版本与 Doris4.0.8同时固定。源码阅读只追四个状态,不需要展开整个 Connector:
begin transaction → write data → preCommit on checkpoint barrier → commit or abort after checkpoint resultDoris 端重点核对 FE 事务管理中的 Label 唯一性、PRECOMMITTED 到 COMMITTED/VISIBLE 的状态迁移,以及重复 COMMIT 的处理。源码价值在于确认失败窗口,而不是证明配置项存在。
生产验收只认端到端证据
| 层次 | 必须留下的证据 |
|---|---|
| Source | Kafka partition/offset 或源库位点 |
| Flink | Checkpoint ID、完成时间、恢复点 |
| Doris | Label、事务状态、可见时间、过滤行数 |
| 业务 | 主键集合、最大版本、Delete 结果、关键金额聚合 |
如果任意一层无法通过同一批次标识串起来,就不能宣称端到端 Exactly-Once。
Exactly-Once 不是一个开关,而是 Source 位点、Checkpoint、事务身份、表模型和业务版本共同闭合的一条证据链。
官方资料
- Transactions
- Load Transactions
- Flink Doris Connector
- Data Update Overview