一、问题从一个小事故说起
有一家做电商结算的小团队,某天中午高峰期,突然接到用户反馈说“下单后订单状态不对”。排查日志发现,订单消息的body内容张冠李戴:本来应该发“订单A已支付”,结果消息里写的是“订单B已支付”。更诡异的是,生产端异步发送的send回调全部返回了成功。也就是说,代码层面“一切正常”,但下游拿到的数据已经乱了。
这种事儿,我们叫作“伪成功”。它比直接报错更可怕,因为系统不会报警,问题要等到业务侧对账才能发现。提到异步发送,很多开发者觉得“反正调了send,然后回调里打印个成功就完事”,实际上一坑接一坑。今天拿一个真实的案例来拆解:线程池饱和和消息对象复用,这俩隐形问题是怎么凑到一起把消息搞丢、搞错的。
二、别踩双坑:线程池饱和 + 消息对象复用
2.1 线程池饱和是怎么发生的
生产端为了不阻塞主流程,经常会把消息发送的动作丢给一个线程池去执行。线程池相当于一个中转站:如果线程有空,就立即干活;如果线程都忙,任务就到队列里排队;如果队列也满了,那就只能选择“拒绝”。RocketMQ 的异步发送本身也有内部线程池来处理网络事件,但业务侧往往还会再包一层线程池,用来做限流、削峰。
问题在于,很多人对线程池的理解停留在“队列满了就抛异常”,但异常一出现,业务代码就把它吞掉了。于是,消息发送任务连 RocketMQ 的 send 方法都没进去,直接就被丢弃了。上层代码呢?它只知道自己“提交了任务”,没有收到任何失败通知,自然认为消息已经进了消息队列。这就是伪成功的起点。
2.2 消息对象复用为什么致命
再看第二个坑。RocketMQ 发送的消息是一个 Message 对象,里面带有 body、topic、tags、keys 等属性。有些同学喜欢把 Message 当成一个公共的“快递盒”,在循环里只 new 一次,然后不停地往里塞不同的内容再 send。看起来省了内存,其实埋了一个大雷。
因为异步发送是“发完就返回”,真正把数据写到网络,是后台线程稍后才做的事。如果这时候你修改了同一个 Message 对象的body或者keys,后台线程可能读到的已经是修改后的值了。更常见的是在 SendCallback 回调里直接读这个共享的 Message,结果十条消息的回调都打印了最后一条消息的内容。这不仅仅是日志错乱,如果 RocketMQ 底层在发送时直接引用了这个 byte[] 数组,那送往 Broker 的数据也会被覆盖。下游收到一堆重复或错乱的消息,而上游还以为是“发送成功”。
三、代码不过关,事故必上门
为了把上面的问题演示清楚,我直接用 Java + RocketMQ 4.8.0 写一段错误代码。技术栈:Java + RocketMQ 4.8.0,JDK 8。
3.1 错误示例:线程池一满,消息就“假装成功”
下面的代码把一个局部变量 message 当成了全局公共品,同时把发送任务交给一个容量很小的线程池。线程池默认的拒绝策略是 AbortPolicy,也就是队列满了就会直接抛异常。但我们又在 catch 里吞掉了异常,最后看起来就像“没发生任何事”。
import org.apache.rocketmq.client.producer.DefaultMQProducer;
import org.apache.rocketmq.client.producer.SendCallback;
import org.apache.rocketmq.client.producer.SendResult;
import org.apache.rocketmq.common.message.Message;
import java.util.concurrent.ArrayBlockingQueue;
import java.util.concurrent.ThreadPoolExecutor;
import java.util.concurrent.TimeUnit;
public class AsyncSendErrorDemo {
public static void main(String[] args) throws Exception {
// 初始化生产者,生产环境建议从配置中心读取
DefaultMQProducer producer = new DefaultMQProducer("error_demo_group");
producer.setNamesrvAddr("127.0.0.1:9876");
producer.start();
// 故意搞一个非常小的线程池:1个核心线程,队列长度1
// 这样连续提交5个任务时,后面3个任务会被拒绝
ThreadPoolExecutor sender = new ThreadPoolExecutor(
1,
1,
0L, TimeUnit.MILLISECONDS,
new ArrayBlockingQueue<>(1)
// 注意:没有设置拒绝策略,默认是AbortPolicy
);
// 错误点1:Message 对象在循环外只创建了一次
Message message = new Message();
message.setTopic("OrderTopic");
for (int i = 0; i < 5; i++) {
final int index = i;
// 每次循环都去修改同一个 message 的 body 和 keys
message.setBody(("订单数据-" + index).getBytes());
message.setKeys("order_key_" + index);
// 把发送任务交给线程池
sender.execute(() -> {
try {
// 这里执行的是 RocketMQ 异步发送
producer.send(message, new SendCallback() {
@Override
public void onSuccess(SendResult sendResult) {
// 错误点2:回调里读取的 message 是共享的,已经被主线程改掉了
System.out.println("发送成功, key=" + message.getKeys()
+ ", body=" + new String(message.getBody()));
}
@Override
public void onException(Throwable e) {
System.err.println("发送失败, key=" + message.getKeys()
+ ", error=" + e.getMessage());
}
});
} catch (Exception e) {
// 错误点3:把异常吞掉了,上层感知不到任何失败
System.err.println("任务被拒绝或异常: " + e.getMessage());
}
});
}
// 等等异步回调完成,再关闭资源
Thread.sleep(3000);
producer.shutdown();
sender.shutdown();
}
}
运行这段代码,你会发现控制台打印出来的回调成功信息中,key和body大多都是“order_key_4 / 订单数据-4”。而实际上前几条消息在发送时,body可能是“订单数据-1”或“订单数据-2”。更闹心的是,有些任务连 send 都没执行,只留下“任务被拒绝或异常”的一行小字,如果不盯着控制台,这些消息就“蒸发”了。
3.2 完美示例:每次用新对象,回调才靠谱
正确的姿势其实很简单:每条消息都 new 一个独立的 Message 对象,同时给线程池配一个不会丢任务的拒绝策略。在 Java 里,ThreadPoolExecutor.CallerRunsPolicy 可以在任务被拒绝时,交给提交任务的线程自己去执行,这样消息不会被随意丢弃,虽然会带来一定的背压,但业务上至少不会“假装成功”。
import org.apache.rocketmq.client.producer.DefaultMQProducer;
import org.apache.rocketmq.client.producer.SendCallback;
import org.apache.rocketmq.client.producer.SendResult;
import org.apache.rocketmq.common.message.Message;
import java.util.concurrent.ArrayBlockingQueue;
import java.util.concurrent.ThreadPoolExecutor;
import java.util.concurrent.TimeUnit;
public class AsyncSendCorrectDemo {
public static void main(String[] args) throws Exception {
// 生产者和上面一样
DefaultMQProducer producer = new DefaultMQProducer("correct_demo_group");
producer.setNamesrvAddr("127.0.0.1:9876");
producer.start();
// 正确做法:核心线程2,最大线程4,队列长度10
// 拒绝策略使用 CallerRunsPolicy:当线程池和队列都满时,直接在提交线程中执行该任务
ThreadPoolExecutor sender = new ThreadPoolExecutor(
2,
4,
30L, TimeUnit.SECONDS,
new ArrayBlockingQueue<>(10),
new ThreadPoolExecutor.CallerRunsPolicy()
);
for (int i = 0; i < 20; i++) {
final int index = i;
// 正确点1:每次循环都创建全新的 Message 对象
// 这样不同线程之间不会相互干扰
Message message = new Message(
"OrderTopic", // topic
"OrderTag", // tag
"order_key_" + index, // keys
("订单数据-" + index).getBytes() // body
);
sender.execute(() -> {
try {
// 异步发送,回调里使用当前的 message,它是这个任务独有的
producer.send(message, new SendCallback() {
@Override
public void onSuccess(SendResult sendResult) {
// 这里打印和业务关联的key,绝对不会串
System.out.println("发送成功, msgId=" + sendResult.getMsgId()
+ ", key=" + message.getKeys()
+ ", body=" + new String(message.getBody()));
}
@Override
public void onException(Throwable e) {
// 正确点2:回调失败必须处理,最好记录并投递到补偿队列
System.err.println("发送失败, key=" + message.getKeys()
+ ", error=" + e.getMessage());
}
});
} catch (Exception e) {
// 正确点3:不要吞异常,这里可以走重试或者本地记录表
// 但因为有了 CallerRunsPolicy,几乎不会走到这里
System.err.println("发送任务异常: " + e.getMessage());
}
});
}
Thread.sleep(5000);
producer.shutdown();
sender.shutdown();
}
}
这段代码虽然只改动了几处,但效果完全不同:每个线程拿到的 message 都是自己独有的,回调里打印的数据不会串号。同时,线程池也不会轻易丢弃任务。万一真的发生瞬时流量极高,CallerRunsPolicy 会让主线程亲自执行发送,相当于把压力往回传,让调用方感知到的速度变慢,从而形成一种自然的限流。
四、如何避免伪成功状态
4.1 从源头切断复用
最重要的一条:永远不要复用 Message 对象。如果你在循环里不断 new Message,虽然会有一些对象开销,但对于业务系统来说可忽略不计。如果你实在想复用,至少要使用 message = new Message() 重新创建,而不是调用 setBody 去覆盖旧值。为了让代码更稳,建议用构造器一次性把 topic、tags、keys、body 都塞进去,保证不可变(RocketMQ 的 Message 本身可变,但我们可以强迫自己不修改)。
4.2 给线程池一个体面的拒绝策略
线程池的拒绝策略决定了任务太多时会发生什么。默认的 AbortPolicy 会抛异常,而业务代码经常把它吞掉;DiscardPolicy 更危险,直接静默丢弃;DiscardOldestPolicy 会丢弃最老的任务,也不适合需要可靠消息的场景。所以建议用 CallerRunsPolicy,让任务在提交方线程中执行,虽然会降低吞吐,但至少不会造成消息静默丢失。如果你要求更高的可靠性,可以放弃自定义线程池,直接使用 RocketMQ 的默认线程池,并且合理设置异步发送使用的信号量。
4.3 回调里不要依赖共享变量
SendCallback 中的 onSuccess 和 onException 是在 RocketMQ 的异步线程里执行的,和你的业务线程完全不是一个线程。如果你在回调里访问一个外部可变的变量,根本不能确定它此时的值是什么。正确的做法是:把当前消息的关键信息作为局部变量传入回调,或者直接使用构造时固定的 message 对象。注意,即使 message 是独立的,也要确保外部没有人再去修改它的引用。
4.4 配合重试和监控,让“伪成功”无处藏身
其实,异步发送从机制上就注定比同步发送难排查。同步发送失败了能立刻 catch 到,异步发送只能靠回调。所以生产环境最好做到三点:
第一,把发送结果指标化,比如统计 onSuccess 次数、onException 次数、回调耗时,接入 Prometheus 之类监控系统。第二,消息发送失败时要有一个补偿通道,比如把失败消息写进一张本地表,定时重试。第三,在关键业务上,即使 onSuccess 也要做下游数据比对,比如“订单消息”是否真的落库,一旦发现缺口立刻告警。
另外,RocketMQ 的 DefaultMQProducer 里有参数 setRetryTimesWhenSendFailed,但这个参数对异步发送默认不生效。异步发送需要自己在回调里做重试,比如在 onException 中重新调用 send。但要注意,重试时也要重新创建 Message 对象,否则会踩到复用的坑。
五、还想再聊聊RocketMQ异步发送的那些注意点
异步发送的优点很明显:不阻塞主线程,适合高吞吐场景。比如订单支付后的通知、积分变动消息、埋点日志等。但它也有明显的缺点:错误发现的延迟变高,链路变长。很多人以为调用 send 之后消息就发出去了,其实 send 只是把数据交给了网络层。真正是否发到 Broker,要看回调。
还有一个容易被忽略的点:消息发送超时。异步发送同样受 sendMsgTimeout 参数控制,默认是 3000 毫秒。当网络抖动或者 Broker 处理慢时,可能会触发超时。超时后回调会收到异常,但如果你在回调里不区分异常类型,可能误判。建议在 onException 里打印异常堆栈,并区分是超时、网络异常还是 Broker 返回错误。
另外,消息对象复用的问题不仅仅存在于发送前,还可能存在于消费端。比如 Consumer 收到消息后,如果直接持有 Message 引用并在线程池中处理,而 Consumer 内部的消息对象可能被复用?实际上 RocketMQ 的消费端每条消息都会创建新的 Message 对象,不会复用,但为了安全,也不要在消费线程池里修改消息属性。
还有一点:如果你的生产者需要发送很多消息,不建议自己再包一层无界队列的线程池,因为那会导致内存无限增长。建议使用有界队列,配合合理的拒绝策略,让流量压力显性地暴露出来,而不是默默堆在内存里最后 OOM。
六、总结
回到我们开头那个“订单状态不对”的案例。它表面上是一个数据错乱问题,背后其实是两个隐形 bug 叠加的结果。一个是线程池饱和后任务被丢弃,另一个是 Message 对象被循环复用。这两个 bug 单独出现时,至少会有异常或日志错乱;但叠加起来,异常被吞,日志又显示成功,就形成了完美的“伪成功状态”。
要避免伪成功,核心就是三句话:每发一条消息,就创建新的 Message;线程池拒绝策略选择 CallerRunsPolicy;回调里永远不要读取共享变量。做到这三点,大概率可以避开这次“事故”。再往深一层,生产系统不能只靠代码逻辑正确,还要有监控和补偿机制。毕竟异步世界里,没有任何一个回调能够保证业务最终成功,只有在监控指标正常、失败补偿到位的情况下,伪成功才会真正无处可藏。
最后注意,所有示例均使用 Java + RocketMQ 4.8.0,如果你用的是 5.x,API 大体相同,但线程池相关配置可能略有变化,请以官方文档为准。
评论
围绕“RocketMQ生产端开启异步发送时回调丢失消息案例,线程池饱和与消息对象复用引发的隐性问题,如何避免伪成功状态?”参与讨论