一、问题从一个小事故说起

有一家做电商结算的小团队,某天中午高峰期,突然接到用户反馈说“下单后订单状态不对”。排查日志发现,订单消息的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 大体相同,但线程池相关配置可能略有变化,请以官方文档为准。