rockermq的java客户端的消息的producer的详解

RocketMQ Producer详解

RocketMQ Producer是RocketMQ消息系统中负责发送消息的组件。它支持多种发送方式,可以满足不同场景下的消息发送需求。

1. Producer基本概念

Producer是RocketMQ中消息的发送方,主要功能包括:

  • 创建和发送消息到指定Topic
  • 支持同步、异步和单向发送方式
  • 支持事务消息发送
  • 负载均衡和故障转移

2. Producer核心组件

2.1 DefaultMQProducer

这是RocketMQ中最常用的Producer实现类,提供了丰富的功能和配置选项。

2.2 Producer Group

生产者组是一组具有相同角色的Producer实例,主要用于事务消息和负载均衡。

3. Producer使用示例

基于RocketMQ 5.x版本的rocketmq-client-java客户端,以下是Producer的使用示例:

package com.serone;

import org.apache.rocketmq.client.apis.ClientConfiguration;
import org.apache.rocketmq.client.apis.ClientConfigurationBuilder;
import org.apache.rocketmq.client.apis.ClientException;
import org.apache.rocketmq.client.apis.ClientServiceProvider;
import org.apache.rocketmq.client.apis.message.Message;
import org.apache.rocketmq.client.apis.producer.Producer;
import org.apache.rocketmq.client.apis.producer.SendReceipt;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;

import java.time.Duration;
import java.util.concurrent.CompletableFuture;

public class RocketMQProducerExample {
    private static final Logger logger = LoggerFactory.getLogger(RocketMQProducerExample.class);

    public static void main(String[] args) {
        // 接入点地址,需要设置成Proxy的地址和端口列表,一般是xxx:8080;xxx:8081
        String endpoint = "localhost:8081";
        // 消息发送的目标Topic名称,需要提前创建。
        String topic = "TestTopic";

        // 创建客户端配置
        ClientServiceProvider provider = ClientServiceProvider.loadService();
        ClientConfigurationBuilder builder = ClientConfiguration.newBuilder().setEndpoints(endpoint);
        ClientConfiguration configuration = builder.build();

        try {
            // 初始化Producer时需要设置通信配置以及预绑定的Topic。
            Producer producer = provider.newProducerBuilder()
                    .setTopics(topic)
                    .setClientConfiguration(configuration)
                    .build();

            // 普通消息发送。
            Message message = provider.newMessageBuilder()
                    .setTopic(topic)
                    // 设置消息索引键,可根据关键字精确查找某条消息。
                    .setKeys("messageKey")
                    // 设置消息Tag,用于消费端根据指定Tag过滤消息。
                    .setTag("messageTag")
                    // 消息体。
                    .setBody("messageBody".getBytes())
                    .build();

            try {
                // 同步发送消息,需要关注发送结果,并捕获失败等异常。
                SendReceipt sendReceipt = producer.send(message);
                logger.info("Send message successfully, messageId={}", sendReceipt.getMessageId());
            } catch (ClientException e) {
                logger.error("Failed to send message", e);
            }

            // 异步发送消息
            CompletableFuture<SendReceipt> future = producer.sendAsync(message);
            future.whenComplete((sendReceipt, throwable) -> {
                if (throwable != null) {
                    logger.error("Failed to send message asynchronously", throwable);
                } else {
                    logger.info("Send message asynchronously successfully, messageId={}", sendReceipt.getMessageId());
                }
            });

            // 等待异步发送完成
            try {
                future.get();
            } catch (Exception e) {
                logger.error("Failed to get send result", e);
            }

            // 关闭生产者
            producer.close();
        } catch (ClientException e) {
            logger.error("Failed to create producer", e);
        }
    }
}

4. 不同版本的Producer使用方式

4.1 rocketmq-client-java (5.x版本)

适用于RocketMQ 5.x版本,使用新的API设计:

package com.serone;

import org.apache.rocketmq.client.apis.ClientException;
import org.apache.rocketmq.client.apis.ClientServiceProvider;
import org.apache.rocketmq.client.apis.message.Message;
import org.apache.rocketmq.client.apis.producer.Producer;
import org.apache.rocketmq.client.apis.producer.SendReceipt;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;

public class SimpleProducer {
    private static final Logger logger = LoggerFactory.getLogger(SimpleProducer.class);

    public static void main(String[] args) {
        // 创建默认的生产者
        String topic = "TestTopic";
        String endpoint = "localhost:8081";
        
        try {
            ClientServiceProvider provider = ClientServiceProvider.loadService();
            Producer producer = provider.newProducerBuilder()
                    .setClientConfiguration(
                            ClientConfiguration.newBuilder()
                                    .setEndpoints(endpoint)
                                    .build()
                    )
                    .setTopics(topic)
                    .build();
            
            // 创建消息
            Message message = provider.newMessageBuilder()
                    .setTopic(topic)
                    .setTag("TagA")
                    .setKeys("Key-123")
                    .setBody("Hello RocketMQ".getBytes())
                    .build();
            
            // 发送消息
            SendReceipt sendReceipt = producer.send(message);
            logger.info("Send message successfully, messageId={}", sendReceipt.getMessageId());
            
            // 关闭生产者
            producer.close();
        } catch (ClientException e) {
            logger.error("Failed to send message", e);
        }
    }
}
4.2 rocketmq-client (传统版本)

适用于RocketMQ 4.x版本:

package com.serone;

import org.apache.rocketmq.client.exception.MQClientException;
import org.apache.rocketmq.client.producer.DefaultMQProducer;
import org.apache.rocketmq.client.producer.SendResult;
import org.apache.rocketmq.common.message.Message;

public class TraditionalProducer {
    public static void main(String[] args) throws MQClientException, InterruptedException {
        // 1. 创建生产者实例,指定生产者组名
        DefaultMQProducer producer = new DefaultMQProducer("producer_group");
        
        // 2. 设置NameServer地址
        producer.setNamesrvAddr("localhost:9876");
        
        // 3. 启动生产者
        producer.start();
        
        try {
            // 4. 创建消息(主题、标签、消息体)
            Message message = new Message(
                    "test_topic",        // 主题
                    "test_tag",          // 标签
                    "Hello RocketMQ".getBytes()   // 消息内容
            );
            
            // 5. 发送消息
            SendResult result = producer.send(message);
            System.out.println("发送结果:" + result);
        } catch (Exception e) {
            e.printStackTrace();
        }
        
        // 6. 关闭生产者
        producer.shutdown();
    }
}

5. Producer发送方式

5.1 同步发送

同步发送会等待服务器响应,适用于重要的通知消息:

// 同步发送消息
SendResult result = producer.send(message);
System.out.printf("发送成功 | MsgId=%s | QueueId=%d%n", result.getMsgId(), result.getMessageQueue().getQueueId());
5.2 异步发送

异步发送不会阻塞线程,适用于对响应时间敏感的场景:

