Kafka批处理监听器中反序列化异常的重试策略与实现

Kafka批处理监听器中反序列化异常的重试策略与实现

本文详细介绍了如何在Spring Kafka批处理监听器中有效处理并重试反序列化异常。通过修改DefaultErrorHandler以取消对DeserializationException的致命标记,并结合监听器内部对带有null载荷的消息进行异常信息提取和重新抛出,实现对整个批次消息的重试,从而提高Kafka应用的鲁棒性。

Kafka批处理监听器反序列化异常重试机制

在spring kafka应用中,当使用批处理监听器处理消息时,如果消费者在反序列化阶段遇到瞬时错误(例如,avro schema注册中心连接问题),默认情况下这些deserializationexception会被视为致命错误,导致消息无法被重试,进而可能丢失或被跳过。为了增强应用的健壮性,我们需要一套机制来捕获这些异常并触发重试。

默认行为与挑战

Spring Kafka的ErrorHandlingDeserializer在反序列化失败时,通常会返回null作为消息载荷,并将原始异常信息存储在消息头中。然而,默认的DefaultErrorHandler会将DeserializationException视为不可重试的异常类型。这意味着即使ErrorHandlingDeserializer捕获了异常,DefaultErrorHandler也不会触发消息的重试。对于批处理监听器,如果批次中的任何一条消息反序列化失败,整个批次都可能受到影响。

实现反序列化异常重试的步骤

要实现反序列化异常的重试,我们需要分两步进行:

修改错误处理器,允许DeserializationException重试。在监听器中识别并重新抛出反序列化异常。

1. 配置DefaultErrorHandler以允许重试

DefaultErrorHandler默认将DeserializationException标记为不可重试。我们需要通过调用removeClassification方法将其从致命异常列表中移除。

