阿里二面:RocketMQ 消费者拉取一批消息,其中部分消费失败了,偏移量怎样更新?

大家好,我是君哥。

最近有读者参加面试时被问了一个问题,如果消费者拉取了一批消息,比如 100 条,第 100 条消息消费成功了,但是第 50 条消费失败,偏移量会怎样更新?就着这个问题,今天来聊一下,如果一批消息有消费失败的情况时,偏移量怎么保存。

1 拉取消息

1.1 封装拉取请求

以 RocketMQ 推模式为例,RocketMQ 消费者启动代码如下:

public static void main(String[] args) throws InterruptedException, MQClientException { DefaultMQPushConsumer consumer = new DefaultMQPushConsumer("CID_JODIE_1"); consumer.subscribe("TopicTest", "*"); consumer.setConsumeFromWhere(ConsumeFromWhere.CONSUME_FROM_FIRST_OFFSET); consumer.setConsumeTimestamp("20181109221800"); consumer.registerMessageListener(new MessageListenerConcurrently() {@Overridepublic ConsumeConcurrentlyStatus consumeMessage(List msgs, ConsumeConcurrentlyContext context){ try{System.out.printf("%s Receive New Messages: %s %n", Thread.currentThread().getName(), msgs); }catch (Exception e){return ConsumeConcurrentlyStatus.RECONSUME_LATER; } return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;} }); consumer.start();}

上面的 DefaultMQPushConsumer 是一个推模式的消费者,启动方法是 start。消费者启动后会触发重平衡线程(RebalanceService),这个线程的任务是在死循环中不停地进行重平衡,最终封装拉取消息的请求到 pullRequestQueue。这个过程涉及到的 UML 类图如下:

☞☞☞AI 智能聊天, 问答助手, AI 智能搜索, 免费无限量使用 DeepSeek R1 模型☜☜☜

图片

Typewise.app Typewise.app

面向客户服务和销售团队的AI写作解决方案。

Typewise.app 39 查看详情 Typewise.app

1.2 处理拉取请求

封装好拉取消息的请求 PullRequest 后,RocketMQ 就会不停地从 pullRequestQueue 获取消息拉取请求进行处理。UML 类图如下:

图片

拉取消息的入口方法是一个死循环,代码如下:

//PullMessageServicepublic void run(){ log.info(this.getServiceName() + " service started"); while (!this.isStopped()) {try { PullRequest pullRequest = this.pullRequestQueue.take(); this.pullMessage(pullRequest);} catch (InterruptedException ignored) {} catch (Exception e) { log.error("Pull Message Service Run Method exception", e);} } log.info(this.getServiceName() + " service end");}

这里拉取到消息后,提交给 PullCallback 这个回调函数进行处理。

拉取到的消息首先被 put 到 ProcessQueue 中的 msgTreeMap 上,然后被封装到 ConsumeRequest 这个线程类来处理。把代码精简后,ConsumeRequest 处理逻辑如下:

//ConsumeMessageConcurrentlyService.javapublic void run(){ MessageListenerConcurrently listener = ConsumeMessageConcurrentlyService.this.messageListener; ConsumeConcurrentlyContext context = new ConsumeConcurrentlyContext(messageQueue); ConsumeConcurrentlyStatus status = null; try {//1.执行消费逻辑,这里的逻辑是在文章开头的代码中定义的status = listener.consumeMessage(Collections.unmodifiableList(msgs), context); } catch (Throwable e) { } if (!processQueue.isDropped()) {//2.处理消费结果ConsumeMessageConcurrentlyService.this.processConsumeResult(status, context, this); } else {log.warn("processQueue is dropped without process consume result. messageQueue={}, msgs={}", messageQueue, msgs); }}

2 处理消费结果

2.1 并发消息

并发消息处理消费结果的代码做精简后如下:

//ConsumeMessageConcurrentlyService.javapublic void processConsumeResult( final ConsumeConcurrentlyStatus status, final ConsumeConcurrentlyContext context, final ConsumeRequest consumeRequest){ int ackIndex = context.getAckIndex(); switch (status) {case CONSUME_SUCCESS: if (ackIndex >= consumeRequest.getMsgs().size()) {ackIndex = consumeRequest.getMsgs().size() - 1; } int ok = ackIndex + 1; int failed = consumeRequest.getMsgs().size() - ok; break;case RECONSUME_LATER: break;default: break; } switch (this.defaultMQPushConsumer.getMessageModel()) {case BROADCASTING: for (int i = ackIndex + 1; i < consumeRequest.getMsgs().size(); i++) { } break;case CLUSTERING: List msgBackFailed = new ArrayList(consumeRequest.getMsgs().size()); for (int i = ackIndex + 1; i = 0 && !consumeRequest.getProcessQueue().isDropped()) {this.defaultMQPushConsumerImpl.getOffsetStore().updateOffset(consumeRequest.getMessageQueue(), offset, true); }}

从上面的代码可以看出,如果处理消息的逻辑是串行的,比如文章开头的代码使用 for 循环来处理消息,那如果在某一条消息处理失败了,直接退出循环,给 ConsumeConcurrentlyContext 的 ackIndex 变量赋值为消息列表中失败消息的位置,这样这条失败消息后面的消息就不再处理了,发送给 Broker 等待重新拉取。代码如下:

public static void main(String[] args) throws InterruptedException, MQClientException { DefaultMQPushConsumer consumer = new DefaultMQPushConsumer("CID_JODIE_1"); consumer.subscribe("TopicTest", "*"); consumer.setConsumeFromWhere(ConsumeFromWhere.CONSUME_FROM_FIRST_OFFSET); consumer.setConsumeTimestamp("20181109221800"); consumer.registerMessageListener(new MessageListenerConcurrently() {@Overridepublic ConsumeConcurrentlyStatus consumeMessage(List msgs, ConsumeConcurrentlyContext context){ for (int i = 0; i < msgs.size(); i++) {try{ System.out.printf("%s Receive New Messages: %s %n", Thread.currentThread().getName(), msgs);}catch (Exception e){ context.setAckIndex(i); return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;} } return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;} }); consumer.start();}

消费成功的消息则从 ProcessQueue 中的 msgTreeMap 中移除,并且返回 msgTreeMap 中最小的偏移量(firstKey)去更新。注意:集群模式偏移量保存在 Broker 端,更新偏移量需要发送消息到 Broker,而广播模式偏移量保存在 Consumer 端,只需要更新本地偏移量就可以。

如果处理消息的逻辑是并行的,处理消息失败后给 ackIndex 赋值是没有意义的,因为可能有多条消息失败,给 ackIndex 变量赋值并不准确。最好的方法就是给 ackIndex 赋值 0,整批消息全部重新消费,这样又可能带来冥等问题。

2.2 顺序消息

对于顺序消息,从 msgTreeMap 取出消息后,先要放到 consumingMsgOrderlyTreeMap 上面,更新偏移量时,是从 consumingMsgOrderlyTreeMap 上取最大的消息偏移量(lastKey)。

3 总结

回到开头的问题,如果一批消息按照顺序消费,是不可能出现第 100 条消息消费成功了,但第 50 条消费失败的情况,因为第 50 条消息失败的时候,应该退出循环,不再继续进行消费。

如果是并发消费,如果出现了这种情况,建议是整批消息全部重新消费,也就是给 ackIndex 赋值 0,这样必须考虑冥等问题。

以上就是阿里二面:RocketMQ 消费者拉取一批消息,其中部分消费失败了,偏移量怎样更新?的详细内容,更多请关注创想鸟其它相关文章!

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

(0)
打赏 微信扫一扫 微信扫一扫 支付宝扫一扫 支付宝扫一扫
如何在电脑上高效浏览微信公众号,让阅读更加便捷
上一篇 2025年11月27日 11:32:56
javascript怎么将16进制转为2进制
下一篇 2025年11月27日 11:32:59