// 异步发送消息
producer.send(message, new SendCallback() {
    @Override
    public void onSuccess(SendResult sendResult) {
        System.out.printf("发送成功 | MsgId=%s | QueueId=%d%n", 
            sendResult.getMsgId(), sendResult.getMessageQueue().getQueueId());
    }

    @Override
    public void onException(Throwable e) {
        System.out.printf("发送失败 | %s%n", e);
    }
});
5.3 单向发送

单向发送不关心发送结果,适用于对可靠性要求不高的场景:

// 单向发送消息
producer.sendOneway(message);

6. Producer重要配置参数

  • producerGroup: 生产者组名称
  • namesrvAddr: NameServer地址
  • sendMsgTimeout: 发送消息超时时间
  • retryTimesWhenSendFailed: 发送失败重试次数
  • retryTimesWhenSendAsyncFailed: 异步发送失败重试次数
  • maxMessageSize: 消息最大大小
  • compressMsgBodyOverHowmuch: 消息压缩阈值

7. 使用注意事项

  1. 资源管理: 使用完Producer后必须调用shutdown()方法释放资源
  2. 线程安全: DefaultMQProducer是线程安全的,可以多线程共享使用
  3. 异常处理: 需要妥善处理发送过程中的各种异常
  4. 日志记录: 建议记录发送结果和异常信息,便于问题排查
  5. 性能优化: 根据业务场景选择合适的发送方式和配置参数

这些示例代码展示了如何在您的项目中使用RocketMQ Producer发送消息。您可以根据实际需求选择合适的发送方式和配置参数。

1.1 生产端消息丢失

网络抖动导致消息发送失败
Broker宕机或不可用
消息发送超时
程序异常退出

日志兜底"策略来处理消息生产端的消息丢失问题。

我理解您选择了"日志兜底"策略来处理消息生产端的消息丢失问题。这是一种简单但有效的兜底方式,让我为您详细介绍如何实现和优化日志兜底策略。

日志兜底策略详解

1. 日志兜底的核心思想

日志兜底策略是通过详细记录消息发送过程中的关键信息,当出现消息丢失时,可以通过日志回溯问题并进行人工补偿。

2. 实现方案

2.1 基础日志记录实现
@Service
public class LogBasedMessageProducer {
    
    private static final Logger logger = LoggerFactory.getLogger(LogBasedMessageProducer.class);
    
    @Autowired
    private RocketMQTemplate rocketMQTemplate;
    
    /**
     * 发送消息并记录详细日志
     */
    public boolean sendMessageWithLog(String topic, String tag, String keys, String messageBody) {
        String traceId = UUID.randomUUID().toString();
        long startTime = System.currentTimeMillis();
        
        try {
            // 记录发送前日志
            logger.info("准备发送消息 | traceId={} | topic={} | tag={} | keys={} | body={}", 
                       traceId, topic, tag, keys, messageBody);
            
            // 构造消息
            Message<String> message = MessageBuilder
                .withPayload(messageBody)
                .setHeader(RocketMQHeaders.KEYS, keys)
                .setHeader(RocketMQHeaders.TAGS, tag)
                .build();
            
            // 发送消息
            SendResult result = rocketMQTemplate.syncSend(topic, message, 3000);
            long endTime = System.currentTimeMillis();
            
            // 记录发送成功日志
            logger.info("消息发送成功 | traceId={} | topic={} | msgId={} | status={} | cost={}ms", 
                       traceId, topic, result.getMsgId(), result.getSendStatus(), 
                       (endTime - startTime));
            
            return true;
        } catch (Exception e) {
            long endTime = System.currentTimeMillis();
            // 记录发送失败日志
            logger.error("消息发送失败 | traceId={} | topic={} | tag={} | keys={} | body={} | cost={}ms | error={}", 
                        traceId, topic, tag, keys, messageBody, (endTime - startTime), 
                        e.getMessage(), e);
            
            return false;
        }
    }
}
2.2 结构化日志记录
@Component
public class StructuredMessageLogger {
    
    private static final Logger businessLogger = LoggerFactory.getLogger("business-message");
    
    /**
     * 记录结构化的消息日志
     */
    public void logMessageInfo(MessageLogInfo logInfo) {
        // 使用JSON格式记录,便于后续分析处理
        businessLogger.info("{}", JSON.toJSONString(logInfo));
    }
    
    /**
     * 消息日志信息类
     */
    public static class MessageLogInfo {
        private String traceId;           // 跟踪ID
        private String businessType;      // 业务类型
        private String topic;             // 主题
        private String tag;               // 标签
        private String keys;              // 消息键
        private String messageBody;       // 消息体
        private String sendStatus;        // 发送状态
        private String msgId;             // 消息ID
        private Long sendTime;            // 发送时间
        private Long costTime;            // 耗时
        private String errorMessage;      // 错误信息
        private Map<String, Object> extInfo; // 扩展信息
        
        // 构造函数、getter、setter省略
        // ...
    }
}
2.3 增强版日志兜底实现
@Service
public class EnhancedLogBasedProducer {
    
    private static final Logger msgLogger = LoggerFactory.getLogger("message-trace");
    private static final Logger errorLogger = LoggerFactory.getLogger("message-error");
    
    @Autowired
    private RocketMQTemplate rocketMQTemplate;
    
    @Autowired
    private StructuredMessageLogger structuredLogger;
    
    /**
     * 带完整日志兜底的消息发送
     */
    public SendResult sendMessage(String topic, String tag, String keys, 
                                 String messageBody, Map<String, Object> businessContext) {
        String traceId = generateTraceId();
        long startTime = System.currentTimeMillis();
        
        // 构建日志信息对象
        StructuredMessageLogger.MessageLogInfo logInfo = new StructuredMessageLogger.MessageLogInfo();
        logInfo.setTraceId(traceId);
        logInfo.setTopic(topic);
        logInfo.setTag(tag);
        logInfo.setKeys(keys);
        logInfo.setMessageBody(messageBody);
        logInfo.setSendTime(startTime);
        logInfo.setExtInfo(businessContext);
        
        try {
            // 记录发送前信息
            msgLogger.info("[MSG_SEND_START] traceId={}, topic={}, tag={}, keys={}", 
                          traceId, topic, tag, keys);
            
            // 构造消息
            Message<String> message = MessageBuilder
                .withPayload(messageBody)
                .setHeader(RocketMQHeaders.KEYS, keys)
                .setHeader(RocketMQHeaders.TAGS, tag)
                .build();
            
            // 发送消息
            SendResult result = rocketMQTemplate.syncSend(topic, message, 5000);
            long endTime = System.currentTimeMillis();
            
            // 更新日志信息
            logInfo.setSendStatus("SUCCESS");
            logInfo.setMsgId(result.getMsgId());
            logInfo.setCostTime(endTime - startTime);
            
            // 记录发送成功日志
            msgLogger.info("[MSG_SEND_SUCCESS] traceId={}, topic={}, msgId={}, status={}, cost={}ms", 
                          traceId, topic, result.getMsgId(), result.getSendStatus(), 
                          (endTime - startTime));
            
            // 记录结构化日志
            structuredLogger.logMessageInfo(logInfo);
            
            return result;
        } catch (Exception e) {
            long endTime = System.currentTimeMillis();
            
            // 更新日志信息
            logInfo.setSendStatus("FAILED");
            logInfo.setErrorMessage(e.getMessage());
            logInfo.setCostTime(endTime - startTime);
            
            // 记录发送失败日志
            errorLogger.error("[MSG_SEND_FAILED] traceId={}, topic={}, tag={}, keys={}, error={}", 
                             traceId, topic, tag, keys, e.getMessage(), e);
            
            // 记录结构化日志
            structuredLogger.logMessageInfo(logInfo);
            
            throw new RuntimeException("消息发送失败", e);
        }
    }
    
