如何从Kafka事务生产者超时引发的ProducerFencedException中恢复?
Kafka事务生产者超时后的恢复方案
我正在测试一个基于Kotlin的Kafka事务生产者应用,覆盖以下场景:
- 正常路径:开启事务 → 发送消息 → 提交事务 → 断言消费者能拉取到消息(运行正常)
- 生产者抛出异常:开启事务 → 发送消息 → 抛出RuntimeException → 中止事务 → 断言消费者拉取不到消息(运行正常)
- 生产者超时(触发ProducerFencedException):开启事务 → 发送消息 → 线程休眠模拟超时 → 断言消费者拉取不到消息,但后续所有操作都会抛出
ProducerFencedException
我知道org.apache.kafka.common.errors.ProducerFencedException是当客户端使用旧Epoch的生产者(即Kafka文档中的"僵尸生产者")时抛出的,但问题是:如何从事务生产者的超时情况中恢复?有没有办法让生产者感知自身已失效并完成重置?
测试代码实现
Main.kt
val producer = Producer() fun main() { happyPath() timeoutTestSync() happyPath() // 原代码此处会失败,因为超时后的生产者已失效 } fun happyPath() { val consumer = Consumer() assert(consumer.readClean().isEmpty) assert(consumer.readDirty().isEmpty) producer.produceInTransaction("test${RandomStringUtils.randomAlphanumeric(1)}", false, false) assert(consumer.readClean().count() == 1) } fun timeoutTestSync() { val consumer = Consumer() assert(consumer.readClean().isEmpty) assert(consumer.readDirty().isEmpty) producer.produceInTransaction("test${RandomStringUtils.randomAlphanumeric(1)}", true, false) assert(consumer.readDirty().isEmpty) assert(consumer.readClean().isEmpty) }
Consumer.kt
private const val INPUT_TOPIC = "output" class Consumer { private val consumer: KafkaConsumer<String, String> = createKafkaConsumer(false) private val dirtConsumer: KafkaConsumer<String, String> = createKafkaConsumer(true) private fun createKafkaConsumer(dirtRead:Boolean): KafkaConsumer<String, String> { val props:MutableMap<String, Any> = HashMap() props[ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG] = "localhost:9099" props[ConsumerConfig.CLIENT_ID_CONFIG] = "client-${RandomStringUtils.randomAlphanumeric(3)}" props[ConsumerConfig.GROUP_ID_CONFIG] = "group-${RandomStringUtils.randomAlphanumeric(3)}" props[ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG] = StringDeserializer::class.java props[ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG] = StringDeserializer::class.java props[ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG] = true props[ConsumerConfig.AUTO_OFFSET_RESET_CONFIG] = "latest" if (dirtRead) props[ConsumerConfig.ISOLATION_LEVEL_CONFIG] = "read_uncommitted" else props[ConsumerConfig.ISOLATION_LEVEL_CONFIG] = "read_committed" val consumer:KafkaConsumer<String, String> = KafkaConsumer<String, String>(props) consumer.subscribe(Collections.singletonList(INPUT_TOPIC)) return consumer } fun readClean(): ConsumerRecords<String, String> { return consumer.poll(Duration.ofMillis(10000)) } fun readDirty(): ConsumerRecords<String, String> { return dirtConsumer.poll(Duration.ofMillis(10000)) } }
Producer.kt(原版本)
private const val OUTPUT_TOPIC = "output" class Producer { private val producer: KafkaProducer<String, String> = createKafkaProducer() private fun createKafkaProducer(): KafkaProducer<String, String> { val props = Properties() props[ProducerConfig.BOOTSTRAP_SERVERS_CONFIG] = "localhost:9099" props[ProducerConfig.ENABLE_IDEMPOTENCE_CONFIG] = "true" props[ProducerConfig.TRANSACTION_TIMEOUT_CONFIG] = 5000 props[ProducerConfig.TRANSACTIONAL_ID_CONFIG] = "prod-${RandomStringUtils.randomAlphanumeric(3)}-" props[ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG] = "org.apache.kafka.common.serialization.StringSerializer" props[ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG] = "org.apache.kafka.common.serialization.StringSerializer" val producerOut:KafkaProducer<String,String> = KafkaProducer(props) producerOut.initTransactions() return producerOut } fun produceInTransaction(message:String, delay:Boolean, exception:Boolean) { try { producer.beginTransaction() val x = producer.send(ProducerRecord<String, String>(OUTPUT_TOPIC, "1", message)) when { delay -> { Thread.sleep(10000) } exception -> { throw RuntimeException("Simulated exception") } } producer.commitTransaction() } catch (e:RuntimeException) { producer.abortTransaction() println("Exception in transaction: ${e}") } } }
超时恢复的实现方案
当事务超时触发ProducerFencedException后,当前生产者实例已被Kafka集群标记为"僵尸生产者",无法再参与任何事务操作,必须销毁旧实例并重建新的事务生产者,同时重新初始化事务上下文。
修改Producer类支持重置
调整Producer类,将生产者实例改为可变变量,并添加重置方法,在捕获到ProducerFencedException时自动重建:
private const val OUTPUT_TOPIC = "output" class Producer { private var producer: KafkaProducer<String, String> = createKafkaProducer() private fun createKafkaProducer(): KafkaProducer<String, String> { val props = Properties() props[ProducerConfig.BOOTSTRAP_SERVERS_CONFIG] = "localhost:9099" props[ProducerConfig.ENABLE_IDEMPOTENCE_CONFIG] = "true" props[ProducerConfig.TRANSACTION_TIMEOUT_CONFIG] = 5000 props[ProducerConfig.TRANSACTIONAL_ID_CONFIG] = "prod-${RandomStringUtils.randomAlphanumeric(3)}-" props[ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG] = "org.apache.kafka.common.serialization.StringSerializer" props[ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG] = "org.apache.kafka.common.serialization.StringSerializer" val producerOut = KafkaProducer<String,String>(props) producerOut.initTransactions() return producerOut } private fun resetProducer() { // 关闭旧生产者释放资源 producer.close(Duration.ofMillis(1000)) // 重建新生产者并初始化事务 producer = createKafkaProducer() } fun produceInTransaction(message:String, delay:Boolean, exception:Boolean) { try { producer.beginTransaction() val x = producer.send(ProducerRecord<String, String>(OUTPUT_TOPIC, "1", message)) when { delay -> { Thread.sleep(10000) } exception -> { throw RuntimeException("Simulated exception") } } producer.commitTransaction() } catch (e: RuntimeException) { try { producer.abortTransaction() } catch (fencedEx: ProducerFencedException) { // 捕获到围栏异常,重置生产者 println("Producer fenced, resetting instance: ${fencedEx}") resetProducer() } catch (otherEx: Exception) { println("Failed to abort transaction: ${otherEx}") } println("Exception in transaction: ${e}") } } }
关键说明
- 捕获ProducerFencedException:在尝试中止事务时,如果抛出该异常,说明当前生产者已失效,必须立即重置。
- 关闭旧生产者:重置前必须调用
close()释放旧实例的资源,避免连接泄漏。 - 重新初始化事务:新生产者创建后必须调用
initTransactions(),这一步会向Kafka集群注册新的事务Epoch,确保新实例能正常参与事务。 - 事务超时配置:
TRANSACTION_TIMEOUT_CONFIG的值不能超过Kafka broker端的transaction.max.timeout.ms(默认900000ms),否则生产者初始化会失败。
测试代码适配
修改后,timeoutTestSync执行完成后,后续的happyPath()调用会使用重置后的生产者,能够正常执行。
内容的提问来源于stack exchange,提问作者marionmaiden
相关产品推荐
相关产品推荐

