ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

把消息从网络里解放出来:Knit本地总线与Outbox补偿机制实践

把消息从网络里解放出来:Knit本地总线与Outbox补偿机制实践 如果只能从一堆网络报错里总结一条经验我会选择这个原则消息语义不应该和网络传输强耦合。很多业务场景里模块之间只是想传递“订单已创建”“支付已完成”这样的领域事件结果却被做成了必须 TCP 建连、必须经过网关、必须等响应的 HTTP 接口。于是只要出现network: unavailable、stream disconnected before completion、error sending request这类网络故障业务也会一起失败。Knit 这类设计试图把两者分开消息仍然发但不依赖当前是否有网络。它偏好的路径是进程内轻量消息总线也就是“给同一进程里的模块之间发信号”让通信回归到业务本身。这篇文章会用一个最小的 Java 工程把 Knit 的本地消息总线、离线场景下的 Outbox 补偿、以及常见网络报错的排查边界完整梳理一遍。1. 先理解“不使用网络的消息传递”到底解决哪类问题1.1 为什么模块间通信不能全部走网络在实际项目里最容易出现的建模错误是把“一次内部业务通知”当成“一次远程服务调用”。比如订单模块收到支付回调后要通知库存模块、通知积分模块、通知大数据埋点。如果这些模块都在同一个应用进程里却各自暴露 HTTP 接口再通过网关互相调用成本会非常高。一次网络调用至少包含这些环节域名解析或直连 IP。建立连接握手。发起请求等待服务处理。网络传输中断或超时。返回响应或抛出异常。只要其中一个环节不稳定消息就送不到。连接池满了会失败对端重启会失败网络抖动会失败防火墙策略调整也会失败。于是业务代码里会出现大量重试、超时调度、熔断逻辑而这些逻辑和“订单创建后要通知库存”这件事并没有直接关系。Knit 的思路很直接如果调用双方就在同一个 JVM、同一个进程里就不要强行造出一条虚拟网络链路。进程内轻量消息总线既能保留“发消息”的语义又把 DNS、端口、网络超时这些问题完全拿掉。1.2 Knit 的本地事件总线模型Knit 可以理解为一套进程内发布订阅模型。发送方不直接调用接收方的方法而是把一条消息发布到总线上。接收方事先订阅感兴趣的主题总线把消息投递给对应的订阅者。用一个关系来表达发布者 - Knit 总线 - 订阅者 A - 订阅者 B - 订阅者 C这里没有 socket没有 IP没有序列化协议没有“网络是否可达”的概念。消息对象在内存中传递所以把它叫做“Messaging Without the Network”非常合适没有网络消息照样可以送出去。和直接方法调用相比这种模型带来的好处是解耦发布者不需要知道谁在处理消息。订阅者不需要知道消息来自哪个具体组件。新增一个订阅者时不需要改动发布者代码。某个订阅者宕掉时只要异常不冒泡就不会拖垮发布者所在线程。这也是 Knit 常用的落地位置领域事件、应用内部状态广播、页面刷新通知、缓存失效通知等。1.3 Knit 的边界它不能替代消息队列需要特别说清楚Knit 处理的是“单个应用进程内部”的消息通信。它不是 MQ不能替代 Kafka、RocketMQ、RabbitMQ 这类跨进程消息中间件。下面这些场景不该让 Knit 承担多个独立服务实例之间需要同步数据。应用重启后消息也需要保证不丢。消息需要被多个进程消费。消息量巨大需要持久化、分区、削峰填谷。存在安全边界不同信任等级的模块之间需要隔离。如果把 Knit 用在这些场景本质上是把单机内存总线当成了分布式队列重启后消息丢失只是时间问题。所以文章后面也会补充一个更可靠的“离线重发”方案让 Knit 和本地数据库配合而不是硬扛跨机器投递。2. 搭建最小 Java 工程把 Knit 核心结构先放好2.1 环境准备和 Maven 依赖下面的示例代码使用 Java 17因为需要record、Map.of这些特性实际项目如果还在用 Java 8也可以改用 Lombok 或普通 POJO。准备环境工具推荐版本作用JDK17 或更高编译运行 Java 代码Maven3.8 或更高管理工程依赖和执行命令IDEIDEA 或 VS Code阅读和调试代码示例工程不需要额外第三方框架pom.xml 中只需要编译插件和运行插件。完整的 Maven 配置文件如下project xmlnshttp://maven.apache.org/POM/4.0.0 xmlns:xsihttp://www.w3.org/2001/XMLSchema-instance xsi:schemaLocationhttp://maven.apache.org/POM/4.0.0 https://maven.apache.org/xsd/maven-4.0.0.xsd modelVersion4.0.0/modelVersion groupIdcom.example/groupId artifactIdknit-local-bus/artifactId version1.0.0-SNAPSHOT/version properties maven.compiler.release17/maven.compiler.release project.build.sourceEncodingUTF-8/project.build.sourceEncoding /properties build plugins plugin groupIdorg.codehaus.mojo/groupId artifactIdexec-maven-plugin/artifactId version3.2.0/version /plugin /plugins /build /project下面所有代码都放在src/main/java/com/example/knit目录下。目录结构如下knit-local-bus/ ├── pom.xml └── src/main/java/com/example/knit/ ├── KnitMessage.java ├── KnitBus.java ├── LocalKnitBus.java └── Application.java2.2 消息对象先约定消息结构再考虑传输任何总线都要有消息载体。网络场景里消息往往是 JSON、Protobuf、XML在 Knit 里消息就是一个普通 Java 对象。为了让消息便于追踪和过滤先定义一个KnitMessage。package com.example.knit; import java.time.Instant; import java.util.Map; import java.util.Objects; import java.util.UUID; public record KnitMessage( String id, String topic, MapString, String headers, Object payload, Instant createdAt) { public KnitMessage { Objects.requireNonNull(id, id must not be null); Objects.requireNonNull(topic, topic must not be null); headers headers null ? Map.of() : Map.copyOf(headers); createdAt createdAt null ? Instant.now() : createdAt; } public static KnitMessage of(String topic, Object payload) { return new KnitMessage( UUID.randomUUID().toString(), topic, Map.of(), payload, Instant.now() ); } }这里有几个细节值得注意id是消息唯一标识即使不走网络也要保留。后续做幂等、日志追踪时都要用它。topic是发布订阅的主题比如order.created、order.payed。headers用来放链路追踪 ID、来源模块、消息版本等信息。payload是真正的业务数据可以是 DTO、Map、业务事件对象。在本地消息总线里payload默认不会像网络传输那样被复制一份。消息在同一个 JVM 内传递时订阅者拿到的是对象引用。所以要尽量避免在订阅者里修改原对象否则后续订阅者可能看到被改过的数据排查起来很难。2.3 总线接口发布、订阅、取消订阅KnitBus只暴露三个核心能力订阅、取消订阅、发布。不需要关心底层是内存 Map 还是 Redis调用方依赖接口即可。package com.example.knit; import java.util.function.Consumer; public interface KnitBus { String subscribe(String topic, ConsumerKnitMessage consumer); void unsubscribe(String topic, String subscriptionId); void publish(String topic, Object payload); }使用ConsumerKnitMessage可以让 Java 8 的 lambda 直接参与订阅。接口命名越简单越好因为调用方只关心“我要订阅什么”和“我要发什么消息”。3. 实现 LocalKnitBus跑通本地订阅发布流程3.1 使用 ConcurrentHashMap 保存主题和订阅者LocalKnitBus是最常用的进程内实现。它用ConcurrentHashMap保存多个主题每个主题下再保存一批订阅者。ConcurrentHashMap能保证并发情况下注册和读取不出现结构损坏。package com.example.knit; import java.util.Map; import java.util.Objects; import java.util.UUID; import java.util.concurrent.ConcurrentHashMap; import java.util.function.Consumer; public class LocalKnitBus implements KnitBus { private final ConcurrentHashMap String, ConcurrentHashMapString, ConsumerKnitMessage subscribers new ConcurrentHashMap(); Override public String subscribe(String topic, ConsumerKnitMessage consumer) { Objects.requireNonNull(topic, topic must not be null); Objects.requireNonNull(consumer, consumer must not be null); String subscriptionId UUID.randomUUID().toString(); subscribers .computeIfAbsent(topic, k - new ConcurrentHashMap()) .put(subscriptionId, consumer); return subscriptionId; } Override public void unsubscribe(String topic, String subscriptionId) { MapString, ConsumerKnitMessage topicSubscribers subscribers.get(topic); if (topicSubscribers ! null) { topicSubscribers.remove(subscriptionId); } } Override public void publish(String topic, Object payload) { KnitMessage message KnitMessage.of(topic, payload); MapString, ConsumerKnitMessage topicSubscribers subscribers.get(topic); if (topicSubscribers null) { return; } for (ConsumerKnitMessage consumer : topicSubscribers.values()) { try { consumer.accept(message); } catch (RuntimeException ex) { System.err.println(knit subscriber error, topic topic , messageId message.id() , error ex.getMessage()); } } } }仔细看publish方法里的异常捕获。订阅者可能因为空指针、数据库查询超时、参数错误等业务问题抛异常。如果不捕获异常会从总线一路向上抛给发布者导致发布者所在的服务也一起失败。这里捕获后至少记录下来让问题暴露在日志里。不过示例代码里只是System.err.println真实项目应该通过统一日志框架记录并保留异常堆栈。订阅者异常也需要有成体系的处理策略不能只打印一行就当作结束。3.2 一个演示应用Application负责创建总线、注册订阅者并发布消息。这里模拟订单创建后订单组件和指标组件同时收到本地事件。package com.example.knit; import java.util.Map; public class Application { public static void main(String[] args) { KnitBus bus new LocalKnitBus(); String orderSubId bus.subscribe(order.created, message - { System.out.println(订单服务收到: message.payload()); }); bus.subscribe(order.created, message - { System.out.println(指标采集收到: message.headers()); }); bus.publish(order.created, Map.of( orderId, 1024, amount, 199 )); bus.unsubscribe(order.created, orderSubId); } }运行命令mvn -q compile exec:java -Dexec.mainClasscom.example.knit.Application输出大致如下订单服务收到:{amount199, orderId1024} 指标采集收到:{}注意Map的打印顺序不固定这不影响业务。第二个订阅者打印的是headers因为我们没有传入 header所以显示为空 Map。如果在真实场景里要传链路信息应该通过KnitMessage里带 headers 的方式构造消息而不是额外调用网络上下文。3.3 关键点本地总线的同步与异步上面实现里publish是同步派发。订阅者收到消息后发布者线程会等待订阅者方法执行完才继续。这样好处是时序清晰坏处是慢订阅者会阻塞发布者。如果某个订阅者要做文件上传、批量查询、生成报表等耗时操作就不该直接放在订阅方法里。常见做法是在订阅者内部把任务提交给独立线程池。真实生产环境不建议使用无界缓存的Executors.newCachedThreadPool而应该使用有界队列和有上限的线程池。import java.util.concurrent.ArrayBlockingQueue; import java.util.concurrent.ThreadPoolExecutor; import java.util.concurrent.TimeUnit; ThreadPoolExecutor executor new ThreadPoolExecutor( 2, 8, 30, TimeUnit.SECONDS, new ArrayBlockingQueue(2000), new ThreadPoolExecutor.CallerRunsPolicy() );CallerRunsPolicy的作用是当线程池任务队列满时不让任务悄悄丢弃而是由调用者线程继续执行。对于本地事件总线这比直接DiscardPolicy安全。4. 网络不可用时的典型报错为什么会消失又为什么会新增4.1 网络报错与本地消息报错的本质差异接触过网络客户端的人对下面这些日志并不陌生network: unavailable stream disconnected before completion: transport error: network error reconnecting... waiting for network connection failed: error sending request这些报错通常来自距离应用很远的网络链路。可能是 DNS 解析失败可能是服务器的连接被重置也可能是代理策略拦截。用 Knit 本地总线传递消息时这些报错不会出现在“业务消息投递”这一步因为消息根本没有离开进程。但“不会出现网络报错”不等于“不会出错”。本地消息总线面临的错误更底层层一些可能包括订阅者抛了运行时异常。订阅者处理消息太慢把发布线程拖垮。订阅关系没有注册成功消息直接落空。总线被多个线程并发读写导致某次发布看到的是旧订阅表。应用重启后内存中的订阅者全部丢失。所以排查时要先问一个问题这是在排查“本地投递”还是在排查“网络传输”两者日志和链路完全不同。4.2 用表格对比网络消息和 Knit 本地消息维度网络消息Knit 本地消息通信距离跨进程、跨机器单进程内底层依赖TCP/IP、端口、DNS、网关无典型报错connect timeout、broken pipe、network unavailable无这类网络报错消息格式JSON、Protobuf 等序列化格式Java 对象或 JVM 内存结构是否保证持久化视 MQ 实现而定默认不保证适用场景分布式系统、服务间通信模块间广播、领域事件排查重点链路、网络策略、对端日志注册关系、线程模型、异常处理这张表可以用于方案选型也可以用于排错定位。只要业务消息已经确定不会跨进程就不必优先怀疑网络故障。4.3 “不走网络”仍然要设置消息追踪标识有些团队把消息改成本地总线后遇到问题反而不容易查。原因是以前跨模块调用时每个请求都有 traceId可以通过日志系统串联。本地总线上推送的事件没有 traceId事后无法知道是哪次用户操作触发了这个事件。因此在KnitMessage中预留headers很重要。发布消息前可以把请求链路里的 traceId、userId、业务单号放进去。这样即使没有网络层也能通过日志关键字找到完整调用链。实际编码时不要直接构造一个没有头消息的Message而是提供一个带有 headers 的工厂方法。例如KnitMessage message KnitMessage.withHeaders( order.payed, Map.of(traceId, traceId, userId, userId), orderEvent );这就是把应用层的诊断信息前置而不是等出了问题再去网络日志里翻。5. 如果“没有网络”只是暂时的要用 Outbox 组合本地总线5.1 本地总线不是持久化队列有一种很典型的误解只要把消息从 HTTP 改成 Knit 本地总线离线状态也能继续发消息那就等于具备了离线持久化能力。事实并不是这样。LocalKnitBus的内存 Map 一旦进程退出所有消息都会消失。如果业务要求“用户在没有网络的环境里也能完成操作等网络恢复后后台自动同步”那必须把待发送消息落盘或写入数据库。这里推荐 Outbox 模式。简单来说本地业务表和待发送消息表放在同一个数据库事务里。业务数据修改成功的同时往 Outbox 表写入一条状态为PENDING的记录。后续后台服务读取 Outbox按状态发送到远端。5.2 设计 Outbox 表结构以订单创建事件为例可以建这样一张发送表CREATE TABLE app_outbox ( id BIGINT AUTO_INCREMENT PRIMARY KEY, topic VARCHAR(128) NOT NULL, payload_json TEXT NOT NULL, trace_id VARCHAR(64) NULL, status VARCHAR(16) NOT NULL DEFAULT PENDING, retry_count INT NOT NULL DEFAULT 0, next_retry_at DATETIME NOT NULL, created_at DATETIME NOT NULL, updated_at DATETIME NOT NULL, KEY idx_status_next_retry (status, next_retry_at) ) ENGINE InnoDB DEFAULT CHARSET utf8mb4;字段含义字段说明topic远端接口对应的事件主题payload_json需要发送的业务数据保存成 JSONtrace_id本地线程中的链路 ID便于远端日志关联statusPENDING、SENDING、SENT、DEAD 四种状态retry_count已经重试的次数next_retry_at下一次允许发送的时间created_at写入时间状态流转如下PENDING - SENDING - SENT \ | \ v - DEAD当网络持续不可用时记录一直停留在 PENDING。等到后台轮询任务看到next_retry_at已经到达再尝试发送。5.3 业务写入代码示意发送方代码如下这里只说明流程具体实现需要结合 Spring 的事务管理或本地事务public void onOrderPaid(OrderPaidEvent event) { // 1. 业务数据写入订单表 orderRepository.updatePaid(event.getOrderId()); // 2. 同一事务里写 outbox防止业务成功但消息漏写 OutboxRecord record OutboxRecord.pending( order.payed, objectMapper.writeValueAsString(event), TraceContext.getTraceId() ); outboxRepository.insert(record); // 3. 本地广播只用于刷新进程内缓存和 UI不作为远端发送保证 knitBus.publish(order.payed.local, event); }这里的关键是knitBus.publish只负责本地通知真正需要发给远端的消息由 Outbox 表保存。业务数据如果发起成功提交本地总线广播失败也只是显示问题不会影响最终一致性。5.4 网络恢复后的重发逻辑后台发送线程可以这样设计查询status PENDING AND next_retry_at now()的数据逐条调用远端接口。调用成功把状态更新为SENT调用失败则增加retry_count并把next_retry_at向后推到新的时间点。ListOutboxRecord records outboxRepository.findPending(now, limit); for (OutboxRecord record : records) { try { remoteClient.post(record.getTopic(), record.getPayloadJson()); outboxRepository.updateStatus(record.getId(), SENT); } catch (IOException ex) { int nextCount record.getRetryCount() 1; LocalDateTime nextRetry now.plusSeconds(backoffSeconds(nextCount)); outboxRepository.markRetry(record.getId(), nextCount, nextRetry); } }重试间隔不要写死推荐指数退避第 1 次失败后等 30 秒。第 2 次失败后等 60 秒。第 3 次失败后等 120 秒。达到最大次数后进入DEAD状态人工或补偿任务处理。这种设计并不违反“Knit 不依赖网络”。Knit 解决的是网络不可用时进程内业务消息仍能传递Outbox 解决的是网络恢复后远端消息最终能被送达。两者职责不同但可以配合使用。6. 本地消息和网络日志结合时按什么顺序排查6.1 订阅者没有触发的排查链路现象很直接publish执行了但某个订阅者的日志没有出现。很多人会先怀疑是不是网络问题实际上本地总线不存在网络问题先从代码执行路径入手。推荐排查顺序确认publish和subscribe使用的是同一个LocalKnitBus实例。如果测试里 new 了两次就会出现订阅者在 A 对象、发布者往 B 对象发消息永远收不到。确认主题字符串完全一致。多了空格、大小写不一致都会导致匹配失败。确认subscribe先于publish执行。如果先发布后订阅内存总线里还没有订阅者这条消息会直接被丢弃。确认订阅方法没有因为异常被总线捕获后只打印日志而你没有看日志。确认unsubscribe没有被提前调用。常见于页面生命周期结束时清理了订阅关系但后续还在发消息。如果使用了线程池异步派发确认任务队列是否已满任务是否被拒绝。6.2 本地消息场景常见问题速查表现象可能原因检查方式处理建议publish 后没有任何订阅者收到主题不一致或没有订阅打印 topic检查注册日志统一主题常量类避免字符串硬编码程序重启后消息丢失只用内存 Map 存消息观察 Outbox 表是否有数据需要持久化时引入数据库或 MQ订阅者异常导致发布流程中断publish 里没有捕获订阅者异常看异常堆栈是否从总线向上抛在每个订阅者调用外捕获异常消息处理很慢接口响应变慢订阅者同步执行耗时操作用 arthas 或 jstack 看线程堆栈将耗时逻辑放入独立线程池同一事件被处理多次重复订阅或重复扫描 Outbox打印 subscriptionId记录消费幂等键用 messageId 做幂等表6.3 网络错误还存在时先分清是哪一层如果代码里既有 Knit 本地总线也有远程 HTTP 调用那么看到network: unavailable时先要看它发生在哪个组件。比如订单模块先写数据库、再本地广播、再调用清结算接口那么网络错误很可能来自清结算接口而不是 Knit 本地广播。排错时可以按这个顺序分层报错来自哪个类、哪个方法。这个方法里面是否创建了 HTTP Client、gRPC Client 或 Redis 连接。如果没有任何网络客户端那么就要怀疑库里隐含了远程调用。如果报错来自本地总线的“订阅者内部”即使日志里出现网络关键字源头也是订阅者自己调了远程方法不是总线投递问题。不要把网络错误和本地总线错误混在一起。先定位责任在谁再决定是重试网络还是修复订阅注册。7. Knit 本地消息总线的最佳实践与扩展方向7.1 代码层面必须守住的原则本地消息总线看似简单但生产环境里的坑往往在细节。下面这些原则可以当作用前检查清单。第一payload必须尽量不可变。前面提到消息在内存中传递如果订阅者 A 修改了列表内容订阅者 B 再收到时可能拿到脏数据。推荐传业务事件 DTO 或不可变对象不要传可变实体。第二每个订阅者都要独立捕获异常。总线层可以做统一捕获但业务层也要针对自己的失败做好降级。比如订阅者判断数据不完整时可以记录 warning 并跳过而不是把异常交给总线。第三订阅关系必须能取消。长期存在的组件订阅后不清理会积累大量无用订阅者。页面或短生命周期对象订阅时要在一个变量里保存subscriptionId退出时调用unsubscribe避免内存泄漏。第四主题命名要规范。推荐形如order.created、payment.succeeded的小写点分法。如果有版本需求可以写成order.created.v1避免未来事件结构变化后新旧消费逻辑互相干扰。第五不要把本地总线直接暴露给不可信外部输入。Knit 只适合受信任模块之间的消息传递如果外部用户能控制 topic 名称可能发出大量不期望的事件治理成本会非常高。7.2 从本地总线向领域事件演进Knit 的写法其实就是领域事件模式的一个轻量实现。团队多人协作时事件拓扑会越来越复杂建议遵守下面这套规范用Component或Service声明事件发布器。用EventListener或KnitSubscriber声明事件订阅器。事件类名用过去式命名比如OrderCreatedEvent、PaymentSucceededEvent。所有跨模块共享事件放入独立的event模块避免循环依赖。这样后续要引入 Spring 的ApplicationEventPublisher、Guava EventBus、或完整 MQ 时迁移成本会比较低。7.3 如果要进入生产环境还差哪些能力当前这个LocalKnitBus是最小实现。真实项目落地前还需要补齐
返回列表