    private String generateTraceId() {
        return "MSG_" + System.currentTimeMillis() + "_" + UUID.randomUUID().toString().substring(0, 8);
    }
}

3. 日志收集与分析

3.1 日志配置示例
# logback-spring.xml 配置
# 消息跟踪日志
<appender name="MESSAGE_TRACE" class="ch.qos.logback.core.rolling.RollingFileAppender">
    <file>logs/message-trace.log</file>
    <rollingPolicy class="ch.qos.logback.core.rolling.TimeBasedRollingPolicy">
        <fileNamePattern>logs/message-trace.%d{yyyy-MM-dd}.%i.log</fileNamePattern>
        <maxFileSize>100MB</maxFileSize>
        <maxHistory>30</maxHistory>
    </rollingPolicy>
    <encoder>
        <pattern>%d{yyyy-MM-dd HH:mm:ss.SSS} [%thread] %-5level %logger{36} - %msg%n</pattern>
    </encoder>
</appender>

# 消息错误日志
<appender name="MESSAGE_ERROR" class="ch.qos.logback.core.rolling.RollingFileAppender">
    <file>logs/message-error.log</file>
    <rollingPolicy class="ch.qos.logback.core.rolling.TimeBasedRollingPolicy">
        <fileNamePattern>logs/message-error.%d{yyyy-MM-dd}.%i.log</fileNamePattern>
        <maxFileSize>100MB</maxFileSize>
        <maxHistory>90</maxHistory>
    </rollingPolicy>
    <encoder>
        <pattern>%d{yyyy-MM-dd HH:mm:ss.SSS} [%thread] %-5level %logger{36} - %msg%n</pattern>
    </encoder>
    <filter class="ch.qos.logback.classic.filter.LevelFilter">
        <level>ERROR</level>
        <onMatch>ACCEPT</onMatch>
        <onMismatch>DENY</onMismatch>
    </filter>
</appender>

<logger name="message-trace" level="INFO" additivity="false">
    <appender-ref ref="MESSAGE_TRACE"/>
</logger>

<logger name="message-error" level="ERROR" additivity="false">
    <appender-ref ref="MESSAGE_ERROR"/>
</logger>
3.2 日志分析脚本示例
#!/bin/bash
# 消息发送失败统计脚本
# usage: ./message_analysis.sh 2023-10-01

DATE=${1:-$(date +%Y-%m-%d)}
LOG_FILE="logs/message-error.log"
OUTPUT_FILE="message_failures_${DATE}.txt"

echo "=== 消息发送失败统计报告 - ${DATE} ===" > $OUTPUT_FILE
echo "统计时间: $(date)" >> $OUTPUT_FILE
echo "" >> $OUTPUT_FILE

# 统计失败消息数量
FAIL_COUNT=$(grep "${DATE}" $LOG_FILE | grep "\[MSG_SEND_FAILED\]" | wc -l)
echo "失败消息总数: $FAIL_COUNT" >> $OUTPUT_FILE
echo "" >> $OUTPUT_FILE

# 按主题统计失败数量
echo "按主题统计:" >> $OUTPUT_FILE
grep "${DATE}" $LOG_FILE | grep "\[MSG_SEND_FAILED\]" | \
    grep -o "topic=[^,]*" | sort | uniq -c | sort -nr >> $OUTPUT_FILE

echo "" >> $OUTPUT_FILE
echo "最近10条失败记录:" >> $OUTPUT_FILE
grep "${DATE}" $LOG_FILE | grep "\[MSG_SEND_FAILED\]" | tail -10 >> $OUTPUT_FILE

echo "分析完成,结果保存在: $OUTPUT_FILE"

4. 人工补偿流程

4.1 补偿工具实现
@Component
public class MessageCompensationTool {
    
    @Autowired
    private EnhancedLogBasedProducer producer;
    
    /**
     * 根据日志信息重新发送消息
     */
    public boolean resendMessageByLog(String traceId) {
        // 从日志中查找消息信息(实际可通过日志分析系统查询)
        MessageLogInfo logInfo = findMessageLogByTraceId(traceId);
        if (logInfo == null) {
            log.warn("未找到traceId={}的消息日志", traceId);
            return false;
        }
        
        try {
            // 重新发送消息
            SendResult result = producer.sendMessage(
                logInfo.getTopic(),
                logInfo.getTag(),
                logInfo.getKeys(),
                logInfo.getMessageBody(),
                logInfo.getExtInfo()
            );
            
            log.info("消息补偿发送成功 | traceId={} | newMsgId={}", traceId, result.getMsgId());
            return true;
        } catch (Exception e) {
            log.error("消息补偿发送失败 | traceId={}", traceId, e);
            return false;
        }
    }
    
    /**
     * 批量补偿特定时间范围内的失败消息
     */
    public void batchCompensateMessages(Date startTime, Date endTime) {
        // 从日志系统中查询指定时间范围内的失败消息
        List<MessageLogInfo> failedMessages = findFailedMessages(startTime, endTime);
        
        int successCount = 0;
        int failCount = 0;
        
        for (MessageLogInfo logInfo : failedMessages) {
            try {
                SendResult result = producer.sendMessage(
                    logInfo.getTopic(),
                    logInfo.getTag(),
                    logInfo.getKeys(),
                    logInfo.getMessageBody(),
                    logInfo.getExtInfo()
                );
                successCount++;
                log.info("批量补偿成功 | traceId={} | msgId={}", logInfo.getTraceId(), result.getMsgId());
            } catch (Exception e) {
                failCount++;
                log.error("批量补偿失败 | traceId={}", logInfo.getTraceId(), e);
            }
        }
        
        log.info("批量补偿完成 | 总数={} | 成功={} | 失败={}", 
                failedMessages.size(), successCount, failCount);
    }
    
