MQ消费失败,自动重试思路


在遇到与第三方系统做对接时,MQ无疑是非常好的解决方案(解耦、异步)。但是如果引入MQ组件,随之要考虑的问题就变多了,如何保证MQ消息能够正常被业务消费。所以引入MQ消费失败情况下,自动重试功能是非常重要的。这里不过细讲MQ有哪些原因会导致失败。

MQ重试,网上有方案一般采用的是,本地消息表+定时任务,不清楚的可以自行了解下。

我这里提供一种另外的思路,供大家参考。方案实现在RabbitMQ(安装延迟队列插件)+.NET CORE 3.1

设计思路为:

内置一个专门做重试的队列,这个队列是一个延迟队列,当业务队列消费失败时,将原始消息投递至重试队列,并设置延迟时间,当延迟时间到达后。重试队列消费会自动将消息重新投递会业务队列,如此便可以实现消息的重试,而且可以根据重试次数来自定义重试时间,比如像微信支付回调一样(第一次延迟3S,第二次延迟10S,第三次延迟60S),上面方案当然要保证MQ消费采用ACK机制。

那么如何让重试队列知道原来的业务队列是哪个,我们定义业务队列时,可以通过MQ的消息头内置一些信息:队列类型(业务队列也有可能是延迟队列)、重试次数(默认为 0)、交换机名称、路由键。业务队列消费失败时,将消息投递至重试队列时,则可以把业务队列的消息头传递至重试队列,那么重试队列消费,重新将消息发送给业务队列时,则可以知道业务队列所需要的所有参数(需要将重试次数+1)。

下面结合代码讲下具体实现:

我们先看看业务队列发送消息时,如何定义

IBasicProperties properties = channel.CreateBasicProperties();
                properties.Persistent = true;
                //初始化,需要内置一些消费异常,自动重试参数 
                if (headers == null)
                {
                    headers = new Dictionary<string, object>();
                }
                //ttlSecond 有值表示消息将投递到延迟队列
                //因为可以自建延迟队列,ttlSecond是业务标识 
                if (ttlSecond.HasValue)
                {
                    if (!headers.ContainsKey("x-delay"))
                    {
                        headers.Add("x-delay", ttlSecond * 1000);
                    }
                    else
                    {
                        headers["x-delay"] = ttlSecond * 1000;
                    }
                    //queueType = 1表示延迟队列 
                    //框架内部重试机制需要此参数,因为重新投递到原始队列时,需要区分普通队列还是延迟队列
                    if (!headers.ContainsKey("queueType"))
                    {
                        headers.Add("queueType", 1);
                    }
                }
                else
                {
                    //queueType = 0表示普通队列
                    if (!headers.ContainsKey("queueType"))
                    {
                        headers.Add("queueType", 0);
                    }
                }
                //重试次数
                if (!headers.ContainsKey("retryCount"))
                {
                    headers.Add("retryCount", 0);
                }
                //原始交换机名称
                if (!headers.ContainsKey("retryExchangeName"))
                {
                    headers.Add("retryExchangeName", exchangeName);
                }
                //原始路由键
                if (!headers.ContainsKey("retryRoutingKey"))
                {
                    headers.Add("retryRoutingKey", routingKey);
                }
                properties.Headers = headers;
                channel.BasicPublish(exchangeName, routingKey, properties, Encoding.UTF8.GetBytes(message));

 这里会内置上面描述的重试队列需要的参数

