ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

Flink 精准一次性(Exactly-Once)交付:两阶段提交与 Kafka 事务对接

Flink 精准一次性(Exactly-Once)交付:两阶段提交与 Kafka 事务对接 Flink 精准一次性Exactly-Once交付两阶段提交与 Kafka 事务对接周二凌晨两点财务对账系统突然亮起了最高级别的红灯日终结算总账与各商户分账户流水汇总出现了高达 182 万元的账实不符财务总监连夜把实时流计算组的同学全部叫醒。值班排查的同学满头大汗地盯着 Flink 监控面板“大喜姐我们 Flink 作业的 Checkpoint 明确配置了CheckpointingMode.EXACTLY_ONCE啊集群昨天由于网络抖动发生了一次 TaskManager 崩溃重启但 Flink 明明是从最近的快照恢复状态了为什么下游 Kafka 和 ClickHouse 里的订单数据居然硬生生重复写入了两次”我拉出他们的 Sink 算子代码看了一眼果然是那个经典的“概念混淆事故”代码里仅仅在 Flink Source 和 State 开启了 Exactly-Once下游输出却使用了一个普通的KafkaSink底层连事务支持Semantic.EXACTLY_ONCE都没开启更没有配套设置 Kafka Broker 的事务超时与下游消费者的隔离级别在分布式流计算体系中端到端精准一次性End-to-End Exactly-Once从来不是单靠 Flink 自身打个勾就能凭空实现的。它必须是一场由“上游可重放 Source 核心状态快照机制Chandy-Lamport 下游两阶段提交2PCSink”共同组成的精密三位一体协同战。一、 经典误区状态的一致性 $\neq$ 端到端的一致性很多开发者对 Flink 的 Exactly-Once 存在灾难级的误解以为 Flink 宣称的“精确一次”意味着数据在全链路网络上绝对只被传输、处理了一次。------------------------------------------------------------- | 误区解密数据物理上必然存在重试与重复传输 | ------------------------------------------------------------- Kafka (Source) --- Flink 算子处理 (产生中间状态) --- Kafka / DB (Sink) | v 发生节点崩溃 / 网络中断 [从最近一次 Checkpoint 强行回滚] Kafka 重放消息再次投递 - Flink 状态恢复正确 - Sink 再次写入同样的数据 | v 如果不做 Sink 控制下游物理介质直接遭遇【数据重复插入】Flink 内部的 Exactly-Once指的是状态State的精确一致。当发生故障回滚时算子内部的累加器、计数器状态保证只反映了每条数据恰好计算一次的结果。端到端End-to-End的 Exactly-Once要求不仅内部状态正确连输出到外部物理介质如 Kafka、MySQL、ClickHouse的数据也必须保证不重不漏。如果下游是一个没有事务机制的普通 Sink故障恢复时的重复重放必然导致外部介质数据重叠二、 两阶段提交2PC协议在 Flink 中的流式解构为了让输出端也能达成精确一次Flink 抽象出了著名的TwoPhaseCommitSinkFunction。该协议与 Flink 的 Checkpoint 周期Barrier 对齐严密咬合分为预提交Pre-commit与正式提交Commit两个核心动作------------------------------------------------------------- | 1. 正常流转阶段 (Stream Processing) | | - Source 消费数据Sink 开启 Kafka 事务 (Tx_1)写入待定数据 | ------------------------------------------------------------- | v Checkpoint Barrier 到达 Sink ------------------------------------------------------------- | 2. 第一阶段预提交 (Pre-Commit) | | - Sink 算子冻结事务 Tx_1调用 flush() 确保数据落盘 Broker | | - 将当前事务 ID (Tx_1) 连同状态一并持久化进 Checkpoint Snapshot| | - 立即为后续数据开启全新的下一个事务 (Tx_2) | ------------------------------------------------------------- | v JobManager 广播 NotifyCheckpointComplete ------------------------------------------------------------- | 3. 第二阶段正式提交 (Formal Commit) | | - Sink 算子收到协调者通知正式向 Kafka 发送 Commit Tx_1 指令 | | - 下游消费者此时才被允许读取到 Tx_1 内的数据 | -------------------------------------------------------------如果在第二阶段提交前系统发生了崩溃Flink 重启后会从 Checkpoint 中取出未完成的事务 IDTx_1重新执行 Commit。如果 Checkpoint 制作失败则向 Kafka 发送 Abort 指令将垃圾数据彻底丢弃。三、 生产级端到端 Exactly-Once 核心实现Java以下是我们线上实时交易账务同步中经过严苛压测验证的 Kafka 事务 Sink 与作业级生产配置代码import org.apache.flink.api.common.serialization.SimpleStringSchema; import org.apache.flink.configuration.Configuration; import org.apache.flink.connector.kafka.sink.KafkaRecordSerializationSchema; import org.apache.flink.connector.kafka.sink.KafkaSink; import org.apache.flink.connector.kafka.sink.DeliveryGuarantee; import org.apache.flink.streaming.api.CheckpointingMode; import org.apache.flink.streaming.api.environment.CheckpointConfig; import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment; import java.time.Duration; import java.util.Properties; public class FinancialExactlyOnceJob { public static void main(String[] args) throws Exception { StreamExecutionEnvironment env StreamExecutionEnvironment.getExecutionEnvironment(); // 1. 核心 Checkpoint 基础配置 CheckpointConfig ckConfig env.getCheckpointConfig(); // 模式必须显式设置为 EXACTLY_ONCE ckConfig.setCheckpointingMode(CheckpointingMode.EXACTLY_ONCE); // Checkpoint 触发间隔生产环境建议 1~3 分钟太短会引发高频事务提交风暴 ckConfig.setCheckpointInterval(60_000L); // Checkpoint 必须在 30 秒内制作完毕否则超时丢弃 ckConfig.setCheckpointTimeout(30_000L); // 两次 Checkpoint 之间的最小物理间歇防止算子无休止做快照被拖死 ckConfig.setMinPauseBetweenCheckpoints(20_000L); // 作业取消时保留快照方便人工追溯与救援 ckConfig.setExternalizedCheckpointCleanup( CheckpointConfig.ExternalizedCheckpointCleanup.RETAIN_ON_CANCELLATION ); // 2. 构建具备生产级事务特性的 KafkaSink (基于 Flink 现代 Connector API) Properties kafkaProducerProps new Properties(); // 核心参数Kafka 事务超时时间必须大于 Flink Checkpoint 的最大超时时间 // 默认 Kafka Broker 事务超时为 15 分钟 (900000ms)客户端必须匹配 kafkaProducerProps.setProperty(transaction.timeout.ms, 900000); KafkaSinkString exactlyOnceKafkaSink KafkaSink.Stringbuilder() .setBootstrapServers(kafka-broker1:9092,kafka-broker2:9092) .setRecordSerializer( KafkaRecordSerializationSchema.builder() .setTopic(dwd_trade_financial_ledger) .setValueSerializationSchema(new SimpleStringSchema()) .build() ) // 核心开关DeliveryGuarantee.EXACTLY_ONCE 将自动挂载两阶段事务提交 .setDeliveryGuarantee(DeliveryGuarantee.EXACTLY_ONCE) .setTransactionalIdPrefix(flink-ledger-sink-tx) // 事务前缀必须全局唯一 .setKafkaProducerConfig(kafkaProducerProps) .build(); // 3. 构建拓扑并将处理好的账务流挂载到事务 Sink // env.fromSource(...).map(...).sinkTo(exactlyOnceKafkaSink); env.execute(FinancialExactlyOnceLedgerStream); } }四、 致命陷阱下游消费者的 isolation.level 设置如果只配置了 Flink 端而在下游消费这批数据的应用端漏掉了一个配置前面所有的努力依然会前功尽弃Kafka 的事务机制在物理上是把数据先写进底层 Log 分区中并在末尾追加一条特殊的Control Record控制标记Commit 或 Abort。默认模式isolation.level read_uncommitted下游消费者根本不管这条数据所属的事务最终有没有成功 Commit只要数据物理写入了分区它就会立即读出来这意味着即使 Flink 后来回滚了该事务下游也早已把脏数据吃进去了必须配置isolation.level read_committed消费者会启动过滤缓冲区严格只消费那些打上了成功 Commit 标记的事务数据。对于未完成或已 Abort 的事务数据消费者指针Fetch Offset会优雅跳过从而在物理上杜绝重复与脏读。# 下游消费端 (如 ClickHouse Kafka Engine 或微服务消费组) 必须强制配置 isolation.levelread_committed五、 架构师血泪避坑清单transaction.timeout.ms倒挂引发的死锁灾难Kafka Broker 端的transaction.max.timeout.ms默认是 15 分钟。如果 Flink 遇到严重反压单次 Checkpoint 延迟达到了 16 分钟Kafka 就会提前认为该事务已死亡并强制触发 Abort。当 Flink 终于慢悠悠发出 Commit 请求时会抛出InvalidTxnStateException导致作业无限崩溃重启。客户端事务超时设置必须与集群配额严格统一。事务 ID 前缀TransactionalIdPrefix严禁重名如果有两个不同的 Flink 作业写入同一个 Kafka 集群且配置了相同的TransactionalIdPrefix后启动的作业会强行夺取前一个作业的事务锁Fencing导致先启动的作业直接报错自杀。每个作业的前缀必须硬编码区分或带上唯一 JobID。幂等性写入Idempotent Sink往往比 2PC 更省心两阶段提交极其精密且脆弱会引入数秒至数分钟的事务提交延迟消费者必须等 Checkpoint 结束才能见数。如果下游存储天然支持唯一主键覆盖如 RedisSET、HBasePut、或 MySQL 的INSERT ... ON DUPLICATE KEY UPDATE优先使用“At-Least-Once 业务唯一键幂等覆写”系统鲁棒性会远超脆弱的分布式事务。
返回列表