Spring XML配置Kafka消费者:异常处理与重试机制求助
Spring XML配置下Kafka消费者异常处理与无限循环问题
背景与需求
- 项目基于Spring XML配置(暂无法迁移至Spring Boot)
- 需为Kafka消费者实现完整的异常处理与重试机制,包括:
- 开箱即用的异常处理能力
- 自定义异常处理逻辑
- 可配置的重试策略
- 缺乏对应的XML配置示例与参考代码
当前问题:垃圾数据触发无限循环
当topic中存在无法反序列化的垃圾数据时,消费者陷入无限循环报错,关键堆栈信息如下:
Caused by: org.apache.kafka.common.errors.RecordDeserializationException: Error deserializing key/value for partition t-employee-0 at offset 4. If needed, please seek past the record to continue consumption. at org.apache.kafka.clients.consumer.internals.CompletedFetch.parseRecord(CompletedFetch.java:309) at org.apache.kafka.clients.consumer.internals.CompletedFetch.fetchRecords(CompletedFetch.java:263) at org.apache.kafka.clients.consumer.internals.AbstractFetch.fetchRecords(AbstractFetch.java:340) at org.apache.kafka.clients.consumer.internals.AbstractFetch.collectFetch(AbstractFetch.java:306) at org.apache.kafka.clients.consumer.KafkaConsumer.pollForFetches(KafkaConsumer.java:1235) at org.apache.kafka.clients.consumer.KafkaConsumer.poll(KafkaConsumer.java:1186) at org.apache.kafka.clients.consumer.KafkaConsumer.poll(KafkaConsumer.java:1159) at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.pollConsumer(KafkaMessageListenerContainer.java:1666) at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.doPoll(KafkaMessageListenerContainer.java:1641) at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.pollAndInvoke(KafkaMessageListenerContainer.java:1439) at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.run(KafkaMessageListenerContainer.java:1330) ... 2 common frames omitted Caused by: org.apache.kafka.common.errors.SerializationException: Can't deserialize data from topic [t-employee] at org.springframework.kafka.support.serializer.JsonDeserializer.deserialize(JsonDeserializer.java:588) at org.apache.kafka.common.serialization.Deserializer.deserialize(Deserializer.java:73) at org.apache.kafka.clients.consumer.internals.CompletedFetch.parseRecord(CompletedFetch.java:300) ... 12 common frames omitted
当前context.xml配置
<bean id="employeeProducer" class="com.example.kafka.producer.EmployeeProducer" /> <bean id="defaultKafkaConsumerFactory" class="org.springframework.kafka.core.DefaultKafkaConsumerFactory"> <constructor-arg> <map> <entry key="spring.json.trusted.packages" value="*" /> <entry key="bootstrap.servers" value="localhost:9092"/> <entry key="auto.offset.reset" value="latest"/> <entry key="group.id" value="group1" /> <entry key="client.id" value="my-client-id" /> <entry key="max.poll.records" value="1"/> <entry key="key.deserializer" value="org.springframework.kafka.support.serializer.ErrorHandlingDeserializer" /> <entry key="value.deserializer" value="org.springframework.kafka.support.serializer.ErrorHandlingDeserializer" /> <entry key="spring.deserializer.key.delegate.class" value="org.springframework.kafka.support.serializer.JsonDeserializer" /> <entry key="spring.deserializer.value.delegate.class" value="org.springframework.kafka.support.serializer.JsonDeserializer" /> <entry key="value.class.name" value="com.example.kafka.model.Employee" /> <entry key="spring.json.value.default.type" value="com.example.kafka.model.Employee" /> <entry key="spring.kafka.producer.properties.spring.json.add.type.headers" value="false" /> </map> </constructor-arg> <constructor-arg> <bean id="keyDeserializer" class="org.apache.kafka.common.serialization.StringDeserializer" /> </constructor-arg> <constructor-arg> <bean id="valueDeserializer" class="org.springframework.kafka.support.serializer.JsonDeserializer" /> </constructor-arg> </bean> <bean id="containerProperties" class="org.springframework.kafka.listener.ContainerProperties"> <constructor-arg name="topics" value="t-employee"/> <property name="groupId" value="group1"/> <property name="messageListener" ref="myListener" /> </bean> <bean id="containerListener" class="org.springframework.kafka.listener.KafkaMessageListenerContainer"> <constructor-arg index="0" ref="defaultKafkaConsumerFactory"/> <constructor-arg index="1" ref="containerProperties" /> </bean> <bean id="myListener" class="com.example.kafka.listener.MyMessageListener" />
核心疑问
- 是否需要同时在ConsumerFactory的构造参数和Properties中配置JsonDeserializer?
- Spring XML中如何正确配置Kafka消费者的异常处理与重试机制?
- 如何解决垃圾数据导致的消费者无限循环问题?
内容的提问来源于stack exchange,提问作者Jeff Cook
相关产品推荐
相关产品推荐

