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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.25 09:15:32