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

如何让Flink中两个Sink在事务中协同实现Exactly-Once语义

问题描述

我有一个Flink流水线,包含两个Sink:一个用于向Kafka发布消息,另一个用于更新数据库状态。我通过Kafka事务和XA JDBC事务来实现Exactly-Once语义,但遇到如下问题:当数据库宕机而Kafka Broker正常运行时,即使数据库在Flink设定的重试次数耗尽后仍未恢复,Flink仍会向Broker发送并提交消息(已验证消费者仅读取已提交消息)。我的需求是确保消息发送与对应记录更新仅执行一次,请问该如何实现?

keyedEventLog流从JDBC源获取输入,随后转换为kafkaProducerRecordStream供Kafka Sink使用。

Kafka Sink代码

CustomKafkaUtil.writeRecords(kafkaProducerRecordStream, writeConfig, StringSerializer::new, "oiwp-event-log");

Custom Kafka Util类代码

CustomKafkaUtil extends KafkaUtil {

    public static <V> void writeRecords(DataStream<KafkaProducerRecord<byte[], V>> producerRecords, EventbusWriteRuntimeConf writeConfig, SerializableSupplier<Serializer<V>> serializerSupplier, String uidSuffix) {
        uidSuffix = uidSuffix != null ? uidSuffix : writeConfig.topic;
        producerRecords
                .sinkTo(getCustomSink(writeConfig))
                .name("kafka-sink-" + uidSuffix)
                .uid("kafka-sink-" + uidSuffix);
    }

    private static KafkaSink<Object> getCustomSink(EventbusWriteRuntimeConf writeConfig) {
        String topic = writeConfig.topic;
        Properties producerProperties = SppConfigUtil.instance.getKafkaProducerConfig(topic, writeConfig.cluster, writeConfig.region);
        if (!writeConfig.kafkaProperties.isEmpty()) {
            producerProperties.putAll(writeConfig.kafkaProperties);
        }
        log.info("producer properties: {}", producerProperties);
        return getSinkBuilder(producerProperties)
                .setRecordSerializer(new KafkaRecordSerializationSchema<Object>() {
                    private static final long serialVersionUID = 5597069351310493251L;

                    public ProducerRecord<byte[], byte[]> serialize(Object element, KafkaRecordSerializationSchema.KafkaSinkContext context, Long timestamp) {
                        KafkaProducerRecord<byte[], byte[]> kafkaRecord = (KafkaProducerRecord<byte[], byte[]>) element;
                        return new ProducerRecord<>(kafkaRecord.topic, kafkaRecord.partition, kafkaRecord.timestamp, kafkaRecord.key, kafkaRecord.value, KafkaUtil.getRecordHeaders(kafkaRecord.headers));
                    }
                })
                .setDeliveryGuarantee(DeliveryGuarantee.EXACTLY_ONCE)
                .setTransactionalIdPrefix("oiwp-kafka-sink-random")
                .build();
    }
}

JDBC Sink代码

keyedEventLogRecord.addSink(
                JdbcSink.exactlyOnceSink(
                        "update " + applicationConfig.getDbConfig().getSchema() + "." + "event_log " + "set event_status = 'SUCCESS' where id = ?",
                        ((preparedStatement, row) -> {
                            ObjectMapper objectMapper = new ObjectMapper();
                            JsonNode jsonNode = null;
                            try {
                                jsonNode = objectMapper.readTree(row.f0);
                            } catch (JsonProcessingException e) {
                                throw new RuntimeException(e);
                            }
                            preparedStatement.setString(1, jsonNode.get("id").asText());
                        }),
                        JdbcExecutionOptions.builder()
                                .withMaxRetries(0)
                                .build(),
                        JdbcExactlyOnceOptions.builder().withAllowOutOfOrderCommits(false)
                                .withTransactionPerConnection(true)
                                .build(),
                        new PGXADataSourceProvider(
                                applicationConfig.getDbConfig().getJdbcUrl() + "/" + applicationConfig.getDbConfig().getDatabase(),
                                applicationConfig.getDbConfig().getUsername(),
                                applicationConfig.getDbConfig().getPassword()
                        ))).name("oiwp-event-log-jdbc-sink").uid("oiwp-event-log-jdbc-sink");
解决方案

问题根源

当前两个Sink的事务未绑定到同一Flink全局事务中,Kafka Sink的事务提交与JDBC Sink的XA事务独立执行。当JDBC事务因数据库宕机失败时,Kafka事务已完成提交,导致数据不一致。

具体实现步骤

  1. 绑定全局事务
    Flink的Exactly-Once依赖两阶段提交(2PC)机制,需确保Kafka和JDBC Sink都参与同一全局事务。当前两者已开启Exactly-Once配置,需确认共享Flink事务协调器,无独立事务配置覆盖。

  2. 调整JDBC重试策略
    移除JdbcExecutionOptions中withMaxRetries(0)的设置,适当增加重试次数(如3次)。若重试后仍失败,Flink会触发Checkpoint失败,进而回滚所有未提交事务(包括Kafka事务),避免单边提交。

  3. 优化Checkpoint配置
    开启Flink Checkpoint,设置合理的间隔(如1分钟)和超时时间(需大于数据库最大恢复窗口)。Checkpoint失败时,Flink会中止所有关联事务,Kafka未提交的事务消息不会被消费者读取。

  4. 统一数据处理边界
    确保kafkaProducerRecordStream与keyedEventLogRecord为同源数据,经过相同分区和处理逻辑,保证单条数据的Kafka发送与数据库更新在同一Checkpoint周期内完成。

  5. 恢复机制保障
    数据库恢复后,Flink从最近成功Checkpoint重启,重新处理未完成批次,确保两项操作要么都成功,要么都失败,最终实现Exactly-Once语义。

内容的提问来源于stack exchange,提问作者Ladu anand

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.19 04:05:57