消息链路拆分最佳实践:钉钉审批异步链路重构【总结】
阿里妹导读
引入消息队列可以帮助我们解耦业务逻辑,提升性能,让主链路更加清晰。但是消息链路的代码腐化和一致性问题也给业务带来了很多困扰,本文阐述了钉 钉审批消息链路重构的设计和解决方案。注:Metaq 是阿里 RocketMQ 消息队列的内网版本。
概述
引入消息队列可以帮助我们解耦业务逻辑,提升性能,让主链路更加清晰。Metaq 也确实可靠,重试机制能够保障足够的一致性。 钉钉审批将审批中的关键事件,比如审批单发起,任务开始,任务结束以及审批单结束等等作为消息发布出去,将审批的主体流程和周边业务清晰地分开。 但是经过数年的产品迭代,周边业务越来越多,消息链路越来越复杂(详见 “消息链路的美好幻想与残酷现实” 章节): 1、 不同的业务逻辑堆砌在一起互相影响 ,单个消息处理方法就能有上千行代码。 2、因为逻辑太复杂无法实现幂等,只能放弃重试, 不一致问题严重 。每个业务都必须额外开发复杂的对账,补偿机制才能避免客诉。 钉钉审批每分钟都会数十万条消息的发出,从监控可以看出平均每分钟都会有百条消息失败(有的时候会有千条),失败率大约 0.5%。
图中只有顶点处存在失败,其他都是 Sunfire 的自动连线,不存在失败。
消息链路的美好幻想与残酷现实
美好幻想 消息队列在设计之初就给业务规划好了一条康庄大道:- 主业务链路作为 Producer 发出消息
- 数十个甚至更多 Consumer 订阅该消息,分别执行自己快速且幂等的原子逻辑
-
逻辑相互影响,修改风险高;
-
链路脆弱,容易中断,一个调用失败,后续所有逻辑将不会执行;
-
没有重试:大泥球无法做到原子和幂等,整体重试代价太大,所以直接异步执行放弃重试
;
-
消息队列引以为傲的重试功能反而会成为故障的温床,导致雪崩。
为什么不直接把大泥球拆分成前面的多个 Consumer 呢? 这确实也是一种方案,但是对于大泥球 Consumer,可能会拆出几十个 Consumer,这会导致非常严重的读扩散。 举个例子,审批单发起的消息中只含有审批单的 id,内容需要从数据库反查,原本在“大泥球”中,只需要查询一次就复用,而拆分后可能要多查几十次。 这还只是众多扩散问题的其中一个,如果为了治理大泥球,却加重了扩散问题,就得不偿失了。
朴素的拆分想法
最简单的想法就是把大泥球中相关的逻辑聚合到一起,不相关的逻辑隔离,拆成一个个独立的业务处理器,俗称 “高内聚低耦合”。 在一个审批单的生命周期中会有很多种类的消息发出,比如审批单发起,审批单结束,审批任务生成,审批任务完成,审批抄送等等,我们可以将这些将这些消息的处理逻辑聚合到一个业务处理器中,只需要实现这个接口的 onXxx 方法,就能将逻辑嵌入整个审批的生命周期中。这个业务处理器接口在审批中就是 ProcessEventListener :public interface ProcessEventListener extends EventListener {/*** 审批单发起事件*/void onProcessInstanceStart(InstanceEventContext instanceEventContext);/*** 审批单结束事件*/void onProcessInstanceFinish(InstanceEventContext instanceEventContext);/*** 审批任务生成事件*/void onTaskActivated(TaskEventContext taskEventContext);// 省略其他事件// ...}
// 接到审批单发起消息InstanceEventContext instanceEventContext = buildContext();for (handler : handlers) {try {handler.onProcessInstanceStart(instanceEventContext);} catch (Exception e) {// 打印监控日志等等// ...}}
现在虽然逻辑内聚了,链路更加健壮了,但是还有很多技术上的问题没有解决:
-
某个处理器因为网络超时失败了,如何重试?:我们仅仅是将逻辑拆开了,执行的时候还是一串 “大泥球”,如果仅仅依靠消息队列本身机制,要重试只能一起重试,这显然无法满足诉求;
- 如何高性能地构建庞大的统一上下文(即 Context 参数):为了满足众多处理器对数据查询的诉求,需要提供庞大的上下文,除了性能风险外,也是代码腐败的温床
精准重试
当处理器因为超时等意外情况失败时,如果业务重要的话,都需要补偿重试。但是不同处理器对重试的诉求可能不同,为此我们在上一节朴素想法的基础上,还实现了 多种重试策略 ,并且支持 单独指定最大重试次数 :-
Ignore:非重要业务,失败就算了,不需要重试;
-
Concurrent:在另外线程池中并发执行,也不会重试。适合一些容易影响后续执行的长耗时的处理器;
-
Retry Now:立即重试。会将任务放到本地的一个延迟队列中,100~500ms 后重试。适合时效性比较强的处理器;
-
Retry Later:重投消息,精确重试失败的处理器,遵从 Metaq 的重投延迟,前三次重试分别是
1s 5s 10s
,因此不适合时效性强的处理器;
为了方便使用,我们将这些策略做成了注解的形式,比如审批抄送时效性没这么强,可以使用 Retry Later 策略:
再比如审批待办同步,是时效性比较强的,间隔太久容易有时效性问题,可以使用 Retry Now 策略:public class CcListener implements ProcessEventListener {// Retry Later 策略, 最多重试两次@Policy(value = PolicyType.RETRY_LATER, retry = 2)@Overridepublic void onProcessInstanceStart(InstanceEventContext instanceEventContext) {//... 逻辑省略}}
前三种策略都比较好理解,就不多说了。 下文重点讨论最后一种策略是如何实现的。public class SyncTodoTaskListener implements ProcessEventListener {// Retry Now 策略, 最多重试三次@Policy(value = PolicyType.RETRY_NOW, retry = 3)@Overridepublic void onProcessInstanceStart(InstanceEventContext instanceEventContext) {//... 逻辑省略}}
要想实现精准重试,只需要记下每个处理器的执行状态,比如重试到第几次,是否成功等等。这种流水数据使用关系数据库记会比较重,因此我们 利用 Metaq 提供的 UserProperty ,将各个业务处理器的执行状态以 json 的形式存储在 RETRY_STORE 这个 UserProperty 中,然后重投到另一个专门的重试 topic 中。 处理器的执行状态存储格式如下:
正常执行结束后,如果存在失败的处理器,并且没有达到最大重试次数,会生成上面这种格式的执行报告,存放到 RETRY_STORE 这个 UserProperty 中再重投消息,简化代码如下:{// 总体重试次数, 第一次重试(第 0 次代表正常执行)"globalCnt": 1,// 每个处理器的执行状态"cntMap": {// handler1 第 1 次重试, 读取该属性可以判断 handler1 是否还有重试机会"handler1": 1,// -1 表示 handler2 已经执行成功, 不需要再执行"handler2": -1,// -2 表示 handler3 已经彻底执行失败(一般是超过了设置的最大重试次数), 不需要再执行"handler3": -2}}
虽然我们另起了一个专门重试用的 topic,但是消费者的逻辑跟原 topic 是完全一样的,理论上直接重投原来的 topic 也是可以的,但还是分离开更加安全。Message message = new Message();// 重投到专门的重试主题message.setTopic("my-retry-topic");// 消息体保持不变message.setBody(preBody);// nextCnt 是重试的次数// 设置 DelayTimeLevel 能够让重投有一定的延时message.setDelayTimeLevel(nextCnt);// 将本次执行状态存储到 user property 中message.putUserProperty("RETRY_STORE", "{\"globalCnt\":1,\"cntMap\":{\"handler1\":1,\"handler2\":-1,\"handler3\":-2}}");// 发送消息mqPublishService.send(message);
总体思路如图:
图中的 重投都是指往 topic 再发一遍消息,而不是 Metaq 自身的重投机制。其中有一些细节问题需要注意,比如某次发布中上线了一个新的处理器,而发布机器刚好接到了一个重试消息,如何才能避免新处理器被意外执行呢?我们之前记录的执行状态 cntMap 中的所有 key 就相当于首次执行时的处理器快照,我们重试时只执行 cntMap 中存在的 key 就好了。 另外 Metaq 的 UserProperty 大小也不是没有限制,它最多能存储 30KB 的数据,粗略计算一下大约可以存储 1500 个处理器的状态,这对大多业务来说都是绰绰有余的了。
统一上下文
另一个问题则是由于我们的编码结构导致的,每个处理器都接收相同的上下文参数(XxxContext)处理业务,这种 "上下文 + 处理器" 的结构在业务系统上很常见,但是它存在几个问题:-
为了满足所有处理器的需求,上下文往往会很庞大,因此
构建性能差。
-
外部
无法感知处理器内部需要使用上下文的哪些字段
,只能一股脑地将所有字段都填充好,传递进去,而且内部很有可能一个字段都不使用,白白损耗了性能。
-
上下文中存在一些
幽灵字段
,在某个处理器中设置进去,又在某几个处理器中读取,也就是它有时候为 null,有时候又有值,维护难度巨大,从中取个值都要战战兢兢。
- 读扩散 问题:每个处理器都去读相同的数据,导致链路数十倍的读扩散。
Lazy 框架的具体实现可以参考另一篇文章 利用 惰性写出高性能且抽象的代码 。public class User {// 用户 idprivate Long uid;// 用户的部门,为了保持示例简单,这里就用普通的字符串// 需要远程调用 通讯录系统 获得private final Lazy<String> department;public User(Long uid, Lazy<String> department) {this.uid = uid;this.department = department;}public Long getUid() {return uid;}public String getDepartment() {return department.get();}}// 构建 User 实体Long uid = 1L;User user = new User(uid, Lazy.of(() -> departmentService.getDepartment(uid)));// 使用 User 实体,部门属性用起来和普通属性一样user.getDepartment();
朴素的懒加载思想在依赖服务正常的情况下能够大大减少读扩散和提升性能,但是在依赖服务大规模异常时,还是会造成读扩散,进一步加大依赖服务的压力:
在层次加载结构中,一旦某个节点被求值,路径上所有的属性都会被缓存,因为这种对象能够自动地优化性能。以 利用惰性写出高性能且抽象的代码 这篇文章中 User 实体为例:// 通过用户获得部门Lazy<String> departmentLazy = Lazy.of(() -> departmentService.getDepartment(uid));// 通过部门获得主管// department -> supervisorLazy<Long> supervisorLazy = departmentLazy.map(department -> SupervisorService.getSupervisor(department));
- 审批单实例上下文
- 审批单发起消息
- 审批单结束消息
- 审批单撤销消息
-
...
- 活动上下文
- 活动开始消息
- 活动结束消息
-
...
- 任务上下文
- 任务开始消息
- 任务结束消息
- 任务取消消息
- ...
如何防止雪崩
在消息链路执行这么一大串处理器,如果发生雪崩还是挺危险的。好在本方案不依赖 Metaq 自身的重试机制,直接不处理任何重复/重试的消息,就能够规避无限重试的问题。 根据踩坑的经验,Metaq 发生雪崩的主要原因在于特定场景下的无限重试 :- rebalance 导致重试次数归 0;
-
消费者执行超时,重试次数从 3 开始继续重投;
不处理重试消息:reconsumeTimes 大于 0 直接返回,不处理重试消息。因为我们不依赖 Metaq 层面的重试机制。private boolean isDuplicateMessage(Message msg) {try {Integer consumedCount = ltairManager.incr("dingflow_mq_consume_" + msg.getMsgId() + "_" + msg.getReconsumeTimes(), 1, 30);if (consumedCount != null && consumedCount > 1) {return true;}} catch (Throwable throwable) {// 打印错误日志// ...}return false;}
Show Me The Code: 框架编码设计
框架的整体设计遵循 代码重构:面向单元测试 中的原则,将副作用和核心逻辑完全抽离。框架的主体部分完全与消息队列无关,而是抽象出了一个通用的 RemoteRetryStore ,用于持久化地存储处理器重试状态。这样我们不仅仅可以基于消息队列重试,也可以基于数据库进行重试,只要传递不同的 RemoteRetryStore 即可。框架在本地就可以模拟各种异常情况进行单元测试,文中就不大段罗列代码了 。监控与感知
有了上面的实践,监控与感知方式基本是不言自明的,只需要在框架层面 catch 异常,在有剩余重试机会时打印重试日志,无重试机会时打印失败日志。结合 Sunfire 的 Top 监控,即可看到每个处理器的失败情况。
如图, OpenEventCallbackListener 处理器因为 NPE 问题,失败量突然大幅度增加。
总结
虽然这篇文章写了很多,但是总体思路非常简单,概括起来只有三点:-
将消息链路拆分成多个处理器;
-
利用 Metaq 的 UserProperty 存储每次处理器的执行状态,精准重试失败的处理器;
-
用懒加载机制构建统一上下文,提升构建性能,降低读扩散,最终设计出最符合业务的上下文;