    // 模拟从日志系统查询的方法
    private MessageLogInfo findMessageLogByTraceId(String traceId) {
        // 实际实现应从专门的日志存储中查询
        return null;
    }
    
    private List<MessageLogInfo> findFailedMessages(Date startTime, Date endTime) {
        // 实际实现应从专门的日志存储中查询
        return new ArrayList<>();
    }
}

5. 监控告警

5.1 失败率监控
@Component
public class MessageFailureMonitor {
    
    private static final Logger monitorLogger = LoggerFactory.getLogger("message-monitor");
    
    // 统计时间窗口(分钟)
    private static final int WINDOW_MINUTES = 5;
    
    // 告警阈值
    private static final double FAILURE_RATE_THRESHOLD = 0.05; // 5%失败率
    
    @Scheduled(fixedRate = 60000) // 每分钟检查一次
    public void checkFailureRate() {
        long currentTime = System.currentTimeMillis();
        long windowStartTime = currentTime - WINDOW_MINUTES * 60 * 1000;
        
        // 统计时间窗口内的消息发送情况
        long totalMessages = countMessagesInWindow(new Date(windowStartTime), new Date(currentTime));
        long failedMessages = countFailedMessagesInWindow(new Date(windowStartTime), new Date(currentTime));
        
        if (totalMessages > 0) {
            double failureRate = (double) failedMessages / totalMessages;
            
            if (failureRate > FAILURE_RATE_THRESHOLD) {
                String alertMessage = String.format(
                    "消息发送失败率过高: %.2f%% (失败%d/总计%d), 时间窗口: %d分钟", 
                    failureRate * 100, failedMessages, totalMessages, WINDOW_MINUTES
                );
                
                monitorLogger.warn("[ALERT] {}", alertMessage);
                // 发送告警通知
                sendAlert(alertMessage);
            }
        }
    }
    
    private long countMessagesInWindow(Date startTime, Date endTime) {
        // 实现日志统计逻辑
        return 0;
    }
    
    private long countFailedMessagesInWindow(Date startTime, Date endTime) {
        // 实现失败日志统计逻辑
        return 0;
    }
    
    private void sendAlert(String message) {
        // 实现告警发送逻辑(邮件、短信、钉钉等)
    }
}

6. 最佳实践建议

  1. 日志完整性:确保记录足够的信息用于问题排查和补偿
  2. 日志隔离:将消息日志与业务日志分离,便于管理和分析
  3. 结构化存储:使用JSON等结构化格式记录日志,便于自动化处理
  4. 定期清理:设置合理的日志保存周期,避免占用过多磁盘空间
  5. 监控告警:建立实时监控机制,及时发现异常情况
  6. 补偿工具:开发便捷的补偿工具,提高问题处理效率
  7. 文档记录:建立完整的操作文档,指导人工处理流程

通过完善的日志兜底策略,即使出现消息丢失,也能通过日志快速定位问题并进行有效补偿,保障业务的连续性。

消息producer的发送模式的案例

我来为您详细介绍RocketMQ Producer的各种发送模式,并提供相应的案例代码。

RocketMQ Producer发送模式详解

1. 同步发送模式 (Sync Send)

同步发送是最常用的消息发送方式,发送方会等待Broker返回确认结果。

@Service
public class SyncMessageProducer {
    
    private static final Logger logger = LoggerFactory.getLogger(SyncMessageProducer.class);
    
    @Autowired
    private RocketMQTemplate rocketMQTemplate;
    
    /**
     * 同步发送普通消息
     */
    public void sendSyncMessage(String topic, String messageBody) {
        try {
            // 方式1: 使用RocketMQTemplate直接发送
            SendResult result = rocketMQTemplate.syncSend(topic, messageBody);
            logger.info("同步发送消息成功 | msgId={} | status={}", 
                       result.getMsgId(), result.getSendStatus());
        } catch (Exception e) {
            logger.error("同步发送消息失败 | topic={} | message={}", topic, messageBody, e);
            throw new RuntimeException("消息发送失败", e);
        }
    }
    
    /**
     * 同步发送带标签的消息
     */
    public void sendSyncMessageWithTag(String topic, String tag, String messageBody) {
        try {
            // 构造带标签的消息
            Message<String> message = MessageBuilder
                .withPayload(messageBody)
                .setHeader(RocketMQHeaders.TAGS, tag)
                .build();
            
            SendResult result = rocketMQTemplate.syncSend(topic, message);
            logger.info("同步发送带标签消息成功 | msgId={} | tag={} | status={}", 
                       result.getMsgId(), tag, result.getSendStatus());
        } catch (Exception e) {
            logger.error("同步发送带标签消息失败 | topic={} | tag={} | message={}", 
                        topic, tag, messageBody, e);
            throw new RuntimeException("消息发送失败", e);
        }
    }
    
    /**
     * 同步发送带键的消息(用于查询和去重)
     */
    public void sendSyncMessageWithKeys(String topic, String keys, String messageBody) {
        try {
            Message<String> message = MessageBuilder
                .withPayload(messageBody)
                .setHeader(RocketMQHeaders.KEYS, keys)
                .build();
            
            SendResult result = rocketMQTemplate.syncSend(topic, message);
            logger.info("同步发送带键消息成功 | msgId={} | keys={} | status={}", 
                       result.getMsgId(), keys, result.getSendStatus());
        } catch (Exception e) {
            logger.error("同步发送带键消息失败 | topic={} | keys={} | message={}", 
                        topic, keys, messageBody, e);
            throw new RuntimeException("消息发送失败", e);
        }
    }
    
    /**
     * 同步发送延迟消息
     */
    public void sendSyncDelayMessage(String topic, String messageBody, int delayLevel) {
        try {
            Message<String> message = MessageBuilder
                .withPayload(messageBody)
                .setHeader(RocketMQHeaders.KEYS, "delay_" + System.currentTimeMillis())
                .setHeader("DELAY", delayLevel) // 延迟级别
                .build();
            
            SendResult result = rocketMQTemplate.syncSend(topic, message);
            logger.info("同步发送延迟消息成功 | msgId={} | delayLevel={} | status={}", 
                       result.getMsgId(), delayLevel, result.getSendStatus());
        } catch (Exception e) {
            logger.error("同步发送延迟消息失败 | topic={} | delayLevel={}", topic, delayLevel, e);
            throw new RuntimeException("延迟消息发送失败", e);
        }
    }
}

2. 异步发送模式 (Async Send)

异步发送不会阻塞当前线程,适合对响应时间敏感的场景。

@Service
public class AsyncMessageProducer {
    
    private static final Logger logger = LoggerFactory.getLogger(AsyncMessageProducer.class);
    
    @Autowired
    private RocketMQTemplate rocketMQTemplate;
    