再来看看业务队列消费如何处理,这里因为会自动重试,所以保证业务队列每次都是消费成功的(MQ才会将消息从队列中删除)

       //每次消费一条
            channel.BasicQos(0, 1, false);

            //定义消费者
            EventingBasicConsumer eventingBasicConsumer = new EventingBasicConsumer(channel);
            eventingBasicConsumer.Received += async (sender, basicConsumer) =>
            {
                string body = Encoding.UTF8.GetString(basicConsumer.Body.ToArray());
                Deadletter deadletter = null;
                try
                {
                    string errorMsg = await action(body);
                    if (!errorMsg.IsNullOrWhiteSpace())
                    {
                        deadletter = new Deadletter() { Body = body, ErrorMsg = errorMsg };
                        _logger.LogError($"业务队列消费异常(已知),消息头:{JsonUtils.Serialize(basicConsumer.BasicProperties.Headers)}{Environment.NewLine}原始消息:{body}{Environment.NewLine}错误:{errorMsg}");
                    }
                }
                catch (Exception ex)
                {
                    deadletter = new Deadletter() { Body = body, ErrorMsg = ex.Message };
                    _logger.LogError(ex, $"业务队列消费异常(未知),消息头:{JsonUtils.Serialize(basicConsumer.BasicProperties.Headers)}{Environment.NewLine}原始消息:{body}");
                }
                //必定应答,不管消费成功还是失败
                channel.BasicAck(basicConsumer.DeliveryTag, false);
                //消费失败,投递消息至重试队列
                if (deadletter != null)
                {
                    PublishRetry(deadletter, basicConsumer.BasicProperties.Headers);
                }
            };

 我们再看看PublishRetry重试队列的推送方法如何实现

IBasicProperties properties = channel.CreateBasicProperties();
                properties.Persistent = true;
                //x-delay为延迟队列的延迟时间
                //如果第一次进行重试,请求头中是不存在延迟时间的,需要新增
                //因为可以进行多次重试,所以第二次时,就会存在延迟时间
                //但因为可以自建用于业务的延迟队列,所以自建的延迟队列,第一次重试也会存在x-delay,但是如果自建的延迟队列失败进行重试时,不能还使用自身的延迟时间,所以需要重新设置为系统默认的失败重试时间
                if (!headers.ContainsKey("x-delay"))
                {
                    headers.Add("x-delay", 0);
                } 
                //重试次数
                int retryCount = Convert.ToInt32(headers["retryCount"]);
                //可以根据重试次数,实现上面说描述的微信回调的重试时间变长效果
                headers["x-delay"] = retryCount * 1000;
                properties.Headers = headers;
                channel.BasicPublish(RETRY_EXCHANGE_NAME, string.Empty, properties, Encoding.UTF8.GetBytes(JsonUtils.Serialize(deadletter)));

重试队列的消费者实现

channel.BasicQos(0, 1, false); 
            EventingBasicConsumer eventingBasicConsumer = new EventingBasicConsumer(channel);
            eventingBasicConsumer.Received += async (sender, basicConsumer) =>
            {
                string message = Encoding.UTF8.GetString(basicConsumer.Body.ToArray());
                Deadletter deadletter = JsonUtils.Deserialize<Deadletter>(message); 
                IDictionary<string, object> headers = basicConsumer.BasicProperties.Headers;
                //请求头中肯定会有如下参数,因为在框架代码中已经内置
                //重试次数
                int retryCount = Convert.ToInt32(headers["retryCount"]);
                //原队列类型,如果原队列本身为延迟队列,重试投递的时候,必须也要为延迟队列,只是不需要延迟时间,投递回原队列后,会立马重新消费
                int queueType = Convert.ToInt32(headers["queueType"]);
                //原队列名称
                string retryExchangeName = Encoding.UTF8.GetString((byte[])headers["retryExchangeName"]);
                //原路由键
                string retryRoutingKey = Encoding.UTF8.GetString((byte[])headers["retryRoutingKey"]);
                if (retryCount <= 10)
                {
                    headers["retryCount"] = retryCount + 1;
                    //原有队列为普通队列,重新投递时,也需要投递为普通队列类型
                    if (queueType == 0)
                    {
                        PublishMessage(retryExchangeName, retryRoutingKey, deadletter.Body, basicConsumer.BasicProperties.Headers);
                    }
                    //原有队列为延迟队列,重新投递时,也需要投递为延迟队列类型
                    else
                    {
                        PublishMessage(retryExchangeName, retryRoutingKey, deadletter.Body, basicConsumer.BasicProperties.Headers, 0);
                    }
                }
                //超过重试最大次数不再处理,交由外部委托来处理死信
                else
                {
                    await deadLetterTask(retryExchangeName, deadletter.Body, deadletter.ErrorMsg);
                }
                //应答
                channel.BasicAck(basicConsumer.DeliveryTag, false);
            };
            //开启监听
            channel.BasicConsume(RETRY_QUEUE_NAME, false, eventingBasicConsumer);

