对于消息队列,生产者通常是入门第一个接触的对象,用于生产消息给消费者消费。本文通过介绍生产者实现类的属性、方法,引出生产者的启动过程、高可靠的实现方式等,主要介绍内容如下:
发送消息的一方被称为生产者,它在整个 RocketMQ 的生产和消费体系中扮演的角色如下图所示。

RocketMQ 客户端中的生产者有两个实现类:org.apache.rocketmq.client.producer.DefaultMQProducer和org.apache.rocketmq.client.producer.TransactionMQProducer。前者用户生产普通消息、顺序消息、单向消息、批量消息、延迟消息,后者主要用户生产事务消息。

如下演示了第一个生产者实例:
public class Producer {
public static void main(String[] args) throws Exception {
final DefaultMQProducer producer = new DefaultMQProducer("pay_group");
producer.setNamesrvAddr("xxxx:9876");
producer.setRetryTimesWhenSendAsyncFailed(2);
producer.start();
System.out.println("Producer Starting...");
// Thread.sleep(5000);
final Message message = new Message("xxx", "xxx", "First Demo".getBytes(RemotingHelper.DEFAULT_CHARSET));
final SendResult sendResult = producer.send(message);
System.out.println(sendResult);
producer.shutdown();
}
}

消息类是发送的主体,如下展示了一些核心字段和方法(省略细节代码):
public class Message implements Serializable {
private static final long serialVersionUID = 8445773977080406428L;
private String topic;
private int flag;
private Map<String, String> properties;
private byte[] body;
private String transactionId;
public void setKeys(String keys) {}
public void setKeys(Collection<String> keys) {}
public void setTags(String tags) {}
public void setDelayTimeLevel(int level) {}
public void setTopic(String topic) {}
public void putUserProperty(final String name,final String value) {}
}
MesssageConst.KEY_SPEARATOR分隔或者直接用另一个重载方法。如果 Broker 中 messageIndexEnable=true 则会根据 key 创建消息的 Hash 索引来帮助用户进行快速查询。Rocket MQ 支持普通消息、分区有序消息、全局有序消息、延迟消息和事务消息。
在发送消息的过程中,客户端、Broker、Namesrv 都有可能发生异常,在发生异常时我们依旧需要保证消息的可靠发送。
**第一种保证机制:重试机制。**RocketMQ 支持同步、异步发送消息,不管使用哪种方式都可以配置失败后重试,如果单个 Broker 发生故障,重试会选择其他的 Broker 保证消息的正常发送。
[scode type=“blue”]配置项retryTimesWhenSendFailed表示同步重试次数,默认为2次,加上正常发送1次,总共3次机会。[/scode]
**第二种保证机制:客户端容错。**RocketMQ Client 会维护一个"Broker-发送延迟"关系,根据这个关系选择一个发送延迟级别较低的 Broker 来发送消息,这样能最大限度地利用 Broker 的能力,剔除已经宕机、不可用或者发送延迟级别较高的 Broker,尽量保证消息的正常发送。
单节点的 Broker 无法保证数据的可靠性,生产环境中建议部署2个Master和2个Slave。主从直接的同步方式分为同步复制和异步复制。

