消息队列rocketmq
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. 使用注意事项
- 资源管理: 使用完Producer后必须调用shutdown()方法释放资源
- 线程安全: DefaultMQProducer是线程安全的,可以多线程共享使用
- 异常处理: 需要妥善处理发送过程中的各种异常
- 日志记录: 建议记录发送结果和异常信息,便于问题排查
- 性能优化: 根据业务场景选择合适的发送方式和配置参数
这些示例代码展示了如何在您的项目中使用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. 最佳实践建议
- 日志完整性:确保记录足够的信息用于问题排查和补偿
- 日志隔离:将消息日志与业务日志分离,便于管理和分析
- 结构化存储:使用JSON等结构化格式记录日志,便于自动化处理
- 定期清理:设置合理的日志保存周期,避免占用过多磁盘空间
- 监控告警:建立实时监控机制,及时发现异常情况
- 补偿工具:开发便捷的补偿工具,提高问题处理效率
- 文档记录:建立完整的操作文档,指导人工处理流程
通过完善的日志兜底策略,即使出现消息丢失,也能通过日志快速定位问题并进行有效补偿,保障业务的连续性。
消息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异步发送线程池的真实配置是:
- 线程数: 固定为CPU核心数(核心线程数=最大线程数)
- 队列容量: 50000个任务
- 线程存活: 60秒空闲后回收
- 队列类型: LinkedBlockingQueue(无界队列,但有容量限制)
这与我在之前回答中提到的配置有所不同,感谢您提供的准确源码信息!这种配置适合大多数场景,但在高并发场景下可能需要自定义线程池来满足性能需求。
更多推荐
所有评论(0)