    /**
     * 异步发送消息
     */
    public void sendAsyncMessage(String topic, String messageBody) {
        try {
            // 构造消息
            Message<String> message = MessageBuilder
                .withPayload(messageBody)
                .setHeader(RocketMQHeaders.KEYS, "async_" + System.currentTimeMillis())
                .build();
            
            // 异步发送
            rocketMQTemplate.asyncSend(topic, message, new SendCallback() {
                @Override
                public void onSuccess(SendResult sendResult) {
                    logger.info("异步发送消息成功 | msgId={} | status={}", 
                               sendResult.getMsgId(), sendResult.getSendStatus());
                    // 可以在这里处理发送成功的业务逻辑
                }
                
                @Override
                public void onException(Throwable throwable) {
                    logger.error("异步发送消息失败 | topic={} | message={}", 
                                topic, messageBody, throwable);
                    // 可以在这里处理发送失败的业务逻辑,如重试、告警等
                }
            });
            
            logger.info("异步发送请求已提交 | topic={}", topic);
        } catch (Exception e) {
            logger.error("异步发送消息异常 | topic={} | message={}", topic, messageBody, e);
            throw new RuntimeException("异步消息发送异常", e);
        }
    }
    
    /**
     * 异步发送带超时控制的消息
     */
    public void sendAsyncMessageWithTimeout(String topic, String messageBody, long timeout) {
        try {
            Message<String> message = MessageBuilder
                .withPayload(messageBody)
                .setHeader(RocketMQHeaders.KEYS, "async_timeout_" + System.currentTimeMillis())
                .build();
            
            rocketMQTemplate.asyncSend(topic, message, new SendCallback() {
                @Override
                public void onSuccess(SendResult sendResult) {
                    logger.info("异步发送消息成功 | msgId={} | status={}", 
                               sendResult.getMsgId(), sendResult.getSendStatus());
                }
                
                @Override
                public void onException(Throwable throwable) {
                    logger.error("异步发送消息失败 | topic={} | message={}", 
                                topic, messageBody, throwable);
                }
            }, timeout);
        } catch (Exception e) {
            logger.error("异步发送消息异常 | topic={} | message={}", topic, messageBody, e);
        }
    }
}

3. 单向发送模式 (Oneway Send)

单向发送不等待Broker的响应,适用于对可靠性要求不高的场景。

@Service
public class OnewayMessageProducer {
    
    private static final Logger logger = LoggerFactory.getLogger(OnewayMessageProducer.class);
    
    @Autowired
    private RocketMQTemplate rocketMQTemplate;
    
    /**
     * 单向发送消息
     */
    public void sendOnewayMessage(String topic, String messageBody) {
        try {
            // 构造消息
            Message<String> message = MessageBuilder
                .withPayload(messageBody)
                .setHeader(RocketMQHeaders.KEYS, "oneway_" + System.currentTimeMillis())
                .build();
            
            // 单向发送
            rocketMQTemplate.sendOneWay(topic, message);
            logger.info("单向发送消息完成 | topic={} | message={}", topic, messageBody);
        } catch (Exception e) {
            logger.error("单向发送消息异常 | topic={} | message={}", topic, messageBody, e);
            // 注意:单向发送失败通常不需要重试,因为其本身就是"尽力而为"的发送方式
        }
    }
    
    /**
     * 批量单向发送消息
     */
    public void sendBatchOnewayMessages(String topic, List<String> messageBodies) {
        try {
            List<Message<String>> messages = new ArrayList<>();
            for (String body : messageBodies) {
                Message<String> message = MessageBuilder
                    .withPayload(body)
                    .setHeader(RocketMQHeaders.KEYS, "batch_" + System.currentTimeMillis())
                    .build();
                messages.add(message);
            }
            
            // 批量单向发送
            rocketMQTemplate.sendOneWay(topic, messages);
            logger.info("批量单向发送消息完成 | topic={} | count={}", topic, messages.size());
        } catch (Exception e) {
            logger.error("批量单向发送消息异常 | topic={} | count={}", topic, messageBodies.size(), e);
        }
    }
}

4. 顺序发送模式 (Ordered Send)

顺序消息保证同一队列内的消息按照发送顺序被消费。

@Service
public class OrderedMessageProducer {
    
    private static final Logger logger = LoggerFactory.getLogger(OrderedMessageProducer.class);
    
    @Autowired
    private RocketMQTemplate rocketMQTemplate;
    
    /**
     * 顺序发送消息(基于消息键的哈希选择队列)
     */
    public void sendOrderlyMessage(String topic, String messageKey, String messageBody) {
        try {
            Message<String> message = MessageBuilder
                .withPayload(messageBody)
                .setHeader(RocketMQHeaders.KEYS, messageKey)
                .build();
            
            // 顺序发送,根据messageKey选择队列
            SendResult result = rocketMQTemplate.syncSendOrderly(topic, message, messageKey);
            logger.info("顺序发送消息成功 | msgId={} | key={} | queueId={} | status={}", 
                       result.getMsgId(), messageKey, 
                       result.getMessageQueue().getQueueId(), 
                       result.getSendStatus());
        } catch (Exception e) {
            logger.error("顺序发送消息失败 | topic={} | key={} | message={}", 
                        topic, messageKey, messageBody, e);
            throw new RuntimeException("顺序消息发送失败", e);
        }
    }
    
    /**
     * 批量顺序发送消息
     */
    public void sendBatchOrderlyMessages(String topic, String messageKey, List<String> messageBodies) {
        try {
            List<Message<String>> messages = new ArrayList<>();
            for (String body : messageBodies) {
                Message<String> message = MessageBuilder
                    .withPayload(body)
                    .setHeader(RocketMQHeaders.KEYS, messageKey + "_" + System.currentTimeMillis())
                    .build();
                messages.add(message);
            }
            
            // 批量顺序发送
            SendResult result = rocketMQTemplate.syncSendOrderly(topic, messages, messageKey);
            logger.info("批量顺序发送消息成功 | msgId={} | key={} | messageCount={} | status={}", 
                       result.getMsgId(), messageKey, messages.size(), result.getSendStatus());
        } catch (Exception e) {
            logger.error("批量顺序发送消息失败 | topic={} | key={} | count={}", 
                        topic, messageKey, messageBodies.size(), e);
            throw new RuntimeException("批量顺序消息发送失败", e);
        }
    }
}

5. 事务消息发送模式 (Transactional Send)

事务消息用于实现分布式事务,确保本地事务和消息发送的一致性。

@Service
public class TransactionalMessageProducer {
    
    private static final Logger logger = LoggerFactory.getLogger(TransactionalMessageProducer.class);
    
    @Autowired
    private RocketMQTemplate rocketMQTemplate;
    