public MQClientInstance getAndCreateMQClientInstance(final ClientConfig clientConfig, RPCHook rpcHook) {
String clientId = clientConfig.buildMQClientId();
MQClientInstance instance = this.factoryTable.get(clientId);
if (null == instance) {
// ConcurrentMap factoryTable
instance =
new MQClientInstance(clientConfig.cloneClientConfig(),
this.factoryIndexGenerator.getAndIncrement(), clientId, rpcHook);
MQClientInstance prev = this.factoryTable.putIfAbsent(clientId, instance);
if (prev != null) {
instance = prev;
log.warn("Returned Previous MQClientInstance for clientId:[{}]", clientId);
} else {
log.info("Created new MQClientInstance for clientId:[{}]", clientId);
}
}
return instance;
}
MQClientInstance 实例与 clientId 是一一对应的,而 clientId 是由 clientIP、instanceName 及 unitName 构成的。一般来讲,为了减少客户端的使用资源,如果将所有的 instanceName 和 unitName 设置为相同的值,就只会创建一个 MQClientInstance 实例。生成 clientId 的代码如下所示:
public String buildMQClientId() {
StringBuilder sb = new StringBuilder();
sb.append(this.getClientIP());
sb.append("@");
sb.append(this.getInstanceName());
if (!UtilAll.isBlank(this.unitName)) {
sb.append("@");
sb.append(this.unitName);
}
return sb.toString();
}
MQClientInstance 实例的功能是管理本实例中全部生产者和消费者的生产和消费行为。核心属性如下所示:
public class MQClientInstance {
...
private final String clientId;
private final long bootTimestamp = System.currentTimeMillis();
private final ConcurrentMap<String/* group */, MQProducerInner> producerTable = new ConcurrentHashMap<String, MQProducerInner>();
private final ConcurrentMap<String/* group */, MQConsumerInner> consumerTable = new ConcurrentHashMap<String, MQConsumerInner>();
private final ConcurrentMap<String/* group */, MQAdminExtInner> adminExtTable = new ConcurrentHashMap<String, MQAdminExtInner>();
private final NettyClientConfig nettyClientConfig;
private final MQClientAPIImpl mQClientAPIImpl;
private final MQAdminImpl mQAdminImpl;
private final ConcurrentMap<String/* Topic */, TopicRouteData> topicRouteTable = new ConcurrentHashMap<String, TopicRouteData>();
private final Lock lockNamesrv = new ReentrantLock();
private final Lock lockHeartbeat = new ReentrantLock();
private final ConcurrentMap<String/* Broker Name */, HashMap<Long/* brokerId */, String/* address */>> brokerAddrTable =
new ConcurrentHashMap<String, HashMap<Long, String>>();
private final ConcurrentMap<String/* Broker Name */, HashMap<String/* address */, Integer>> brokerVersionTable =
new ConcurrentHashMap<String, HashMap<String, Integer>>();
private final ScheduledExecutorService scheduledExecutorService = Executors.newSingleThreadScheduledExecutor(new ThreadFactory() {
@Override
public Thread newThread(Runnable r) {
return new Thread(r, "MQClientFactoryScheduledThread");
}
});
private final ClientRemotingProcessor clientRemotingProcessor;
private final PullMessageService pullMessageService;
private final RebalanceService rebalanceService;
private final DefaultMQProducer defaultMQProducer;
private final ConsumerStatsManager consumerStatsManager;
...
}
核心方法如下所示:
public class MQClientInstance {
...
public void updateTopicRouteInfoFromNameServer(){}
private void cleanOfflineBroker() {}
public void checkClientInBroker() throws MQClientException {}
public void sendHeartbeatToAllBrokerWithLock() {}
public boolean updateTopicRouteInfoFromNameServer(final String topic) {}
public boolean updateTopicRouteInfoFromNameServer(final String topic, boolean isDefault, DefaultMQProducer defaultMQProducer) {}
public boolean registerConsumer(final String group,MQConsumerInner consumer) {}
public void unregisterConsumer(final String group) {}
public boolean registerProducer(final String group, DefaultMQProducerImpl producer) {}
public void unregisterProducer(final String group) {}
public boolean registerAdminExt(final String group, MQAdminExtInner admin) {}
public void unregisterAdminExt(final String group) {}
public void rebalanceImmediately() {}
public void doRebalance() {}
public FindBrokerResult findBrokerAddressInAdmin(final String brokerName) {}
public String findBrokerAddressInPublish(final String brokerName) {}
public FindBrokerResult findBrokerAddressInSubscribe(final String brokerName, final long brokerId, final boolean onlyThisBroker) {}
public List<String> findConsumerIdList(final String topic, final String group) {}
public String findBrokerAddrByTopic(final String topic) {}
public void resetOffset(String topic, String group, Map<MessageQueue, Long> offsetTable) {}
public Map<MessageQueue, Long> getConsumerStatus(String topic, String group) {}
public concurrent.ConcurrentMap<String, TopicRouteData> getTopicRouteTable() {}
public ConsumeMessageDirectlyResult consumeMessageDirectly(final MessageExt msg, final String consumerGroup, final String brokerName) {}
public ConsumerRunningInfo consumerRunningInfo(final String consumerGroup) {}
...
}
RocketMQ 客户端的消息发送通常分为以下3层:
业务层:通常指直接调用 RocketMQ Client 发送 API 的业务代码
消息处理层:指 RocketMQ Client 获取业务发送的消息对象后,一系列的参数检查、消息发送准备、参数包装等操作
通信层:指 RocketMQ 基于 Netty 封装的一个 RPC 通信服务,RocketMQ 的各个组件之间的通信全部使用该通信层。