然后在系统中,内置重试队列消费者

//注册框架内自动重试
            _rabbitMQClient.SubscribeRetry(async (exchangeName, message, errorMsg) =>
            {
                string content = $"原始交换机名称:{exchangeName}{Environment.NewLine}" +
                             $"原始消息内容:{message}{Environment.NewLine}" +
                             $"错误消息:{errorMsg}";

                await PushWeChatMessage(content);
            });

 上述为我们MQ实现自动重试的一种方案,当然中间包括每次如果消费失败都可以发送通知,来通知业务人员关注消费失败的情况。可以自定义最大重试次数、重试间隔时间、死信的处理,这里仅仅是MQ重试机制的一种思路而已,大家如果有更好的方案,欢迎多多沟通。

引用来源:https://www.cnblogs.com/jiangbiao/p/15748007.html

版权声明:本文为YES开发框架网发布内容,转载请附上原文出处连接
管理员
上一篇:函数防抖、节流
下一篇:使用.NET 6开发TodoList应用(10)——实现DELETE请求以及HTTP请求幂等性
评论列表

发表评论

评论内容
昵称:
验证码:
验证码
关联文章

MQ消费失败自动思路
nginx配置http自动定向到https
IIS中应用程序池自动停止,启报错
.NET中大型项目开发必备(12)--使用MQ消息队列
IIS中URL定向配置
修改CPU型号(启依然有效)
服务器ntlmssp攻击防御措施,windows server大量审核失败问题
GZUpdate自动升级服务 .NET C/S Winform客户端程序自动升级演示
SAP S/4HANA MM模块培训 41 - 发票计划:周期性、里程碑与自动结算
CentOS7 nginx SSL证书申请并自动续期
openVPN客户端windows开机自动启动
SAP S/4HANA MM模块培训 49 - 自动科目确定:MM-FI集成、事务键与OBYC
记一次.Net Core程序启动失败的排查过程
IIS URL定向VUE配置
GZUpdate自动升级程序客户端演示
Python使用UUID模块云服务器上获取MAC地址,启后就不一样了
fastreport明细自动高度
SAP S/4HANA MM模块培训 30 - 采购订单带分类审批配置与失败排查
.NET 微服务——CI/CD(3):镜像自动分发
.NET 微服务——CI/CD(2):自动打包镜像

