You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.24 03:40:12