Flink Checkpoint 和 Doris 2PC 都成功,订单9001却从已支付退回待支付。没有重复事务,也没有丢失事件;一条更旧的 CDC 记录最后到达,并被精确提交了一次。
2PC 保证批次提交边界,Unique Key 保证一主键一行,只有可靠 Sequence 才能决定乱序事件谁获胜。
四层机制各自只管一件事
| 层次 | 负责 | 不负责 |
|---|---|---|
| Flink Checkpoint | 算子状态与 Source 位点一致恢复 | 业务事件新旧 |
| Doris Stream Load 2PC | Checkpoint 对应事务只提交一次 | 主键设计和列值正确 |
| Unique Key MoW | 相同 Key 只暴露一个逻辑版本 | 没有 Sequence 时的业务顺序 |
| Sequence | 同 Key 按业务版本选胜者 | 跨 Key 事务约束和漏事件 |
把这四层都叫 Exactly-Once,事故发生后就只能互相甩锅。
复现一次旧状态回写
错误表只有 Unique Key,没有 Sequence:
CREATETABLEorder_state_bad(order_idBIGINT,op_versionBIGINT,statusVARCHAR(16),amountDECIMAL(12,2))UNIQUEKEY(order_id)DISTRIBUTEDBYHASH(order_id)BUCKETS4PROPERTIES("enable_unique_key_merge_on_write"="true");先写版本 105 的已支付,再写迟到的版本 103:
INSERTINTOorder_state_badVALUES(9001,105,'PAID',100.00);INSERTINTOorder_state_badVALUES(9001,103,'PENDING',100.00);后到行成为可见版本。正确表必须把业务版本映射为 Sequence:
CREATETABLEorder_state_good(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");同样顺序写入后,版本 105 仍然获胜。
Sequence 选择比字段本身更重要
| 候选 Sequence | 优点 | 经典失败 |
|---|---|---|
| 应用更新时间 | 直观 | 时钟回拨、同秒冲突、不同服务不可比 |
| 数据库自增版本 | 单表稳定 | 分库后不一定全序 |
| Binlog 文件与 Position | 源库提交顺序可靠 | 需要编码成可比较值,跨源无天然顺序 |
| Kafka Offset | 单 Partition 单调 | 跨 Partition 不可直接比较 |
| 业务状态序号 | 符合领域规则 | 需要业务明确状态机和补偿语义 |
不存在脱离 Source 拓扑的万能 Sequence。一个订单的所有事件如果可能分散到多个 Kafka Partition,Offset 就不能直接做全局版本。
DDL 演进会制造另一种旧数据
上游增加字段后,Flink 状态、Connector 映射和 Doris Schema 的生效顺序不一致,可能产生:新字段丢失、默认值覆盖、部分更新失败或整批过滤。
安全顺序应是:
目标表先兼容新字段 → Connector 能识别旧与新两种事件 → 上游开始发送新 Schema → 对账新字段非空率与默认值 → 最后清理旧兼容逻辑不要用 Checkpoint 成功证明 Schema Evolution 成功。Checkpoint 只保存作业当时接受到的状态。
排障需要两张表
Unique Key 服务表只能看胜者,失败版本可能已被 Delete Bitmap 隐藏。保留一张 Duplicate Key 审计表,记录 Source 位点、事件类型、业务版本和原始载荷:
SELECTorder_id,op_version,op_type,source_offset,ingest_timeFROMorder_event_auditWHEREorder_id=9001ORDERBYop_version,ingest_time;再与服务表比较:
SELECTorder_id,op_version,statusFROMorder_state_goodWHEREorder_id=9001;审计表回答发生过什么,服务表回答现在是什么。没有审计表时,状态被覆盖后很难证明错误来自源端乱序还是 Sink 合并。
源码证据落在 Delete Bitmap
Doris4.0.8的BaseTablet::calc_segment_delete_bitmap会查找历史 Rowset 的同 Key 行。历史行 Sequence 更大时,新到行自身进入 Delete Bitmap;新行 Sequence 更大时,历史位置进入 Delete Bitmap。
Checkpoint 决定这批是否提交 Sequence 决定同 Key 哪行获胜 Delete Bitmap 决定查询跳过哪行 Compaction 最终回收失败版本这条链解释了为什么不用原地更新也能立即查询唯一状态,也解释了写入侧要承担主键查找和位图维护成本。
修复不能只补一个属性
现有错误表不能靠修改属性安全变成有 Sequence 的正确表。生产修复通常需要新建目标表、从审计事件按最大业务版本重建、双写对账、切换读流量,再停止旧表。
验收至少覆盖 Insert、Update、Delete、同版本冲突、跨 Partition 乱序、Checkpoint 恢复和 DDL 变化。
Flink 可以把每条事件可靠地送到终点,但只有业务版本才能告诉 Doris 哪条事件代表现在。
官方资料与源码
- Flink Doris Connector
- Data Update Overview
- Unique Key Model
base_tablet.cppat 4.0.8