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

Spring Kafka生产者maxAge配置无效引发InvalidPidMappingException问题

问题:Spring Kafka事务生产者长时间无流量后触发InvalidPidMappingException错误

使用Spring Kafka事务生产者时,若系统超过7天无流量,会出现transactional id元数据丢失的情况,后续写入操作会触发如下错误:

org.apache.kafka.common.errors.InvalidPidMappingException: The producer attempted to use a producer id which is not currently assigned to its transactional id.

目前临时解决方法是重启Kubernetes中的所有实例。

查阅文档得知,Spring Kafka 2.5.8版本新增了maxAge属性,设置其值小于7天可刷新Broker端元数据,但测试后未达预期效果。

为模拟该场景搭建了演示项目,将Broker端事务元数据过期时间设为10秒,相关配置如下:

KAFKA_TRANSACTIONAL_ID_EXPIRATION_MS: 10000
KAFKA_PRODUCER_ID_EXPIRATION_CHECK_INTERVAL_MS: 100
KAFKA_PRODUCER_ID_EXPIRATION_MS: 10000
KAFKA_TRANSACTION_ABORT_TIMED_OUT_TRANSACTION_CLEANUP_INTERVAL_MS: 10000
KAFKA_TRANSACTION_REMOVE_EXPIRED_TRANSACTION_CLEANUP_INTERVAL_MS: 10000

设置maxAge为7秒,期望在元数据过期前(10秒)完成刷新,但仍持续触发InvalidPidMappingException错误:

2023-11-19T13:22:57.234+01:00 DEBUG 28648 --- [fix_localhost-0] o.a.k.c.p.internals.TransactionManager   : [Producer clientId=producer-transactionTestPrefix_localhost-0, transactionalId=transactionTestPrefix_localhost-0] Transition from state COMMITTING_TRANSACTION to error state ABORTABLE_ERROR

org.apache.kafka.common.errors.InvalidPidMappingException: The producer attempted to use a producer id which is not currently assigned to its transactional id.

若未处理ABORTABLE_ERROR(如使用AfterRollbackProcessor或死信队列DLQ),还会导致消息丢失。

使用版本

  • confluentinc/cp-kafka:latest(Kafka 3.4)
  • Spring Boot 3.1.5(JDK 17)
  • Kafka Clients 3.4.1
  • Spring Kafka 3.0.12
  • Spring Stream Binder Kafka 4.0.4

解决方案

需为每个Kafka binder单独设置maxAge属性,示例代码如下:

@Bean
KafkaTransactionManager customKafkaTransactionManager() {
    KafkaMessageChannelBinder kafka = (KafkaMessageChannelBinder)this.binderFactory.getBinder("kafka1", MessageChannel.class);
    DefaultKafkaProducerFactory<byte[], byte[]> producerFactory = (DefaultKafkaProducerFactory<byte[], byte[]>) kafka.getTransactionalProducerFactory();

    producerFactory.setMaxAge(Duration.ofSeconds(60));

    return new KafkaTransactionManager(producerFactory);
}

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.06 06:06:13