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

如何从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}")
        }
    }
}

关键说明

  1. 捕获ProducerFencedException:在尝试中止事务时,如果抛出该异常,说明当前生产者已失效,必须立即重置。
  2. 关闭旧生产者:重置前必须调用close()释放旧实例的资源,避免连接泄漏。
  3. 重新初始化事务:新生产者创建后必须调用initTransactions(),这一步会向Kafka集群注册新的事务Epoch,确保新实例能正常参与事务。
  4. 事务超时配置:TRANSACTION_TIMEOUT_CONFIG的值不能超过Kafka broker端的transaction.max.timeout.ms(默认900000ms),否则生产者初始化会失败。

测试代码适配

修改后,timeoutTestSync执行完成后,后续的happyPath()调用会使用重置后的生产者,能够正常执行。


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.09 21:05:21