private SendResult sendDefaultImpl(Message msg,final CommunicationMode communicationMode,final SendCallback sendCallback,final long timeout)。communicationMode:通信模式,同步、异步还是单向。sendCallback:对于异步模式,需要设置发送完成后的回调。 sendDefaultImpl方法是发送消息的核心方法,执行过程分为5步:
maxMessageSize 进行设置。tryToFindTopicPublishInfo方法,获取 Topic 路由信息,如果不存在则发出异常提醒用户。如果本地缓存没有路由信息,就通过 Namesrv 获取路由信息,更新到本地,再返回。具体代码如下所示:private TopicPublishInfo tryToFindTopicPublishInfo(final String topic) {
TopicPublishInfo topicPublishInfo = this.topicPublishInfoTable.get(topic);
// topicPublishInfo.ok(): null != this.messageQueueList && !this.messageQueueList.isEmpty()
if (null == topicPublishInfo || !topicPublishInfo.ok()) {
this.topicPublishInfoTable.putIfAbsent(topic, new TopicPublishInfo());
this.mQClientFactory.updateTopicRouteInfoFromNameServer(topic);
topicPublishInfo = this.topicPublishInfoTable.get(topic);
}
if (topicPublishInfo.isHaveTopicRouterInfo() || topicPublishInfo.ok()) {
return topicPublishInfo;
} else {
this.mQClientFactory.updateTopicRouteInfoFromNameServer(topic, true, this.defaultMQProducer);
topicPublishInfo = this.topicPublishInfoTable.get(topic);
return topicPublishInfo;
}
}
timesTotal,同步重试和异步重试的执行方式是不同的。selectOneMessageQueue()。根据队列对象中保存的上次发送消息的 Broker 的名字和 Topic 路由,选择(轮询)一个 Queue 将消息发送到 Broker。我们可以通过 sendLatencyFaultEnable来设置是否总是发送到延迟级别较低的 Broker,默认值为 Flase。sendKernelImpl方法。该方法是发送消息的核心方法,主要用于准备通信层的入参(比如Broker地址、请求体等),将请求传递给通信层,内部实现是基于Netty的,在封装为通信层 request 对象 RemotingCommand 前,会设置 RequestCode 表示当前请求时发送单个消息还是批量消息。同步发送消息是,根据 HashKey 将消息发送到指定的分区中,每个分区中的消息都是按照发送顺序保存的,即分区有序。如果 Topic 的分区被设置为1,这个 Topic 的消息就是全局有序的。注意:顺序消息的发送必须是单线程,多线程将不再有序。顺序消息的消费和普通消息的消费方式不同。
public class OrderMessageProducer {
public static void main(String[] args) throws MQClientException, InterruptedException, UnsupportedEncodingException, RemotingException, MQBrokerException {
final DefaultMQProducer producer = new DefaultMQProducer("pay_group");
producer.setNamesrvAddr("192.168.233.145:9876");
producer.setRetryTimesWhenSendAsyncFailed(2);
producer.start();
System.out.println("Producer Starting...");
// Thread.sleep(5000);
String[] tags = new String[] {"TagA", "TagB", "TagC", "TagD", "TagE"};
for (int i = 0; i < 100; i++) {
int orderId = i % 10;
//Create a message instance, specifying topic, tag and message body.
Message msg = new Message("catwinghu", tags[i % tags.length], "KEY" + i,
("Hello RocketMQ " + i).getBytes(RemotingHelper.DEFAULT_CHARSET));
SendResult sendResult = producer.send(msg, new MessageQueueSelector() {
@Override
public MessageQueue select(List<MessageQueue> mqs, Message msg, Object orderId) {
Integer id = (Integer) orderId;
int index = id % mqs.size();
return mqs.get(index);
}
}, orderId);
System.out.printf("%s%n", sendResult);
}
producer.shutdown();
}
}
如下图展示了发送结果,不难发现orderId%5为0 的order都发送到了0号队列,实现了顺序性。

