如何让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事务已完成提交,导致数据不一致。
具体实现步骤
绑定全局事务
Flink的Exactly-Once依赖两阶段提交(2PC)机制,需确保Kafka和JDBC Sink都参与同一全局事务。当前两者已开启Exactly-Once配置,需确认共享Flink事务协调器,无独立事务配置覆盖。调整JDBC重试策略
移除JdbcExecutionOptions中withMaxRetries(0)的设置,适当增加重试次数(如3次)。若重试后仍失败,Flink会触发Checkpoint失败,进而回滚所有未提交事务(包括Kafka事务),避免单边提交。优化Checkpoint配置
开启Flink Checkpoint,设置合理的间隔(如1分钟)和超时时间(需大于数据库最大恢复窗口)。Checkpoint失败时,Flink会中止所有关联事务,Kafka未提交的事务消息不会被消费者读取。统一数据处理边界
确保kafkaProducerRecordStream与keyedEventLogRecord为同源数据,经过相同分区和处理逻辑,保证单条数据的Kafka发送与数据库更新在同一Checkpoint周期内完成。恢复机制保障
数据库恢复后,Flink从最近成功Checkpoint重启,重新处理未完成批次,确保两项操作要么都成功,要么都失败,最终实现Exactly-Once语义。
内容的提问来源于stack exchange,提问作者Ladu anand