    /**
     * 发送事务消息
     */
    @Transactional
    public void sendTransactionMessage(String topic, String messageBody, Order order) {
        try {
            // 1. 先执行本地事务(如保存订单)
            // 这里已经在@Transactional注解的事务中
            
            // 2. 发送半消息(预发送)
            Message<String> message = MessageBuilder
                .withPayload(messageBody)
                .setHeader(RocketMQHeaders.KEYS, "tx_" + order.getId())
                .build();
            
            // 3. 发送事务消息
            SendResult result = rocketMQTemplate.sendMessageInTransaction(
                topic, message, order);
            
            logger.info("事务消息发送成功 | msgId={} | orderId={} | status={}", 
                       result.getMsgId(), order.getId(), result.getSendStatus());
        } catch (Exception e) {
            logger.error("事务消息发送失败 | topic={} | orderId={}", topic, order.getId(), e);
            throw new RuntimeException("事务消息发送失败", e);
        }
    }
    
    /**
     * 事务消息监听器
     */
    @RocketMQTransactionListener
    public class TransactionListenerImpl implements RocketMQLocalTransactionListener {
        
        @Autowired
        private OrderService orderService;
        
        /**
         * 执行本地事务
         */
        @Override
        public RocketMQLocalTransactionState executeLocalTransaction(
                Message msg, Object arg) {
            try {
                Order order = (Order) arg;
                // 执行本地业务逻辑
                orderService.processOrder(order);
                
                // 本地事务执行成功,提交消息
                return RocketMQLocalTransactionState.COMMIT;
            } catch (Exception e) {
                logger.error("本地事务执行失败", e);
                // 本地事务执行失败,回滚消息
                return RocketMQLocalTransactionState.ROLLBACK;
            }
        }
        
        /**
         * 检查本地事务状态(用于事务状态回查)
         */
        @Override
        public RocketMQLocalTransactionState checkLocalTransaction(Message msg) {
            String keys = (String) msg.getHeaders().get(RocketMQHeaders.KEYS);
            // 根据消息key查询本地事务状态
            Order order = orderService.getOrderByMessageKey(keys);
            if (order != null && order.getStatus() == OrderStatus.PROCESSED) {
                return RocketMQLocalTransactionState.COMMIT;
            } else if (order != null && order.getStatus() == OrderStatus.FAILED) {
                return RocketMQLocalTransactionState.ROLLBACK;
            } else {
                // 状态未知,稍后再次检查
                return RocketMQLocalTransactionState.UNKNOWN;
            }
        }
    }
}

6. 批量发送模式 (Batch Send)

批量发送可以提高发送效率,减少网络交互次数。

@Service
public class BatchMessageProducer {
    
    private static final Logger logger = LoggerFactory.getLogger(BatchMessageProducer.class);
    
    @Autowired
    private RocketMQTemplate rocketMQTemplate;
    
    /**
     * 批量发送消息
     */
    public void sendBatchMessages(String topic, List<String> messageBodies) {
        try {
            List<Message<String>> messages = new ArrayList<>();
            for (String body : messageBodies) {
                Message<String> message = MessageBuilder
                    .withPayload(body)
                    .setHeader(RocketMQHeaders.KEYS, "batch_" + System.currentTimeMillis())
                    .build();
                messages.add(message);
            }
            
            // 批量发送
            SendResult result = rocketMQTemplate.syncSend(topic, messages);
            logger.info("批量发送消息成功 | msgId={} | messageCount={} | status={}", 
                       result.getMsgId(), messages.size(), result.getSendStatus());
        } catch (Exception e) {
            logger.error("批量发送消息失败 | topic={} | count={}", topic, messageBodies.size(), e);
            throw new RuntimeException("批量消息发送失败", e);
        }
    }
    
    /**
     * 带标签的批量发送
     */
    public void sendBatchMessagesWithTag(String topic, String tag, List<String> messageBodies) {
        try {
            List<Message<String>> messages = new ArrayList<>();
            for (String body : messageBodies) {
                Message<String> message = MessageBuilder
                    .withPayload(body)
                    .setHeader(RocketMQHeaders.TAGS, tag)
                    .setHeader(RocketMQHeaders.KEYS, "batch_" + System.currentTimeMillis())
                    .build();
                messages.add(message);
            }
            
            SendResult result = rocketMQTemplate.syncSend(topic, messages);
            logger.info("带标签批量发送消息成功 | msgId={} | tag={} | messageCount={} | status={}", 
                       result.getMsgId(), tag, messages.size(), result.getSendStatus());
        } catch (Exception e) {
            logger.error("带标签批量发送消息失败 | topic={} | tag={} | count={}", 
                        topic, tag, messageBodies.size(), e);
            throw new RuntimeException("带标签批量消息发送失败", e);
        }
    }
}

7. 使用示例和场景选择

@RestController
@RequestMapping("/message")
public class MessageController {
    
    @Autowired
    private SyncMessageProducer syncProducer;
    
    @Autowired
    private AsyncMessageProducer asyncProducer;
    
    @Autowired
    private OnewayMessageProducer onewayProducer;
    
    @Autowired
    private OrderedMessageProducer orderedProducer;
    
    @Autowired
    private TransactionalMessageProducer transactionalProducer;
    
    @Autowired
    private BatchMessageProducer batchProducer;
    
    /**
     * 同步发送 - 适用于重要通知、订单创建等场景
     */
    @PostMapping("/sync")
    public ResponseEntity<String> sendSyncMessage(@RequestBody MessageRequest request) {
        syncProducer.sendSyncMessage(request.getTopic(), request.getMessageBody());
        return ResponseEntity.ok("同步消息发送成功");
    }
    
    /**
     * 异步发送 - 适用于日志记录、行为跟踪等场景
     */
    @PostMapping("/async")
    public ResponseEntity<String> sendAsyncMessage(@RequestBody MessageRequest request) {
        asyncProducer.sendAsyncMessage(request.getTopic(), request.getMessageBody());
        return ResponseEntity.ok("异步发送请求已提交");
    }
    
    /**
     * 单向发送 - 适用于大量日志、监控数据等场景
     */
    @PostMapping("/oneway")
    public ResponseEntity<String> sendOnewayMessage(@RequestBody MessageRequest request) {
        onewayProducer.sendOnewayMessage(request.getTopic(), request.getMessageBody());
        return ResponseEntity.ok("单向发送完成");
    }
    
    /**
     * 顺序发送 - 适用于订单状态变更、库存操作等需要顺序处理的场景
     */
    @PostMapping("/orderly")
    public ResponseEntity<String> sendOrderlyMessage(@RequestBody OrderlyMessageRequest request) {
        orderedProducer.sendOrderlyMessage(
            request.getTopic(), 
            request.getMessageKey(), 
            request.getMessageBody()
        );
        return ResponseEntity.ok("顺序消息发送成功");
    }
    
