6 消息事务样例

事务消息共有三种状态,提交状态、回滚状态、中间状态:

  • TransactionStatus.CommitTransaction: 提交事务,它允许消费者消费此消息。
  • TransactionStatus.RollbackTransaction: 回滚事务,它代表该消息将被删除,不允许被消费。
  • TransactionStatus.Unknown: 中间状态,它代表需要检查消息队列来确定状态。

6.1 发送事务消息样例

1、创建事务性生产者

使用 TransactionMQProducer类创建生产者,并指定唯一的 ProducerGroup,就可以设置自定义线程池来处理这些检查请求。执行本地事务后、需要根据执行结果对消息队列进行回复。回传的事务状态在请参考前一节。

  1. import org.apache.rocketmq.client.consumer.DefaultMQPushConsumer;
  2. import org.apache.rocketmq.client.consumer.listener.ConsumeConcurrentlyContext;
  3. import org.apache.rocketmq.client.consumer.listener.ConsumeConcurrentlyStatus;
  4. import org.apache.rocketmq.client.consumer.listener.MessageListenerConcurrently;
  5. import org.apache.rocketmq.common.message.MessageExt;
  6. import java.util.List;
  7. public class TransactionProducer {
  8. public static void main(String[] args) throws MQClientException, InterruptedException {
  9. TransactionListener transactionListener = new TransactionListenerImpl();
  10. TransactionMQProducer producer = new TransactionMQProducer("please_rename_unique_group_name");
  11. ExecutorService executorService = new ThreadPoolExecutor(2, 5, 100, TimeUnit.SECONDS, new ArrayBlockingQueue<Runnable>(2000), new ThreadFactory() {
  12. @Override
  13. public Thread newThread(Runnable r) {
  14. Thread thread = new Thread(r);
  15. thread.setName("client-transaction-msg-check-thread");
  16. return thread;
  17. }
  18. });
  19. producer.setExecutorService(executorService);
  20. producer.setTransactionListener(transactionListener);
  21. producer.start();
  22. String[] tags = new String[] {"TagA", "TagB", "TagC", "TagD", "TagE"};
  23. for (int i = 0; i < 10; i++) {
  24. try {
  25. Message msg =
  26. new Message("TopicTest1234", tags[i % tags.length], "KEY" + i,
  27. ("Hello RocketMQ " + i).getBytes(RemotingHelper.DEFAULT_CHARSET));
  28. SendResult sendResult = producer.sendMessageInTransaction(msg, null);
  29. System.out.printf("%s%n", sendResult);
  30. Thread.sleep(10);
  31. } catch (MQClientException | UnsupportedEncodingException e) {
  32. e.printStackTrace();
  33. }
  34. }
  35. for (int i = 0; i < 100000; i++) {
  36. Thread.sleep(1000);
  37. }
  38. producer.shutdown();
  39. }
  40. }

2、实现事务的监听接口

当发送半消息成功时,我们使用 executeLocalTransaction 方法来执行本地事务。它返回前一节中提到的三个事务状态之一。checkLocalTransaction 方法用于检查本地事务状态,并回应消息队列的检查请求。它也是返回前一节中提到的三个事务状态之一。

  1. public class TransactionListenerImpl implements TransactionListener {
  2. private AtomicInteger transactionIndex = new AtomicInteger(0);
  3. private ConcurrentHashMap<String, Integer> localTrans = new ConcurrentHashMap<>();
  4. @Override
  5. public LocalTransactionState executeLocalTransaction(Message msg, Object arg) {
  6. int value = transactionIndex.getAndIncrement();
  7. int status = value % 3;
  8. localTrans.put(msg.getTransactionId(), status);
  9. return LocalTransactionState.UNKNOW;
  10. }
  11. @Override
  12. public LocalTransactionState checkLocalTransaction(MessageExt msg) {
  13. Integer status = localTrans.get(msg.getTransactionId());
  14. if (null != status) {
  15. switch (status) {
  16. case 0:
  17. return LocalTransactionState.UNKNOW;
  18. case 1:
  19. return LocalTransactionState.COMMIT_MESSAGE;
  20. case 2:
  21. return LocalTransactionState.ROLLBACK_MESSAGE;
  22. }
  23. }
  24. return LocalTransactionState.COMMIT_MESSAGE;
  25. }
  26. }

6.2 事务消息使用上的限制

  1. 事务消息不支持延时消息和批量消息。
  2. 为了避免单个消息被检查太多次而导致半队列消息累积,我们默认将单个消息的检查次数限制为 15 次,但是用户可以通过 Broker 配置文件的 transactionCheckMax参数来修改此限制。如果已经检查某条消息超过 N 次的话( N = transactionCheckMax ) 则 Broker 将丢弃此消息,并在默认情况下同时打印错误日志。用户可以通过重写 AbstractTransactionalMessageCheckListener 类来修改这个行为。
  3. 事务消息将在 Broker 配置文件中的参数 transactionTimeout 这样的特定时间长度之后被检查。当发送事务消息时,用户还可以通过设置用户属性 CHECK_IMMUNITY_TIME_IN_SECONDS 来改变这个限制,该参数优先于 transactionTimeout 参数。
  4. 事务性消息可能不止一次被检查或消费。
  5. 提交给用户的目标主题消息可能会失败,目前这依日志的记录而定。它的高可用性通过 RocketMQ 本身的高可用性机制来保证,如果希望确保事务消息不丢失、并且事务完整性得到保证,建议使用同步的双重写入机制。
  6. 事务消息的生产者 ID 不能与其他类型消息的生产者 ID 共享。与其他类型的消息不同,事务消息允许反向查询、MQ服务器能通过它们的生产者 ID 查询到消费者。