生产者发送消息后,消费者在指定时间才能消费消息,这类消息被称为延迟消息或定时消息。生产者发送延迟消息前需要设置延迟级别,目前开源版本支持18个延迟级别:1s 5s 10s 30s 1m 2m 3m 4m 5m 6m 7m 8m 9m 10m 20m 30m 1h 2h。
public class DelayMessageProducer {
public static void main(String[] args) throws Exception{
final DefaultMQProducer producer = new DefaultMQProducer("pay_group");
producer.setNamesrvAddr("192.168.233.145:9876");
producer.setRetryTimesWhenSendAsyncFailed(2);
producer.start();
Thread.sleep(5000);
System.out.println("Producer Starting...");
final Message message = new Message("catwinghu", "taga", "DelayMessage".getBytes(RemotingHelper.DEFAULT_CHARSET));
// 1s 5s 10s 30s 1m 2m 3m 4m 5m 6m 7m 8m 9m 10m 20m 30m 1h 2h
// 设置延迟级别
message.setDelayTimeLevel(3);
final SendResult sendResult = producer.send(message);
System.out.println(sendResult);
producer.shutdown();
}
}
Broker 在接收到用户发送的消息后,首先将消息保存到名为SCHEDULE_TOPIC_XXX的 Topic 中。此时消费者无法消费到该延迟消息。然后,由 Broker 端的定时投递任务定时投递给消费者。保存延迟消息的实现逻辑见org.apache.rocketmq.store.schedule.ScheduleMessageService类。按照配置的延迟级别初始化多个任务,每秒执行一次,若消息投递满足时间条件,则将消息投递到原始的 Topic 中。发送(投递)延迟消息定时任务(DeliverDelayedMessageTimerTask)代码如下所示:
class DeliverDelayedMessageTimerTask extends TimerTask {
/**
* 延迟级别
*/
private final int delayLevel;
/**
* 位置
*/
private final long offset;
public DeliverDelayedMessageTimerTask(int delayLevel, long offset) {
this.delayLevel = delayLevel;
this.offset = offset;
}
@Override
public void run() {
try {
this.executeOnTimeup();
} catch (Exception e) {
// XXX: warn and notify me
log.error("ScheduleMessageService, executeOnTimeup exception", e);
ScheduleMessageService.this.timer.schedule(new DeliverDelayedMessageTimerTask(
this.delayLevel, this.offset), DELAY_FOR_A_PERIOD);
}
}
/**
* 纠正可投递时间。
* 因为发送级别对应的发送间隔可以调整,如果超过当前间隔,则修正成当前配置,避免后面的消息无法发送。
*
* @param now 当前时间
* @param deliverTimestamp 投递时间
* @return 纠正结果
*/
private long correctDeliverTimestamp(final long now, final long deliverTimestamp) {
long result = deliverTimestamp;
long maxTimestamp = now + ScheduleMessageService.this.delayLevelTable.get(this.delayLevel);
if (deliverTimestamp > maxTimestamp) {
result = now;
}
return result;
}
public void executeOnTimeup() {
ConsumeQueue cq = ScheduleMessageService.this.defaultMessageStore.findConsumeQueue(SCHEDULE_TOPIC, delayLevel2QueueId(delayLevel));
long failScheduleOffset = offset;
if (cq != null) {
SelectMappedBufferResult bufferCQ = cq.getIndexBuffer(this.offset);
if (bufferCQ != null) {
try {
long nextOffset = offset;
int i = 0;
for (; i < bufferCQ.getSize(); i += ConsumeQueue.CQ_STORE_UNIT_SIZE) {
long offsetPy = bufferCQ.getByteBuffer().getLong();
int sizePy = bufferCQ.getByteBuffer().getInt();
long tagsCode = bufferCQ.getByteBuffer().getLong();
long now = System.currentTimeMillis();
long deliverTimestamp = this.correctDeliverTimestamp(now, tagsCode);
nextOffset = offset + (i / ConsumeQueue.CQ_STORE_UNIT_SIZE);
long countdown = deliverTimestamp - now;
if (countdown <= 0) { // 消息到达可发送时间
MessageExt msgExt = ScheduleMessageService.this.defaultMessageStore.lookMessageByOffset(offsetPy, sizePy);
if (msgExt != null) {
try {
// 发送消息
MessageExtBrokerInner msgInner = this.messageTimeup(msgExt);
PutMessageResult putMessageResult = ScheduleMessageService.this.defaultMessageStore.putMessage(msgInner);
if (putMessageResult != null && putMessageResult.getPutMessageStatus() == PutMessageStatus.PUT_OK) { // 发送成功
continue;
} else { // 发送失败
// XXX: warn and notify me
log.error("ScheduleMessageService, a message time up, but reput it failed, topic: {} msgId {}", msgExt.getTopic(), msgExt.getMsgId());
// 安排下一次任务
ScheduleMessageService.this.timer.schedule(new DeliverDelayedMessageTimerTask(this.delayLevel, nextOffset), DELAY_FOR_A_PERIOD);
// 更新进度
ScheduleMessageService.this.updateOffset(this.delayLevel, nextOffset);
return;
}
} catch (Exception e) {
// XXX: warn and notify me
log.error("ScheduleMessageService, messageTimeup execute error, drop it. msgExt="
+ msgExt + ", nextOffset=" + nextOffset + ",offsetPy=" + offsetPy + ",sizePy=" + sizePy, e);
}
}
} else {
// 安排下一次任务
ScheduleMessageService.this.timer.schedule(new DeliverDelayedMessageTimerTask(this.delayLevel, nextOffset), countdown);
// 更新进度
ScheduleMessageService.this.updateOffset(this.delayLevel, nextOffset);
return;
}
} // end of for
nextOffset = offset + (i / ConsumeQueue.CQ_STORE_UNIT_SIZE);
// 安排下一次任务
ScheduleMessageService.this.timer.schedule(new DeliverDelayedMessageTimerTask(this.delayLevel, nextOffset), DELAY_FOR_A_WHILE);
// 更新进度
ScheduleMessageService.this.updateOffset(this.delayLevel, nextOffset);
return;
} finally {
bufferCQ.release();
}
} // end of if (bufferCQ != null)
else { // 消费队列已经被删除部分,跳转到最小的消费进度
long cqMinOffset = cq.getMinOffsetInQueue();
if (offset < cqMinOffset) {
failScheduleOffset = cqMinOffset;
log.error("schedule CQ offset invalid. offset=" + offset + ", cqMinOffset="
+ cqMinOffset + ", queueId=" + cq.getQueueId());
}
}
} // end of if (cq != null)
ScheduleMessageService.this.timer.schedule(new DeliverDelayedMessageTimerTask(this.delayLevel, failScheduleOffset), DELAY_FOR_A_WHILE);
}
/**
* 设置消息内容
*
* @param msgExt 消息
* @return 消息
*/
private MessageExtBrokerInner messageTimeup(MessageExt msgExt) {
MessageExtBrokerInner msgInner = new MessageExtBrokerInner();
msgInner.setBody(msgExt.getBody());
msgInner.setFlag(msgExt.getFlag());
MessageAccessor.setProperties(msgInner, msgExt.getProperties());
TopicFilterType topicFilterType = MessageExt.parseTopicFilterType(msgInner.getSysFlag());
long tagsCodeValue =
MessageExtBrokerInner.tagsString2tagsCode(topicFilterType, msgInner.getTags());
msgInner.setTagsCode(tagsCodeValue);
msgInner.setPropertiesString(MessageDecoder.messageProperties2String(msgExt.getProperties()));
msgInner.setSysFlag(msgExt.getSysFlag());
msgInner.setBornTimestamp(msgExt.getBornTimestamp());
msgInner.setBornHost(msgExt.getBornHost());
msgInner.setStoreHost(msgExt.getStoreHost());
msgInner.setReconsumeTimes(msgExt.getReconsumeTimes());
msgInner.setWaitStoreMsgOK(false);
MessageAccessor.clearProperty(msgInner, MessageConst.PROPERTY_DELAY_TIME_LEVEL);
msgInner.setTopic(msgInner.getProperty(MessageConst.PROPERTY_REAL_TOPIC));
String queueIdStr = msgInner.getProperty(MessageConst.PROPERTY_REAL_QUEUE_ID);
int queueId = Integer.parseInt(queueIdStr);
msgInner.setQueueId(queueId);
return msgInner;
}
}
单向消息的生产者只管发送过程,不管发送结果。单向消息主要用于日志传输等消息允许丢失的场景。
public class OnewayProducer {
public static void main(String[] args) throws Exception {
final DefaultMQProducer producer = new DefaultMQProducer("pay_group");
producer.setNamesrvAddr("192.168.233.145:9876");
producer.setRetryTimesWhenSendAsyncFailed(2);
producer.start();
System.out.println("Producer Starting...");
Thread.sleep(5000);
final Message message = new Message("catwinghu", "taga", "OneWay Demo".getBytes(RemotingHelper.DEFAULT_CHARSET));
producer.sendOneway(message);
producer.shutdown();
}
}
批量消息发送能提高发送效率,提升系统吞吐量。批量消息发送有以下3点注意事项:
public class BatchMsgProducer {
public static void main(String[] args) throws Exception {
final DefaultMQProducer producer = new DefaultMQProducer("batch_group");
producer.setNamesrvAddr("192.168.233.145:9876");
producer.setRetryTimesWhenSendAsyncFailed(2);
producer.start();
System.out.println("Producer Starting...");
Thread.sleep(5000);
final List<Message> messages = Arrays.asList(new Message("catwinghu", "taga", "Order001".getBytes(RemotingHelper.DEFAULT_CHARSET)),
new Message("catwinghu", "taga", "Order002".getBytes(RemotingHelper.DEFAULT_CHARSET)),
new Message("catwinghu", "taga", "Order003".getBytes(RemotingHelper.DEFAULT_CHARSET)));
final SendResult sendResult = producer.send(messages);
System.out.println(sendResult);
producer.shutdown();
}
}