    /**
     * 事务消息 - 适用于分布式事务场景
     */
    @PostMapping("/transaction")
    public ResponseEntity<String> sendTransactionMessage(@RequestBody TransactionMessageRequest request) {
        Order order = new Order();
        order.setId(request.getOrderId());
        // 设置其他订单属性...
        
        transactionalProducer.sendTransactionMessage(
            request.getTopic(), 
            request.getMessageBody(), 
            order
        );
        return ResponseEntity.ok("事务消息发送成功");
    }
    
    /**
     * 批量发送 - 适用于批量数据处理场景
     */
    @PostMapping("/batch")
    public ResponseEntity<String> sendBatchMessages(@RequestBody BatchMessageRequest request) {
        batchProducer.sendBatchMessages(request.getTopic(), request.getMessageBodies());
        return ResponseEntity.ok("批量消息发送成功");
    }
}

8. 发送模式选择指南

发送模式特点适用场景性能可靠性
同步发送等待响应,可靠性高重要业务消息、订单通知等中等高
异步发送不阻塞线程,有回调日志记录、行为跟踪等高中等
单向发送不等待响应,性能最高大量日志、监控数据等最高低
顺序发送保证顺序性订单状态变更、库存操作等中等高
事务发送保证分布式事务一致性支付、转账等金融业务较低最高
批量发送提高发送效率批量数据处理高高

选择合适的发送模式需要综合考虑业务需求、性能要求和可靠性要求。

异步提交消息的线程池

1、系统默认行为
感谢您提供这段关键的源码!这确实是RocketMQ DefaultMQProducerImpl 构造函数中关于异步发送线程池的实现。让我基于这段真实的源码来分析异步发送的线程池配置。

基于源码分析,异步发送线程池的真实配置是:

  • 核心线程数: Runtime.getRuntime().availableProcessors() (CPU核心数)
  • 最大线程数: Runtime.getRuntime().availableProcessors() (CPU核心数)
  • 队列容量: 50000
  • 空闲线程存活时间: 60秒
  • 队列类型: LinkedBlockingQueue
  • 线程命名: AsyncSenderExecutor_1, AsyncSenderExecutor_2, …

RocketMQ异步发送线程池真实配置分析

1. 线程池配置详解

根据您提供的源码,RocketMQ异步发送线程池的真实配置如下:

public class DefaultMQProducerImpl {
    // 异步发送线程池队列,容量50000
    private final LinkedBlockingQueue<Runnable> asyncSenderThreadPoolQueue;
    
    // 默认异步发送线程池
    private final ThreadPoolExecutor defaultAsyncSenderExecutor;
    
    public DefaultMQProducerImpl(final DefaultMQProducer defaultMQProducer, RPCHook rpcHook) {
        this.defaultMQProducer = defaultMQProducer;
        this.rpcHook = rpcHook;

        // 1. 创建异步发送任务队列,容量50000
        this.asyncSenderThreadPoolQueue = new LinkedBlockingQueue<Runnable>(50000);
        
        // 2. 创建异步发送线程池
        this.defaultAsyncSenderExecutor = new ThreadPoolExecutor(
            Runtime.getRuntime().availableProcessors(),  // corePoolSize: CPU核心数
            Runtime.getRuntime().availableProcessors(),  // maximumPoolSize: CPU核心数
            1000 * 60,                                   // keepAliveTime: 60秒
            TimeUnit.MILLISECONDS,                       // 时间单位
            this.asyncSenderThreadPoolQueue,             // 工作队列
            new ThreadFactory() {                        // 线程工厂
                private AtomicInteger threadIndex = new AtomicInteger(0);

                @Override
                public Thread newThread(Runnable r) {
                    return new Thread(r, "AsyncSenderExecutor_" + this.threadIndex.incrementAndGet());
                }
            });
    }
}

2. 真实配置参数总结

基于源码分析,异步发送线程池的真实配置是:

  • 核心线程数: Runtime.getRuntime().availableProcessors() (CPU核心数)
  • 最大线程数: Runtime.getRuntime().availableProcessors() (CPU核心数)
  • 队列容量: 50000
  • 空闲线程存活时间: 60秒
  • 队列类型: LinkedBlockingQueue
  • 线程命名: AsyncSenderExecutor_1, AsyncSenderExecutor_2, …

3. 修正后的异步发送Demo

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 org.apache.rocketmq.remoting.common.RemotingHelper;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;

import java.util.concurrent.ThreadPoolExecutor;
import java.util.concurrent.TimeUnit;

/**
 * 基于真实源码配置的异步发送Demo
 */
public class RealConfigAsyncProducerDemo {
    
    private static final Logger logger = LoggerFactory.getLogger(RealConfigAsyncProducerDemo.class);
    
    public static void main(String[] args) throws Exception {
        // 1. 创建生产者实例
        DefaultMQProducer producer = new DefaultMQProducer("RealConfigAsyncProducerGroup");
        producer.setNamesrvAddr("localhost:9876");
        producer.start();
        
        // 2. 查看实际的线程池配置
        inspectThreadPoolConfig(producer);
        
        // 3. 异步发送消息
        for (int i = 0; i < 10; i++) {
            sendMessageAsync(producer, "TestTopic", "RealAsyncMessage " + i);
        }
        
        // 4. 等待异步回调完成
        Thread.sleep(5000);
        
        // 5. 再次查看线程池状态
        inspectThreadPoolConfig(producer);
        
        producer.shutdown();
    }
    
    /**
     * 查看线程池配置信息
     */
    public static void inspectThreadPoolConfig(DefaultMQProducer producer) {
        try {
            // 通过反射获取内部的线程池(仅用于演示,实际项目中不建议这样做)
            java.lang.reflect.Field field = producer.getClass().getDeclaredField("defaultMQProducerImpl");
            field.setAccessible(true);
            Object producerImpl = field.get(producer);
            
            java.lang.reflect.Field executorField = producerImpl.getClass().getDeclaredField("defaultAsyncSenderExecutor");
            executorField.setAccessible(true);
            ThreadPoolExecutor executor = (ThreadPoolExecutor) executorField.get(producerImpl);
            
            logger.info("线程池状态: 核心线程数={}, 最大线程数={}, 活跃线程数={}, 队列大小={}", 
                       executor.getCorePoolSize(),
                       executor.getMaximumPoolSize(),
                       executor.getActiveCount(),
                       executor.getQueue().size());
                       
        } catch (Exception e) {
            logger.warn("无法获取线程池信息: {}", e.getMessage());
        }
    }
    
