如何为Spring Kafka默认生产者工厂配置maxAge属性?
配置Spring Kafka生产者工厂的maxAge属性
根据Spring Kafka文档说明:
从2.5.8版本开始,您现在可以在生产者工厂上配置maxAge属性。这对于使用可能在broker的transactional.id.expiration.ms时间内处于空闲状态的事务型生产者非常有用。在当前kafka-clients中,这可能会在没有重新平衡的情况下导致ProducerFencedException。通过将maxAge设置为小于transactional.id.expiration.ms的值,工厂会在生产者超过其最大时长时刷新它。
具体配置位置分两种场景:
1. 自定义Java配置类
如果通过Java代码手动创建生产者工厂,直接在DefaultKafkaProducerFactory实例上调用setMaxAge()方法:
import org.springframework.kafka.core.DefaultKafkaProducerFactory; import java.util.Map; // 构建包含bootstrap.servers等基础配置的参数集合 Map<String, Object> producerConfigs = ...; DefaultKafkaProducerFactory<String, Object> producerFactory = new DefaultKafkaProducerFactory<>(producerConfigs); // 设置maxAge(单位毫秒),需小于broker的transactional.id.expiration.ms值 producerFactory.setMaxAge(3600000); // 示例:1小时
2. Spring Boot自动配置(配置文件)
如果使用Spring Boot自动配置的默认生产者工厂,直接在配置文件中添加工厂属性:
application.properties:
# 基础生产者配置 spring.kafka.producer.bootstrap-servers=your-broker-address # 设置生产者工厂的maxAge,单位毫秒 spring.kafka.producer.factory.max-age=3600000
application.yaml:
spring: kafka: producer: bootstrap-servers: your-broker-address factory: max-age: 3600000
注意:配置的maxAge值必须小于Kafka Broker端的transactional.id.expiration.ms参数值(默认15分钟/900000毫秒),才能有效避免ProducerFencedException问题。
内容的提问来源于stack exchange,提问作者Andras Hatvani
相关产品推荐
相关产品推荐

