把消息从网络里解放出来:Knit本地总线与Outbox补偿机制实践
2026/9/4 6:18:19 网站建设 项目流程

如果只能从一堆网络报错里总结一条经验,我会选择这个原则:消息语义不应该和网络传输强耦合。很多业务场景里,模块之间只是想传递“订单已创建”“支付已完成”这样的领域事件,结果却被做成了必须 TCP 建连、必须经过网关、必须等响应的 HTTP 接口。于是只要出现network: unavailablestream disconnected before completionerror 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,因为需要recordMap.of这些特性,实际项目如果还在用 Java 8,也可以改用 Lombok 或普通 POJO。

准备环境:

工具推荐版本作用
JDK17 或更高编译运行 Java 代码
Maven3.8 或更高管理工程依赖和执行命令
IDEIDEA 或 VS Code阅读和调试代码

示例工程不需要额外第三方框架,pom.xml 中只需要编译插件和运行插件。完整的 Maven 配置文件如下:

<project xmlns="http://maven.apache.org/POM/4.0.0" xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance" xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 https://maven.apache.org/xsd/maven-4.0.0.xsd"> <modelVersion>4.0.0</modelVersion> <groupId>com.example</groupId> <artifactId>knit-local-bus</artifactId> <version>1.0.0-SNAPSHOT</version> <properties> <maven.compiler.release>17</maven.compiler.release> <project.build.sourceEncoding>UTF-8</project.build.sourceEncoding> </properties> <build> <plugins> <plugin> <groupId>org.codehaus.mojo</groupId> <artifactId>exec-maven-plugin</artifactId> <version>3.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.java

2.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, Map<String, 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.createdorder.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, Consumer<KnitMessage> consumer); void unsubscribe(String topic, String subscriptionId); void publish(String topic, Object payload); }

使用Consumer<KnitMessage>可以让 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, ConcurrentHashMap<String, Consumer<KnitMessage>>> subscribers = new ConcurrentHashMap<>(); @Override public String subscribe(String topic, Consumer<KnitMessage> 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) { Map<String, Consumer<KnitMessage>> topicSubscribers = subscribers.get(topic); if (topicSubscribers != null) { topicSubscribers.remove(subscriptionId); } } @Override public void publish(String topic, Object payload) { KnitMessage message = KnitMessage.of(topic, payload); Map<String, Consumer<KnitMessage>> topicSubscribers = subscribers.get(topic); if (topicSubscribers == null) { return; } for (Consumer<KnitMessage> 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.mainClass=com.example.knit.Application

输出大致如下:

订单服务收到:{amount=199, orderId=1024} 指标采集收到:{}

注意,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需要发送的业务数据,保存成 JSON
trace_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向后推到新的时间点。

List<OutboxRecord> 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执行了,但某个订阅者的日志没有出现。很多人会先怀疑是不是网络问题,实际上本地总线不存在网络问题,先从代码执行路径入手。

推荐排查顺序:

  1. 确认publishsubscribe使用的是同一个LocalKnitBus实例。如果测试里 new 了两次,就会出现订阅者在 A 对象、发布者往 B 对象发,消息永远收不到。
  2. 确认主题字符串完全一致。多了空格、大小写不一致都会导致匹配失败。
  3. 确认subscribe先于publish执行。如果先发布后订阅,内存总线里还没有订阅者,这条消息会直接被丢弃。
  4. 确认订阅方法没有因为异常被总线捕获后只打印日志,而你没有看日志。
  5. 确认unsubscribe没有被提前调用。常见于页面生命周期结束时清理了订阅关系,但后续还在发消息。
  6. 如果使用了线程池异步派发,确认任务队列是否已满,任务是否被拒绝。

6.2 本地消息场景常见问题速查表

现象可能原因检查方式处理建议
publish 后没有任何订阅者收到主题不一致或没有订阅打印 topic,检查注册日志统一主题常量类,避免字符串硬编码
程序重启后消息丢失只用内存 Map 存消息观察 Outbox 表是否有数据需要持久化时引入数据库或 MQ
订阅者异常导致发布流程中断publish 里没有捕获订阅者异常看异常堆栈是否从总线向上抛在每个订阅者调用外捕获异常
消息处理很慢,接口响应变慢订阅者同步执行耗时操作用 arthas 或 jstack 看线程堆栈将耗时逻辑放入独立线程池
同一事件被处理多次重复订阅或重复扫描 Outbox打印 subscriptionId,记录消费幂等键用 messageId 做幂等表

6.3 网络错误还存在时,先分清是哪一层

如果代码里既有 Knit 本地总线,也有远程 HTTP 调用,那么看到network: unavailable时,先要看它发生在哪个组件。比如订单模块先写数据库、再本地广播、再调用清结算接口,那么网络错误很可能来自清结算接口,而不是 Knit 本地广播。

排错时可以按这个顺序分层:

  1. 报错来自哪个类、哪个方法。
  2. 这个方法里面是否创建了 HTTP Client、gRPC Client 或 Redis 连接。
  3. 如果没有任何网络客户端,那么就要怀疑库里隐含了远程调用。
  4. 如果报错来自本地总线的“订阅者内部”,即使日志里出现网络关键字,源头也是订阅者自己调了远程方法,不是总线投递问题。

不要把网络错误和本地总线错误混在一起。先定位责任在谁,再决定是重试网络还是修复订阅注册。

7. Knit 本地消息总线的最佳实践与扩展方向

7.1 代码层面必须守住的原则

本地消息总线看似简单,但生产环境里的坑往往在细节。下面这些原则可以当作用前检查清单。

第一,payload必须尽量不可变。前面提到,消息在内存中传递,如果订阅者 A 修改了列表内容,订阅者 B 再收到时可能拿到脏数据。推荐传业务事件 DTO 或不可变对象,不要传可变实体。

第二,每个订阅者都要独立捕获异常。总线层可以做统一捕获,但业务层也要针对自己的失败做好降级。比如订阅者判断数据不完整时,可以记录 warning 并跳过,而不是把异常交给总线。

第三,订阅关系必须能取消。长期存在的组件订阅后不清理,会积累大量无用订阅者。页面或短生命周期对象订阅时,要在一个变量里保存subscriptionId,退出时调用unsubscribe,避免内存泄漏。

第四,主题命名要规范。推荐形如order.createdpayment.succeeded的小写点分法。如果有版本需求,可以写成order.created.v1,避免未来事件结构变化后新旧消费逻辑互相干扰。

第五,不要把本地总线直接暴露给不可信外部输入。Knit 只适合受信任模块之间的消息传递,如果外部用户能控制 topic 名称,可能发出大量不期望的事件,治理成本会非常高。

7.2 从本地总线向领域事件演进

Knit 的写法其实就是领域事件模式的一个轻量实现。团队多人协作时,事件拓扑会越来越复杂,建议遵守下面这套规范:

  • @Component@Service声明事件发布器。
  • @EventListener@KnitSubscriber声明事件订阅器。
  • 事件类名用过去式命名,比如OrderCreatedEventPaymentSucceededEvent
  • 所有跨模块共享事件放入独立的event模块,避免循环依赖。

这样后续要引入 Spring 的ApplicationEventPublisher、Guava EventBus、或完整 MQ 时,迁移成本会比较低。

7.3 如果要进入生产环境,还差哪些能力

当前这个LocalKnitBus是最小实现。真实项目落地前,还需要补齐:

    需要专业的网站建设服务?

    联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

    立即咨询