相关推荐

  • VSCode的悬浮提示信息如何自定义?

    通过JSDoc或docstring添加注释可直接影响VSCode悬浮提示内容,如JavaScript/TypeScript中使用/* /格式、Python中使用三引号文档字符串,配合Pylance等扩展增强显示;安装语言支持扩展可提升提示丰富度;高级场景可通过开发自定义语言服务器,在textDocu…

    2026年9月20日
    500
  • AI赋能蛋白质研究:SaprotHub让蛋白质AI模型训练和调用不再有门槛!

    AI赋能蛋白质研究:SaprotHub让蛋白质AI模型训练和调用不再有门槛!AI赋能蛋白质研究:SaprotHub让蛋白质AI模型训练和调用不再有门槛!AI赋能蛋白质研究:SaprotHub让蛋白质AI模型训练和调用不再有门槛!AI赋能蛋白质研究:SaprotHub让蛋白质AI模型训练和调用不再有门槛!

    编辑 | scienceai 近年来,AI 技术在蛋白质研究领域发挥了越来越重要的作用。从 AlphaFold2 在结构预测任务上的脱颖而出,到各类蛋白质语言模型(PLMs)在功能预测方面的重大进展,生物研究者们可以利用各式各样的 AI 模型来辅助他们的研究。 然而,随着模型变得越来越复杂,如何训练…

    2026年9月7日 用户投稿
    200
  • 清华大学研究组首次在双重编码量子比特之间实现高保真度量子纠缠门

    清华大学科研团队在双重编码量子比特量子逻辑门研究中取得突破性进展,成功在钡离子超精细能级中实现了双重编码量子比特间的直接纠缠逻辑门。这项研究成果显著降低了量子线路复杂度,为未来量子纠错和离子-光子量子网络应用铺平了道路。 离子阱技术是构建大规模通用量子计算机的热门选择,已成功实现高保真量子比特操作。…

    2026年9月6日
    000
  • 实现5Å全原子RMSD,普渡大学深度学习方法准确预测RNA三级结构,登Nature子刊

    实现5Å全原子RMSD,普渡大学深度学习方法准确预测RNA三级结构,登Nature子刊实现5Å全原子RMSD,普渡大学深度学习方法准确预测RNA三级结构,登Nature子刊实现5Å全原子RMSD,普渡大学深度学习方法准确预测RNA三级结构,登Nature子刊实现5Å全原子RMSD,普渡大学深度学习方法准确预测RNA三级结构,登Nature子刊

    ☞☞☞AI 智能聊天, 问答助手, AI 智能搜索, 免费无限量使用 DeepSeek R1 模型☜☜☜ 编辑 | 萝卜皮 非编码 RNA 在各种生物功能中发挥着调控作用,并且与人类健康、药物设计等领域息息相关。 了解功能的机械机制需要三级结构信息,然而,通过实验确定 RNA 三维结构成本高昂且耗时…

    2026年9月5日 用户投稿
    100
  • AI for Science:北大、东方理工等团队用人工智能在实验数据中挖掘潜在规律

    AI for Science:北大、东方理工等团队用人工智能在实验数据中挖掘潜在规律AI for Science:北大、东方理工等团队用人工智能在实验数据中挖掘潜在规律AI for Science:北大、东方理工等团队用人工智能在实验数据中挖掘潜在规律AI for Science:北大、东方理工等团队用人工智能在实验数据中挖掘潜在规律

    ☞☞☞AI 智能聊天, 问答助手, AI 智能搜索, 免费无限量使用 DeepSeek R1 模型☜☜☜ 编辑 | ScienceAI‍‍ 科学研究的核心之一是发现能够描述自然现象的规律性方程。这些方程不仅能加深我们对自然的理解,还能为复杂问题的解决提供明确指导。 然而,许多领域,尤其是材料和化学等…

    2026年9月4日 用户投稿
    300
  • 70年AI研究得出了《苦涩的教训》:为什么说AI创业也在重复其中的错误?

    70年AI研究得出了《苦涩的教训》:为什么说AI创业也在重复其中的错误?70年AI研究得出了《苦涩的教训》:为什么说AI创业也在重复其中的错误?70年AI研究得出了《苦涩的教训》:为什么说AI创业也在重复其中的错误?70年AI研究得出了《苦涩的教训》:为什么说AI创业也在重复其中的错误?

    人人都在做垂直 AI 产品,为什么要反其道而行? Scaling Laws 是否失灵,这个话题从 2024 年年尾一直讨论至今,也没有定论。 Ilya Sutskever 在 NeurIPS 会上直言:大模型预训练这条路可能已经走到头了。上周的 CES 2025,黄仁勋有提到,在英伟达看来,Scal…

    2026年9月4日 用户投稿
    200
  • 刚刚,奥特曼剧透GPT-4.5、GPT-5重大更新,o3取消独立发布

    刚刚,奥特曼剧透GPT-4.5、GPT-5重大更新,o3取消独立发布刚刚,奥特曼剧透GPT-4.5、GPT-5重大更新,o3取消独立发布刚刚,奥特曼剧透GPT-4.5、GPT-5重大更新,o3取消独立发布刚刚,奥特曼剧透GPT-4.5、GPT-5重大更新,o3取消独立发布

    奥特曼深夜一则推文,在网络上掀起了讨论狂潮。 没有一点点预告,奥特曼亲自公布自家产品路线图,并承认公司最近发布的一些产品有些混乱。 推文透露,OpenAI 的下一步是发布 GPT-4.5,这是其最后一个非思维链 (CoT) 模型。它的内部代号是 Orion,不过奥特曼没有解释 GPT-4.5 将具有…

    2026年9月2日 用户投稿
    100
  • 从概念到应用,清华团队开发DeepTFBU工具包助力基因表达精准调控

    从概念到应用,清华团队开发DeepTFBU工具包助力基因表达精准调控从概念到应用,清华团队开发DeepTFBU工具包助力基因表达精准调控从概念到应用,清华团队开发DeepTFBU工具包助力基因表达精准调控从概念到应用,清华团队开发DeepTFBU工具包助力基因表达精准调控

    ☞☞☞AI 智能聊天, 问答助手, AI 智能搜索, 免费无限量使用 DeepSeek R1 模型☜☜☜ 编辑 | 萝卜皮 增强子通过与转录因子 (TF) 相互作用,在各种生物过程中充当基因表达的关键调节器。虽然转录因子结合位点 (TFBS) 被广泛认为是 TF 结合和增强子活性的关键决定因素,但其…

    2026年8月31日 用户投稿
    800
  • 注册Deepseek账号时,哪些信息是必填的?

    必填信息:1、邮箱注册;2、手机号码注册;3、第三方社交平台注册。注册成功后通常还需要填写一些基本个人信息,如昵称、性别、生日等。 1、DeepSeek官网入口☜☜☜☜☜点击保存 2、DeepSeek官方登录入口☜☜☜☜☜点击保存 ☞☞☞AI 智能聊天, 问答助手, AI 智能搜索, 免费无限量使用…

    2026年8月29日
    100
  • 耶鲁、剑桥等开发MindLLM,将脑成像直接转换为文本

    耶鲁、剑桥等开发MindLLM,将脑成像直接转换为文本耶鲁、剑桥等开发MindLLM,将脑成像直接转换为文本耶鲁、剑桥等开发MindLLM,将脑成像直接转换为文本耶鲁、剑桥等开发MindLLM,将脑成像直接转换为文本

    ☞☞☞AI 智能聊天, 问答助手, AI 智能搜索, 免费无限量使用 DeepSeek R1 模型☜☜☜ 编辑 | 萝卜皮 将功能性磁共振成像 (fMRI) 信号解码为文本一直是神经科学界面临的一项重大挑战,它有望推动脑机接口的发展,并加深对大脑机制的了解。然而,现有的方法往往存在预测性能不佳、任务…

    2026年8月28日 用户投稿
    000
  • 高通牵手诺基亚贝尔实验室,展示无线网络中可互操作的多厂商AI的价值

    高通牵手诺基亚贝尔实验室,展示无线网络中可互操作的多厂商AI的价值高通牵手诺基亚贝尔实验室,展示无线网络中可互操作的多厂商AI的价值高通牵手诺基亚贝尔实验室,展示无线网络中可互操作的多厂商AI的价值高通牵手诺基亚贝尔实验室,展示无线网络中可互操作的多厂商AI的价值

    ai赋能无线网络:高通与诺基亚联合展示多厂商ai互操作性 本文介绍高通和诺基亚在2024年世界移动通信大会(MWC 2024)上展示的成果:一个基于AI的、可互操作的多厂商无线网络系统。该系统显著提升了网络吞吐量,并在不同物理环境中展现出强大的鲁棒性。 关键发现: AI模型在各种物理环境(包括室内外…

    2026年8月28日 用户投稿
    100
  • 除了Transformer架构,还有哪些常用的大模型架构

    常见大模型架构多样。RNN 处理序列,却因梯度问题难应对长序列;其变体 LSTM 借门控机制改善,GRU 则简化结构提效率。CNN 从计算机视觉起步,借卷积等提取特征,后拓展应用。GAN 用于生成,借生成与判别对抗训练。VAE 融合自编码器与变分推断生成多样样本 。 ☞☞☞AI 智能聊天, 问答助手…

    2026年8月25日
    000
  • java中的consumer关键字用途 消费者Consumer的2个典型应用

    java中的consumer关键字用途 消费者Consumer的2个典型应用java中的consumer关键字用途 消费者Consumer的2个典型应用java中的consumer关键字用途 消费者Consumer的2个典型应用java中的consumer关键字用途 消费者Consumer的2个典型应用

    java中的consumer接口用于定义不返回结果的操作,其核心目的是简化代码并提升可读性与维护性。1. 它常用于集合的foreach方法,实现更简洁的遍历操作;2. 在stream api中通过peek和foreach方法支持中间处理与最终操作;3. 可自定义多参数consumer接口以满足特定需…

    2026年8月25日 用户投稿
    000
  • 如何用HTML插入标签云组件_HTML CSS3变换与随机颜色生成算法

    使用HTML构建标签结构,CSS3添加旋转与过渡效果,JavaScript生成随机HSL颜色并设置字体大小,实现动态交互的标签云组件。 要在网页中实现一个动态的标签云组件,结合 HTML、CSS3 变换和随机颜色生成算法,可以按照以下步骤操作。这个组件不仅能提升页面视觉效果,还能通过色彩和旋转增加交…

    2025年12月23日
    300
  • 如何在Go Gin应用中集成前端JavaScript模块(如Sentry)

    本文探讨了在Go Gin框架下,通过HTML模板服务前端页面时,如何有效集成JavaScript模块(如Sentry)。针对浏览器不直接支持Node.js模块导入语法的问题,文章详细阐述了利用CDN引入Sentry SDK的解决方案,并提供了具体的代码示例,帮助开发者实现前端错误监控功能,避免了复杂…

    2025年12月23日
    000
  • html官网浏览入口_html网站设计免费平台

    html官网浏览入口在https://www.codepen.io,该平台支持实时预览代码、创建Pen项目、Fork开源示例,可添加外部资源,具备点赞评论收藏等社区互动功能,设有挑战活动与作品集分类,开放API接口,界面简洁适合初学者,在线编写无需配置环境,支持多种预处理器和响应式测试。 html官…

    2025年12月23日
    000
  • html如何修改日期样式

    在html中,可以使用“::-webkit-datetime-edit”伪元素选择器来修改日期格式,只需要用该选择器选中元素,在设置具体样式即可,具体语法为“::-webkit-datetime-edit{属性:属性值}”。 本教程操作环境:windows7系统、CSS3&&HTML…

    2025年12月21日
    100
  • 单选框的type属性值为什么

    单选框的type属性值为“radio”。html type属性可以规定要显示的输入框“”元素的类型;值为“radio”时显示为单选框、“checkbox”时显示为复选框、“select”时显示为下拉式选框等等。 本教程操作环境:windows7系统、HTML5版、Dell G3电脑。 在HTML中,…

    2025年12月21日
    000
  • HTML中type是什么意思

    在HTML中,type是类型的意思,是一个标签属性,主要用于定义标签元素的类型或文档(脚本)的MIME类型;例在input标签中type属性可以规定input元素的类型,在script标签中type属性可以规定脚本的MIME类型。 本教程操作环境:windows7系统、html5版、Dell G3电…

    2025年12月21日
    200
  • HTML中ul标签如何去掉点?HTML无序列表的样式实例解析

    HTML中ul标签如何去掉点?HTML无序列表的样式实例解析HTML中ul标签如何去掉点?HTML无序列表的样式实例解析HTML中ul标签如何去掉点?HTML无序列表的样式实例解析HTML中ul标签如何去掉点?HTML无序列表的样式实例解析

    本篇文章主要讲述的是关于html中的ul标签的默认小点给取消掉,还有关于html的无序列表ul标签的样式解释,给出了ul标签中的type属性三种值的介绍。现在就让我们一起来看本篇文章吧 首先这篇文章一开始我们就开始介绍在html中是怎么把ul标签的点给去掉的: 大家应该都使用过ul无序列表标签,ul…

    2025年12月21日 用户投稿
    500

发表回复

登录后才能评论
关注微信