热门标签
.NET Core .NET Reactor ag-grid AI发布 api安全 ASP.NET Core C#DLL加密 C#播放声音 C#代码混淆 C#代码加密 ChromeDriver Codex DateTime DBeaver devexpress devTool DLL混淆 edge.js EF EFCore Electron element-ui el-form el-table excel FastReport FileStream FolderBrowerDialog FolderSelectDialog form提交 git gridcontrol gridview input javascript json字符串 JS转换对象JSON jwt JWT授权 linq log Math MCP mitmproxy MVC MySQL Navicat netstat nginx node_modules NSwag Nuget Nuget镜像 number PowerShell pyinstaller python pythoncom python爬虫 python抓包 pywin32 redis Requests-html RestSharp Selenium sql SQL Server Swagger to-cms Visual Studio VSCode vue VueRouter vue路由 VUE页面通讯 Webpack Windows Windows服务 winform wmi xlrd yaml YESCMS YESWEB开发框架 白象 表单提交 播放声音 打开URL 代码混淆 弹窗提醒 端口占用 对象转换 分布式 公共字典 机器码 进程排查 静态资源 开发指南 路由参数 密钥 配置教程 配置文件 权限 人工智能 任务 任务调度 日期间隔 日志 日志记录 省市区 授权验证 数据库 四舍五入 文案 文件读取 文件夹选择 文件目录选择 问题排查 行政区域数据 页面通讯 中间件 CSharp 事务锁 工单系统 并发控制 重复提交 CMS Markdig Markdown markdown-it marked 技术选型 VS Code 开发工具 源代码管理 版本控制 Docker PostgreSQL 时区 部署排查 CMS架构 EF Core 主题系统 二次开发 插件系统 容器 运维命令 镜像清理 Linux NAS 远程挂载 飞牛 fnOS S/4HANA SAP GUI SAP HANA SAP R/3 SAP入门 SAP版本 ERP SAP SAP MM 库存管理 物料管理 采购管理 入门教程 SAP S/4HANA SPRO 企业结构 采购组织 MM01 物料主数据 物料类型 BP分组 业务伙伴 供应商主数据 ME41 RFQ 库存物料 采购流程 ME51 消耗性物料 科目分配 采购申请 AC03 ML81N 外部服务 服务主数据 Business Partner SAP培训 ME51N MM模块 Lean Services MM-SRV 外部服务采购 PIR 供应来源 采购主数据 采购信息记录 ME31K 框架协议 计划协议 采购合同 ME01 供应来源确定 货源清单 MEQ1 供应源确定 配额安排 配额评分 MD04 MD21 MRP 计划文件 需求计划 批量程序 MD01N MD02 MRP Live MD05 MM 物料计划 优化采购 供应源 采购订单 ME2A 供应商确认 采购监控 Flexible Workflow 凭证释放 采购审批 释放策略 实地盘点 物料凭证 货物移动 MIGO 收货 移动类型 已撤回 供应商退货 货物发出 STO 库存转储 转移过账 生产订单 预留 GR/IR MIRO 供应商发票 物流发票校验 OMR2 税码 FI PP SD 实操教程 MRBR OMR6 发票差异 交货成本 后续借记 MI01 实物盘点 盘点差异 公司代码 工厂 组织结构 OMS2 主数据定制 自动科目确定 BP角色 CVI 伙伴确定 编号范围 凭证类型 字段选择 FBN1 OMBT OMC2 会计凭证 OMJJ BOM 委外加工 项目类别L MRKO 供应商寄售 特殊库存K MRKON PIPE Pipeline 特殊库存P ERS MRIS 发票计划 周期性结算 里程碑付款 变更追踪 版本管理 采购凭证 SFTP WebDAV 网盘 飞牛fnOS AMPL HERS MPN 中文教程 库存确定 可用性检查 缺件检查 Output Management 消息确定 输出确定 分割评估 库存计价 评估类别 评估类型 PB00 RM0000 条件技术 采购定价 MM-FI集成 OBYC 库存估价 文本类型 文本采用 EFB EVO MSV SU3 用户参数 发票校验 合同参照 履约保留款 特别总账 预付款 Fiori Launchpad SAP Fiori 应用导航 用户体验 LSMW LTMC Migration Cockpit 数据迁移 BRFplus OPD Output Control My Inbox 审批流程 灵活工作流 SAP PP 外部加工 SAP QM 检验批 质量信息记录 采购收货 SAP PM 维护BOM 维护订单 SAP SD SAP Service 端到端流程 MM模块培训 FI-MM集成 供应商管理 审批配置 FICO入门 SAP FICO 财务配置 供应商税务 预扣税 House Bank 银行对账 客户清账 应收账款 FI控制 验证与替代 印度 GST 税务配置 F110 FBZP EWM入门 SAP EWM 仓库管理 OX14 成本核算 物料评估 后勤配置 物料组 价值更新 数量更新 PP-PI 流程制造 生产计划 容差配置 SAP事务码 SAP基础 TCODE Basis 事务代码 MMNR 编号区间 采购实操 组织架构 OMSF SAP实操 FI配置
联系我们
联系电话:15090125178(微信同号)
电子邮箱:garson_zhang@163.com
站长微信二维码
微信二维码