Flink单事务内写入DB与Kafka/Axon的实现方案咨询
现有代码是否满足需求?
不满足。你当前的代码是把同一个数据流分别绑定了两个独立的Sink,JDBC Sink和Kafka Sink各自维护自己的事务周期,两者的提交/回滚动作完全独立。举个例子:如果JDBC插入失败回滚,但Kafka Sink可能已经完成了事务提交;反过来Kafka发送失败时,JDBC的更新可能已经生效,完全做不到“要么都成功,要么都失败”的原子性要求。
怎么实现同一事务内的双Sink操作?
要实现跨JDBC和Kafka的事务原子性,核心是利用Flink的两阶段提交(2PC)机制,通过自定义TwoPhaseCommitSinkFunction来统一管理两个外部系统的事务生命周期。
第一步:写一个自定义的双事务Sink
你需要实现一个继承自TwoPhaseCommitSinkFunction的Sink类,把JDBC和Kafka的事务逻辑整合到一起,统一处理事务的开启、预提交、提交和回滚:
public class DualExactlyOnceSink extends TwoPhaseCommitSinkFunction<FooModel, DualTransactionState, Void> { private transient JdbcConnection jdbcConn; private transient KafkaProducer<String, FooModel> kafkaProducer; private final String kafkaTopic; public DualExactlyOnceSink(String kafkaTopic) { super(TypeInformation.of(DualTransactionState.class), TypeInformation.of(Void.class)); this.kafkaTopic = kafkaTopic; } // 初始化JDBC连接和Kafka生产者 @Override public void open(Configuration parameters) throws Exception { super.open(parameters); // 复用你原来FooJdbcSink的连接逻辑 jdbcConn = FooJdbcSink.createConnection(); // 初始化Kafka生产者,开启事务支持 Properties kafkaProps = new Properties(); kafkaProps.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "你的Kafka broker地址"); // 每个并行子任务用独立的transactional.id,避免冲突 kafkaProps.put(ProducerConfig.TRANSACTIONAL_ID_CONFIG, "dual-sink-tx-" + getRuntimeContext().getIndexOfThisSubtask()); kafkaProducer = new KafkaProducer<>(kafkaProps, new StringSerializer(), new JsonSerializer<>()); kafkaProducer.initTransactions(); } // 开启新事务:同时启动JDBC和Kafka的事务 @Override protected DualTransactionState beginTransaction() throws Exception { jdbcConn.setAutoCommit(false); kafkaProducer.beginTransaction(); // 用时间戳作为事务标识,方便跟踪 return new DualTransactionState(System.currentTimeMillis()); } // 处理每条数据:同时写入JDBC和Kafka(暂存,不提交) @Override protected void invoke(DualTransactionState transaction, FooModel value, Context context) throws Exception { // 复用原来FooJdbcSink的插入/更新逻辑 FooJdbcSink.executeInsertOrUpdate(jdbcConn, value); // 把数据发送到Topic2,此时消息只在Kafka事务中暂存 kafkaProducer.send(new ProducerRecord<>(kafkaTopic, value.getId(), value)); } // 预提交:确保两个系统的操作都已就绪,不会在提交时失败 @Override protected void preCommit(DualTransactionState transaction) throws Exception { // Flush Kafka生产者,确保消息已经到达Kafka broker kafkaProducer.flush(); // JDBC这边不需要额外操作,因为还没执行commit } // 提交事务:同时提交JDBC和Kafka的事务 @Override protected void commit(DualTransactionState transaction) { try { jdbcConn.commit(); kafkaProducer.commitTransaction(); } catch (Exception e) { throw new RuntimeException("提交事务失败", e); } } // 回滚事务:任一系统失败时,同时回滚两个事务 @Override protected void abort(DualTransactionState transaction) { try { jdbcConn.rollback(); kafkaProducer.abortTransaction(); } catch (Exception e) { throw new RuntimeException("回滚事务失败", e); } } // 关闭资源 @Override public void close() throws Exception { super.close(); if (jdbcConn != null) jdbcConn.close(); if (kafkaProducer != null) kafkaProducer.close(); } // 自定义事务状态类,用来跟踪事务 public static class DualTransactionState implements Serializable { private final long transactionId; public DualTransactionState(long transactionId) { this.transactionId = transactionId; } } }
第二步:替换原有代码中的双Sink
把原来两个独立的Sink删掉,换成这个自定义的双事务Sink:
DataStream<AxonMessage> stream = env.fromSource(axon.source(Constants.CONSUMER_TOPIC_NAME, Constants.CONSUMER_GROUP_ID), WatermarkStrategy.noWatermarks(), "foo-kafka-source") .map(axonMessage -> (FooModel) axonMessage.getPayload()); // 使用自定义的双事务Sink,传入Topic2的名称 stream.addSink(new DualExactlyOnceSink(Constants.TOPIC2_NAME)) .name("dual-exactly-once-sink") .uid("dual-exactly-once-sink");
必配的Flink检查点设置
两阶段提交依赖Flink的检查点机制,所以必须开启检查点并配置Exactly-Once模式:
// 每5秒触发一次检查点,时间间隔根据你的业务调整 env.enableCheckpointing(5000); // 设置检查点模式为Exactly-Once env.getCheckpointConfig().setCheckpointingMode(CheckpointingMode.EXACTLY_ONCE); // 两个检查点之间至少间隔3秒,避免资源占用过高 env.getCheckpointConfig().setMinPauseBetweenCheckpoints(3000);
为什么不能用两个独立的ExactlyOnceSink?
每个ExactlyOnceSink都是独立参与Flink检查点流程的,它们的事务提交时机完全独立。哪怕两个Sink都配置了Exactly-Once,也没法保证它们的提交/回滚动作完全同步,自然做不到跨系统的原子性。只有通过自定义的TwoPhaseCommitSinkFunction,把两个外部系统的事务绑定到同一个检查点周期里,才能实现“要么都成,要么都败”的要求。
内容的提问来源于stack exchange,提问作者Khushboo
相关产品推荐
相关产品推荐