事务消息的发送、消费流程和延迟消息类似,都是先发送到一个对消费者不可见的 Topic 中。当事务被业务提交后,会被二次投递到原始的 Topic 中,此时消费者正常消费,事务消息的发送具体分为以下两个步骤:

整体交互流程图如下所示:

import org.apache.rocketmq.client.exception.MQClientException;
import org.apache.rocketmq.client.producer.LocalTransactionState;
import org.apache.rocketmq.client.producer.TransactionListener;
import org.apache.rocketmq.client.producer.TransactionMQProducer;
import org.apache.rocketmq.common.message.Message;
import org.apache.rocketmq.common.message.MessageExt;
import org.apache.rocketmq.remoting.common.RemotingHelper;
import java.io.UnsupportedEncodingException;
import java.util.concurrent.*;
public class TransactionProducer {
private String producerGroup = "transaction_group";
//事务监听器
private TransactionListener listener = new TransactionCheckListenerImpl();
private TransactionMQProducer producer = null;
private ExecutorService executorService = new ThreadPoolExecutor(2,5,100, TimeUnit.MILLISECONDS, new ArrayBlockingQueue<Runnable>(2000),new ThreadFactory(){
@Override
public Thread newThread(Runnable r){
final Thread thread = new Thread(r);
thread.setName("client-transaction-msg-check-thread");
return thread;
}
});
public TransactionProducer(){
producer = new TransactionMQProducer(producerGroup);
producer.setNamesrvAddr("192.168.233.145:9876");
producer.setTransactionListener(listener);
producer.setExecutorService(executorService);
start();
}
public void start(){
try {
this.producer.start();
Thread.sleep(5000);
System.out.println("Producer Starting...");
} catch (Exception e) {
e.printStackTrace();
}
}
public void shutdown(){
this.producer.shutdown();
}
public TransactionMQProducer getProducer(){
return producer;
}
public static void main(String[] args) throws MQClientException, UnsupportedEncodingException {
final TransactionProducer transactionProducer = new TransactionProducer();
final TransactionMQProducer producer = transactionProducer.getProducer();
final Message trmsg1 = new Message("catwinghu", "taga","1111", "TransactionMsg1".getBytes(RemotingHelper.DEFAULT_CHARSET));
final Message trmsg2 = new Message("catwinghu", "taga","2222", "TransactionMsg2".getBytes(RemotingHelper.DEFAULT_CHARSET));
final Message trmsg3 = new Message("catwinghu", "taga","3333", "TransactionMsg3".getBytes(RemotingHelper.DEFAULT_CHARSET));
producer.sendMessageInTransaction(trmsg1,1);
producer.sendMessageInTransaction(trmsg2,2);
producer.sendMessageInTransaction(trmsg3,3);
}
}
class TransactionCheckListenerImpl implements TransactionListener {
@Override
public LocalTransactionState executeLocalTransaction(Message msg, Object arg) {
System.out.println("执行本地事务");
final String body = new String(msg.getBody());
final String keys = msg.getKeys();
final String transactionId = msg.getTransactionId();
System.out.printf("transactionId=%s key=%s body=%s\n",transactionId,keys,body);
final int status = Integer.parseInt(arg.toString());
if(status == 1){
System.out.println("提交");
return LocalTransactionState.COMMIT_MESSAGE;
}
if(status == 2){
System.out.println("回滚");
return LocalTransactionState.ROLLBACK_MESSAGE;
}
// UNKNOW 会回查
System.out.println("回查");
return LocalTransactionState.UNKNOW;
}
@Override
public LocalTransactionState checkLocalTransaction(MessageExt msg) {
System.out.println("事务回查");
final String body = new String(msg.getBody());
final String keys = msg.getKeys();
final String transactionId = msg.getTransactionId();
System.out.printf("transactionId=%s key=%s body=%s\n",transactionId,keys,body);
// 只有Commit、Rollback
//可以根据key去检查本地事务消息是否完成
return LocalTransactionState.COMMIT_MESSAGE;
}
}