    /**
     * 异步发送消息
     */
    public static void sendMessageAsync(DefaultMQProducer producer, String topic, String messageBody) {
        try {
            Message msg = new Message(topic, "RealAsyncTag", "RealKey_" + System.currentTimeMillis(),
                                    messageBody.getBytes(RemotingHelper.DEFAULT_CHARSET));
            
            String mainThread = Thread.currentThread().getName();
            logger.info("主线程 {} 提交异步发送任务", mainThread);
            
            producer.send(msg, new SendCallback() {
                @Override
                public void onSuccess(SendResult sendResult) {
                    String callbackThread = Thread.currentThread().getName();
                    logger.info("异步发送成功 - 回调线程: {} | msgId: {}", callbackThread, sendResult.getMsgId());
                }
                
                @Override
                public void onException(Throwable throwable) {
                    String callbackThread = Thread.currentThread().getName();
                    logger.error("异步发送失败 - 回调线程: {} | error: {}", callbackThread, throwable.getMessage());
                }
            });
            
        } catch (Exception e) {
            logger.error("异步发送提交异常", e);
        }
    }
}

4. 线程池行为演示

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 org.apache.rocketmq.remoting.common.RemotingHelper;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;

import java.util.ArrayList;
import java.util.List;
import java.util.concurrent.CountDownLatch;

/**
 * 演示线程池行为的异步发送Demo
 */
public class ThreadPoolBehaviorDemo {
    
    private static final Logger logger = LoggerFactory.getLogger(ThreadPoolBehaviorDemo.class);
    
    public static void main(String[] args) throws Exception {
        DefaultMQProducer producer = new DefaultMQProducer("ThreadPoolBehaviorProducerGroup");
        producer.setNamesrvAddr("localhost:9876");
        producer.start();
        
        // 模拟高并发异步发送,观察线程池行为
        simulateHighLoadAsyncSend(producer, 100);
        
        Thread.sleep(10000);
        producer.shutdown();
    }
    
    /**
     * 模拟高负载异步发送
     */
    public static void simulateHighLoadAsyncSend(DefaultMQProducer producer, int messageCount) 
            throws InterruptedException {
        CountDownLatch latch = new CountDownLatch(messageCount);
        List<String> callbackThreads = new ArrayList<>();
        
        logger.info("开始高负载异步发送测试,消息数量: {}", messageCount);
        logger.info("CPU核心数: {}", Runtime.getRuntime().availableProcessors());
        
        for (int i = 0; i < messageCount; i++) {
            final int index = i;
            try {
                Message msg = new Message("TestTopic", "HighLoadTag", "HighLoadKey_" + index,
                                        ("High Load Message " + index).getBytes(RemotingHelper.DEFAULT_CHARSET));
                
                producer.send(msg, new SendCallback() {
                    @Override
                    public void onSuccess(SendResult sendResult) {
                        String threadName = Thread.currentThread().getName();
                        synchronized (callbackThreads) {
                            if (!callbackThreads.contains(threadName)) {
                                callbackThreads.add(threadName);
                                logger.info("发现新的回调线程: {}", threadName);
                            }
                        }
                        logger.debug("消息 {} 发送成功: {}", index, sendResult.getMsgId());
                        latch.countDown();
                    }
                    
                    @Override
                    public void onException(Throwable throwable) {
                        logger.error("消息 {} 发送失败: {}", index, throwable.getMessage());
                        latch.countDown();
                    }
                });
                
                // 控制发送速率,避免瞬间提交过多任务
                if (i % 10 == 0) {
                    Thread.sleep(10);
                }
                
            } catch (Exception e) {
                logger.error("发送消息 {} 异常", index, e);
                latch.countDown();
            }
        }
        
        latch.await();
        logger.info("所有消息发送完成,使用的回调线程数量: {}", callbackThreads.size());
        callbackThreads.forEach(thread -> logger.info("回调线程: {}", thread));
    }
}

5. 自定义线程池替代默认配置

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 org.apache.rocketmq.remoting.common.RemotingHelper;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;

import java.util.concurrent.*;

/**
 * 使用自定义线程池的异步发送Demo
 */
public class CustomThreadPoolAsyncProducer {
    
    private static final Logger logger = LoggerFactory.getLogger(CustomThreadPoolAsyncProducer.class);
    
    public static void main(String[] args) throws Exception {
        DefaultMQProducer producer = new DefaultMQProducer("CustomThreadPoolProducerGroup");
        producer.setNamesrvAddr("localhost:9876");
        
        // 设置自定义异步发送线程池
        setCustomAsyncExecutor(producer);
        
        producer.start();
        
        // 发送异步消息
        for (int i = 0; i < 5; i++) {
            sendMessageAsync(producer, "TestTopic", "CustomPoolMessage " + i);
        }
        
        Thread.sleep(3000);
        producer.shutdown();
    }
    
    /**
     * 设置自定义异步发送线程池
     */
    public static void setCustomAsyncExecutor(DefaultMQProducer producer) {
        ThreadPoolExecutor customExecutor = new ThreadPoolExecutor(
            10,                                    // 核心线程数
            50,                                    // 最大线程数
            120,                                   // 空闲线程存活时间(秒)
            TimeUnit.SECONDS,
            new LinkedBlockingQueue<>(100000),     // 大容量队列
            new ThreadFactory() {
                private AtomicInteger threadIndex = new AtomicInteger(0);
                
                @Override
                public Thread newThread(Runnable r) {
                    return new Thread(r, "CustomAsyncSender_" + threadIndex.incrementAndGet());
                }
            },
            new ThreadPoolExecutor.CallerRunsPolicy()  // 调用者运行策略
        );
        
        producer.setAsyncSenderExecutor(customExecutor);
        logger.info("已设置自定义异步发送线程池");
    }
    
    /**
     * 异步发送消息
     */
    public static void sendMessageAsync(DefaultMQProducer producer, String topic, String messageBody) {
        try {
            Message msg = new Message(topic, "CustomTag", "CustomKey_" + System.currentTimeMillis(),
                                    messageBody.getBytes(RemotingHelper.DEFAULT_CHARSET));
            
            producer.send(msg, new SendCallback() {
                @Override
                public void onSuccess(SendResult sendResult) {
                    logger.info("自定义线程池发送成功 - 线程: {} | msgId: {}", 
                               Thread.currentThread().getName(), sendResult.getMsgId());
                }
                
                @Override
                public void onException(Throwable throwable) {
                    logger.error("自定义线程池发送失败 - 线程: {} | error: {}", 
                                Thread.currentThread().getName(), throwable.getMessage());
                }
            });
            
        } catch (Exception e) {
            logger.error("异步发送异常", e);
        }
    }
}

关键结论

根据您提供的源码,RocketMQ异步发送线程池的真实配置是:

  1. 线程数: 固定为CPU核心数(核心线程数=最大线程数)
  2. 队列容量: 50000个任务
  3. 线程存活: 60秒空闲后回收
  4. 队列类型: LinkedBlockingQueue(无界队列,但有容量限制)

这与我在之前回答中提到的配置有所不同,感谢您提供的准确源码信息!这种配置适合大多数场景,但在高并发场景下可能需要自定义线程池来满足性能需求。

Logo

北京人形旗下天工造物具身智能开源社区,聚焦具身天工与慧思开物两大平台

更多推荐