import org.springframework.context.annotation.Bean;import org.springframework.context.annotation.Configuration;import org.springframework.kafka.annotation.EnableKafka;import org.springframework.kafka.config.ConcurrentKafkaListenerContainerFactory;import org.springframework.kafka.core.DefaultKafkaConsumerFactory;import org.springframework.kafka.listener.ContainerProperties;import org.springframework.kafka.listener.DefaultErrorHandler;import org.springframework.kafka.support.serializer.ErrorHandlingDeserializer;import org.springframework.kafka.support.serializer.SerializationUtils;import org.apache.kafka.common.serialization.Deserializer;import org.apache.kafka.common.serialization.StringDeserializer;import org.springframework.kafka.support.serializer.DeserializationException;import org.springframework.boot.autoconfigure.kafka.KafkaProperties;import java.util.List;import java.util.Map;@Configuration@EnableKafkapublic class KafkaConfiguration {    @Bean("myContainerFactory")    public ConcurrentKafkaListenerContainerFactory createFactory(            KafkaProperties properties    ) {        ConcurrentKafkaListenerContainerFactory factory = new ConcurrentKafkaListenerContainerFactory();        factory.setConsumerFactory(                new DefaultKafkaConsumerFactory(                        properties.buildConsumerProperties(),                        new StringDeserializer(),                        new ErrorHandlingDeserializer(new MyDeserializer()) // 使用自定义反序列化器                )        );        factory.getContainerProperties().setAckMode(                ContainerProperties.AckMode.MANUAL_IMMEDIATE        );        // 创建并配置DefaultErrorHandler        DefaultErrorHandler errorHandler = new DefaultErrorHandler();        // 移除DeserializationException的致命标记,使其可以重试        errorHandler.removeClassification(DeserializationException.class);        factory.setCommonErrorHandler(errorHandler);        return factory;    }    // 模拟偶尔失败的反序列化器    static class MyDeserializer implements Deserializer {        private int retries = 0;        @Override        public void configure(Map configs, boolean isKey) {            // No-op        }        @Override        public String deserialize(String topic, byte[] bytes) {            String s = new String(bytes);            // 模拟第一次遇到包含"7"的字符串时抛出异常,第二次成功            if (s.contains("7") && retries == 0) {                retries = 1; // 标记已尝试一次                System.out.println("Simulating deserialization error for: " + s);                throw new DeserializationException("Simulated deserialization failure for: " + s, bytes, false);            }            retries = 0; // 重置计数器            System.out.println("Deserialized successfully: " + s);            return s;        }        @Override        public void close() {            // No-op        }    }}

在上述配置中,我们创建了一个DefaultErrorHandler实例,并通过removeClassification(DeserializationException.class)方法明确指示Kafka,当遇到DeserializationException时,不应将其视为致命错误,而是应该尝试重试。

2. 在监听器中处理null载荷并重新抛出异常

即使DefaultErrorHandler被配置为允许重试DeserializationException,ErrorHandlingDeserializer在反序列化失败时仍然会向监听器发送一个null载荷的消息。为了触发批次的重试,监听器需要检查这些null载荷,从消息头中提取原始的反序列化异常,并将其重新抛出。

序列猴子开放平台 序列猴子开放平台

具有长序列、多模态、单模型、大数据等特点的超大规模语言模型

序列猴子开放平台 0 查看详情 序列猴子开放平台

批处理监听器需要接收List<Message>而不是List,以便能够访问消息头。

import org.springframework.stereotype.Component;import org.springframework.kafka.annotation.KafkaListener;import org.springframework.kafka.support.Acknowledgment;import org.springframework.kafka.support.KafkaHeaders;import org.springframework.kafka.support.serializer.DeserializationException;import org.springframework.kafka.support.serializer.SerializationUtils;import org.springframework.messaging.Message;import org.springframework.messaging.handler.annotation.Header;import org.springframework.kafka.listener.ListenerUtils;import java.util.List;import java.util.Objects;@Componentpublic class StringListener {    @KafkaListener(            topics = {"string-test"},            groupId = "test",            batch = "true",            containerFactory = "myContainerFactory"    )    public void listen(List<Message> messages, Acknowledgment acknowledgment) {        boolean hasDeserializationError = false;        for (Message message : messages) {            String payload = message.getPayload();            if (payload == null) {                // 载荷为null,检查是否是反序列化异常                byte[] exceptionHeader = (byte[]) message.getHeaders().get(SerializationUtils.VALUE_DESERIALIZER_EXCEPTION_HEADER);                if (exceptionHeader != null) {                    DeserializationException deserializationException =                         ListenerUtils.byteArrayToDeserializationException(exceptionHeader);                    if (deserializationException != null) {                        System.err.println("Detected deserialization error for a message in batch: " + deserializationException.getMessage());                        hasDeserializationError = true;                        // 这里不直接抛出,而是标记,待循环结束后统一处理,确保检查完所有消息                    }                }            } else {                System.out.println("Processed message: " + payload);            }        }        if (hasDeserializationError) {            // 如果批次中存在反序列化错误,则重新抛出异常,触发整个批次的重试            System.err.println("Batch contains deserialization errors. Re-throwing to trigger retry.");            // 抛出任意RuntimeException即可,DefaultErrorHandler会根据配置进行重试            throw new RuntimeException("Batch failed due to deserialization error(s).");        }        // 如果没有反序列化错误,或者所有错误都已处理且不需重试,则提交偏移量        acknowledgment.acknowledge();        System.out.println("Batch processed and acknowledged successfully.");    }}

在上述监听器代码中:

我们接收List<Message>以便访问消息头。遍历批次中的每条消息。如果payload为null,则尝试从SerializationUtils.VALUE_DESERIALIZER_EXCEPTION_HEADER头中提取原始的DeserializationException。使用ListenerUtils.byteArrayToDeserializationException()辅助方法将字节数组转换回DeserializationException对象。如果检测到DeserializationException,则设置hasDeserializationError标志。在遍历完所有消息后,如果hasDeserializationError为true,则重新抛出一个RuntimeException。这个异常会被DefaultErrorHandler捕获,由于我们已将DeserializationException从致命列表中移除,DefaultErrorHandler会根据其配置的重试策略(例如,指数退避)来重试整个批次。

注意事项:

批次重试的粒度: 重新抛出异常会导致整个批次的消息被重试,而不是仅仅重试失败的那一条消息。这意味着批次中已经成功处理的消息也会被重新处理。在设计业务逻辑时需要考虑幂等性。Spring Kafka版本: 确保使用的Spring Kafka版本支持DefaultErrorHandler的removeClassification方法(通常在2.8.4及更高版本中可用)。对于批处理错误处理器的分类行为,较新的版本(如2.9.x及更高版本)对FallbackBatchErrorHandler的异常分类处理更为完善。异常类型: 即使是从消息头中提取的DeserializationException,最终在监听器中重新抛出的可以是任何RuntimeException。DefaultErrorHandler会根据其内部的异常分类规则来决定是否重试。

总结

通过上述配置和代码实现,我们成功地为Kafka批处理监听器添加了反序列化异常的重试机制。这使得应用程序能够更优雅地处理瞬时的数据反序列化问题,避免了消息丢失,并提高了系统的容错能力。在实际应用中,应根据业务需求仔细调整DefaultErrorHandler的重试策略(如退避间隔、最大重试次数等),并确保消息处理的幂等性。

以上就是Kafka批处理监听器中反序列化异常的重试策略与实现的详细内容,更多请关注创想鸟其它相关文章!

版权声明:本文内容由互联网用户自发贡献,该文观点仅代表作者本人。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。
如发现本站有涉嫌抄袭侵权/违法违规的内容, 请发送邮件至 chuangxiangniao@163.com 举报,一经查实,本站将立刻删除。
发布者:程序猿,转转请注明出处:https://www.chuangxiangniao.com/p/216590.html

赞 (0)
打赏 微信扫一扫 微信扫一扫 支付宝扫一扫 支付宝扫一扫
赛力斯拟在港交所上市,H 股发行获中国证监会备案
上一篇 2025年11月3日 15:08:55
创意改造(创意改造让家居环保时尚更具个性)
下一篇 2025年11月3日 15:09:00

相关推荐

  • LINUX连接不上WiFi怎么办_LINUX系统WiFi连接失败排查指南

    LINUX连接不上WiFi怎么办_LINUX系统WiFi连接失败排查指南LINUX连接不上WiFi怎么办_LINUX系统WiFi连接失败排查指南LINUX连接不上WiFi怎么办_LINUX系统WiFi连接失败排查指南LINUX连接不上WiFi怎么办_LINUX系统WiFi连接失败排查指南

    首先检查无线网卡是否被系统识别,通过lspci或lsusb命令确认硬件存在;若识别正常但无法连接,需安装对应驱动如firmware-iwlwifi或rtl88x2bu-dkms;确保NetworkManager服务已启动并启用;使用nmcli命令扫描并连接WiFi网络;若仍失败,可手动编辑Netpl…

    2026年9月26日 • 用户投稿
    400
  • sublime怎么快速打开最近的项目_sublime访问历史项目的快捷方法

    sublime怎么快速打开最近的项目_sublime访问历史项目的快捷方法sublime怎么快速打开最近的项目_sublime访问历史项目的快捷方法sublime怎么快速打开最近的项目_sublime访问历史项目的快捷方法sublime怎么快速打开最近的项目_sublime访问历史项目的快捷方法

    Sublime Text 通过命令面板访问最近项目,按 Ctrl+Shift+P 输入“project”选择“Project: Switch Project”即可打开历史项目,结合保存项目文件可高效管理多工作区。 Sublime Text 没有直接打开“最近项目”的独立快捷键,但可以通过命令面板快速…

    2026年9月26日 • 用户投稿
    000
  • Java 方法中数组参数的正确调用方式

    Java 方法中数组参数的正确调用方式Java 方法中数组参数的正确调用方式Java 方法中数组参数的正确调用方式Java 方法中数组参数的正确调用方式

    本文旨在阐述如何在 Java 方法中正确传递和使用数组参数。通过一个实际的例子,我们将详细讲解如何创建数组、将其作为参数传递给方法,以及如何在方法内部访问和操作数组元素。掌握这些技巧对于编写高效且易于维护的 Java 代码至关重要。 在 Java 编程中,方法经常需要接收数组作为参数,以便对一组数据…

    2026年9月26日 • 用户投稿
    000
  • 安装 Windows 10 时,提示 “计算机的磁盘空间不足”,如何清理?

    安装 Windows 10 时,提示 “计算机的磁盘空间不足”,如何清理?安装 Windows 10 时,提示 “计算机的磁盘空间不足”,如何清理?安装 Windows 10 时,提示 “计算机的磁盘空间不足”,如何清理?安装 Windows 10 时,提示 “计算机的磁盘空间不足”,如何清理?

    首先需明确是全新安装还是升级安装,通常全新安装更易解决空间不足问题。在Windows 10安装界面按Shift+F10打开命令提示符,输入diskpart进入分区工具,执行list disk查看磁盘,select disk X选择目标磁盘(X为磁盘编号),再通过list partition查看分区情…

    2026年9月26日 • 用户投稿
    100
  • win8桌面图标不见了_Win8桌面图标恢复

    win8桌面图标不见了_Win8桌面图标恢复win8桌面图标不见了_Win8桌面图标恢复win8桌面图标不见了_Win8桌面图标恢复win8桌面图标不见了_Win8桌面图标恢复

    首先检查桌面图标显示设置,右键桌面选择“查看”并勾选“显示桌面图标”;若无效,通过任务管理器重启Windows资源管理器进程;如仍无改善,可删除%localappdata%目录下的IconCache.db文件以重建图标缓存;最后使用系统自带的桌面疑难解答工具进行自动修复。 如果您发现Windows …

    2026年9月26日 • 用户投稿
    000
  • 从Scanner读取单个字符时处理空格的问题

    从Scanner读取单个字符时处理空格的问题从Scanner读取单个字符时处理空格的问题从Scanner读取单个字符时处理空格的问题从Scanner读取单个字符时处理空格的问题

    本文旨在解决Java中使用Scanner读取用户输入时,由于Scanner默认以空格作为分隔符,导致读取单个字符时出现的问题。我们将深入探讨Scanner的工作原理,并提供使用Scanner.nextLine()方法读取整行输入来解决此问题的方案,确保程序能够正确处理包含空格的输入。 在使用Java…

    2026年9月26日 • 用户投稿
    100
  • grokAI平台官方网站主页 grokAI 智能助手入口官方直达地址

    GrokAI平台官方网站主页是https://grok.com/,用户可直接访问该网址进入。新用户无需注册即可点击“Start Chatting”体验基础功能,登录X账号则可使用高级服务。 ☞☞☞AI 智能聊天, 问答助手, AI 智能搜索, 免费无限量使用 DeepSeek R1 模型☜☜☜ Gr…

    2026年9月26日
    100
  • Z790主板比B760主板强在哪里?

    Z790主板相比B760在超频支持、供电能力与扩展性上更强。1. Z790支持带“K”后缀CPU超频,B760不支持;2. Z790供电模组更豪华,散热设计更强,可应对高功耗CPU长时间满载;3. Z790拥有8条DMI通道,扩展接口更丰富,支持更多高速设备;4. Z790对高频内存支持更好,内存超…

    2026年9月26日
    100
  • 从 0 开始学 V8 漏洞利用之 V8 通用利用链(二)

    作者:hcamael@知道创宇404实验室 相关阅读:从 0 开始学 V8 漏洞利用之环境搭建(一)经过一段时间的研究,先进行一波总结,不过因为刚开始研究没多久,也许有一些局限性,以后如果发现了,再进行修正。 概述 ‍我认为,在搞漏洞利用前都得明确目标。比如打CTF做二进制的题目,大部分情况下,目标…

    2026年9月26日
    100
  • windows怎么关闭快速启动_快速启动功能关闭教程

    windows怎么关闭快速启动_快速启动功能关闭教程windows怎么关闭快速启动_快速启动功能关闭教程windows怎么关闭快速启动_快速启动功能关闭教程windows怎么关闭快速启动_快速启动功能关闭教程

    关闭快速启动可解决Windows启停慢或硬件兼容问题,需通过控制面板、电源高级设置或命令提示符操作。2. 控制面板路径:运行control→硬件和声音→电源选项→选择电源按钮功能→更改不可用设置→取消勾选启用快速启动→保存。3. 高级设置路径:电源选项→更改计划设置→更改高级电源设置→展开电源按钮和…

    2026年9月26日 • 用户投稿
    200
  • x浏览器如何拦截弹窗广告_x浏览器弹窗广告拦截教程

    x浏览器如何拦截弹窗广告_x浏览器弹窗广告拦截教程x浏览器如何拦截弹窗广告_x浏览器弹窗广告拦截教程x浏览器如何拦截弹窗广告_x浏览器弹窗广告拦截教程x浏览器如何拦截弹窗广告_x浏览器弹窗广告拦截教程

    开启x浏览器广告拦截功能可有效屏蔽弹窗广告。首先在设置中启用“广告过滤”并选择强力模式;其次通过自定义规则添加已知广告域名进行精准拦截;接着在隐私与安全设置中开启“阻止弹出窗口”开关,阻断脚本触发的弹窗;最后可使用轻阅读模式简化网页结构,避免广告加载,提升浏览体验。 如果您在浏览网页时频繁遇到弹窗广…

    2026年9月26日 • 用户投稿
    300
  • 强!荣耀 Magic V5 官宣搭载 6100mAh 青海湖刀片电池

    强!荣耀 Magic V5 官宣搭载 6100mAh 青海湖刀片电池强!荣耀 Magic V5 官宣搭载 6100mAh 青海湖刀片电池强!荣耀 Magic V5 官宣搭载 6100mAh 青海湖刀片电池强!荣耀 Magic V5 官宣搭载 6100mAh 青海湖刀片电池

    官方消息透露,7 月 2 日晚 19:00,荣耀将召开 magic v5 及 ai 终端生态发布会。届时,荣耀 magic v5 等多款旗舰新品将同步登场。早在 6 月 25 日,荣耀就已为 magic v5 开启预热宣传。据 cnmo 掌握的信息,这款折叠屏手机搭载了容量高达 6100mah 的青…

    2026年9月26日 • 用户投稿
    100
  • MAC连接的移动硬盘速度很慢_Mac外置硬盘读写速度慢问题排查

    MAC连接的移动硬盘速度很慢_Mac外置硬盘读写速度慢问题排查MAC连接的移动硬盘速度很慢_Mac外置硬盘读写速度慢问题排查MAC连接的移动硬盘速度很慢_Mac外置硬盘读写速度慢问题排查MAC连接的移动硬盘速度很慢_Mac外置硬盘读写速度慢问题排查

    答案:Mac连接移动硬盘速度慢可能因存储不足、接口问题或硬盘故障等导致。应清理硬盘空间至10%-15%以上,更换为USB 3.0及以上数据线并直连主机端口,使用“磁盘工具”检查健康状况并修复错误,通过“活动监视器”终止高占用进程,并确保硬盘采用APFS或exFAT等合适文件系统以提升性能。 如果您在…

    2026年9月26日 • 用户投稿
    300
  • NVMe驱动器的SLC缓存用完后性能下降多少?

    NVMe驱动器的SLC缓存用完后性能下降多少?NVMe驱动器的SLC缓存用完后性能下降多少?NVMe驱动器的SLC缓存用完后性能下降多少?NVMe驱动器的SLC缓存用完后性能下降多少?

    NVMe驱动器在SLC缓存耗尽后写入速度会骤降至数十到两百MB/s,具体取决于NAND类型、容量和主控方案,QLC型号甚至可能低于机械硬盘速度。 NVMe驱动器在SLC缓存耗尽后,性能会经历显著的下降,通常写入速度会从数百甚至数千MB/s骤降至数十到两百MB/s的水平,具体取决于驱动器采用的NAND…

    2026年9月26日 • 用户投稿
    100
  • 什么是线程池?为什么使用线程池?ThreadPoolExecutor有哪些核心参数?

    什么是线程池?为什么使用线程池?ThreadPoolExecutor有哪些核心参数?什么是线程池?为什么使用线程池?ThreadPoolExecutor有哪些核心参数?什么是线程池?为什么使用线程池?ThreadPoolExecutor有哪些核心参数?什么是线程池?为什么使用线程池?ThreadPoolExecutor有哪些核心参数?

    线程池通过复用预先创建的线程,避免频繁创建销毁带来的开销,提升系统性能与稳定性。ThreadPoolExecutor是Java中实现线程池的核心类,其核心参数包括corePoolSize(核心线程数)、maximumPoolSize(最大线程数)、keepAliveTime(非核心线程空闲存活时间)…

    2026年9月26日 • 用户投稿
    100
  • 小米 MIX Flip 2 小折叠等深微曲屏,屏幕更平整

    小米 MIX Flip 2 小折叠等深微曲屏,屏幕更平整小米 MIX Flip 2 小折叠等深微曲屏,屏幕更平整小米 MIX Flip 2 小折叠等深微曲屏,屏幕更平整小米 MIX Flip 2 小折叠等深微曲屏,屏幕更平整

    小米宣布将于 6 月 26 日晚 7 点召开人车家全生态新品发布会,届时将正式推出小米 mix flip 2 小折叠屏手机。 今日,官方率先曝光了该机的外观设计与核心配置信息。 据介绍,小米 MIX Flip 2 采用三面等深微曲屏幕设计,搭配金属磨砂中框;内置全新转轴结构,官方强调“屏幕平整度令人…

    2026年9月26日 • 用户投稿
    100
  • 伊津野英昭腾讯原创3A新情报:融合鬼泣、龙信精华!

    伊津野英昭腾讯原创3A新情报:融合鬼泣、龙信精华!伊津野英昭腾讯原创3A新情报:融合鬼泣、龙信精华!伊津野英昭腾讯原创3A新情报:融合鬼泣、龙信精华!伊津野英昭腾讯原创3A新情报:融合鬼泣、龙信精华!

    据automatonmedia报道,《鬼泣》系列总监、《龙之信条》系列主导者伊津野英昭近日在接受《fami通》采访时,分享了他离开卡普空后首个新项目的最新进展。 伊津野在卡普空工作长达30年,于2024年8月正式离职,并加入腾讯,出任光子工作室日本分部负责人。他目前正在主导开发的首款作品,是一款面向…

    2026年9月26日 • 用户投稿
    000
  • win8怎么映射网络驱动器_Win8网络驱动器映射方法

    win8怎么映射网络驱动器_Win8网络驱动器映射方法win8怎么映射网络驱动器_Win8网络驱动器映射方法win8怎么映射网络驱动器_Win8网络驱动器映射方法win8怎么映射网络驱动器_Win8网络驱动器映射方法

    首先通过计算机界面或运行命令映射网络驱动器,选择驱动器号并输入UNC路径,可勾选登录时重新连接;其次确保网络发现、文件共享及相关服务(Server、Workstation、TCP/IP NetBIOS Helper)已启用,以保证共享访问正常。 如果您需要在Windows 8系统中访问局域网中的共享…

    2026年9月26日 • 用户投稿
    200
  • sublime怎么使用多光标编辑_sublime多光标操作技巧详解

    sublime怎么使用多光标编辑_sublime多光标操作技巧详解sublime怎么使用多光标编辑_sublime多光标操作技巧详解sublime怎么使用多光标编辑_sublime多光标操作技巧详解sublime怎么使用多光标编辑_sublime多光标操作技巧详解

    掌握Sublime Text多光标编辑可大幅提升效率:1. 按Ctrl/Cmd点击或多行列选添加多个光标;2. 用Ctrl+D逐个选中相同词批量修改;3. Ctrl+Shift+Alt+G全选所有匹配项全局替换;4. Alt+Shift拖动进行列选择,实现块状编辑,适合对齐或加前缀后缀。 Subli…

    2026年9月26日 • 用户投稿
    100
  • debian邮件服务器如何实现自动回复

    debian邮件服务器如何实现自动回复debian邮件服务器如何实现自动回复debian邮件服务器如何实现自动回复debian邮件服务器如何实现自动回复

    在debian系统搭建自动回复邮件服务器,只需简单几步即可实现。本文将指导您配置postfix邮件服务器,实现自动回复功能。 一、安装Postfix 首先,确认Debian系统已安装Postfix。若未安装,请执行以下命令: sudo apt updatesudo apt install postf…

    2026年9月26日 • 用户投稿
    300

发表回复

登录后才能评论
关注微信