Rocketmq: Rocketmq4.8 official transaction message example does not call the transaction callback interface

Created on 25 Mar 2021  ·  3Comments  ·  Source: apache/rocketmq

rocketmq4.8官方事务消息例子不调用事务回查接口(checkLocalTransaction)

public class TransactionListenerImpl implements TransactionListener {
private AtomicInteger transactionIndex = new AtomicInteger(0);
private ConcurrentHashMap localTrans = new ConcurrentHashMap<>();

@Override
public LocalTransactionState executeLocalTransaction(Message msg, Object arg) {
    System.out.println("执行本地事务,msg="+ JSON.toJSONString(msg));
    int value = transactionIndex.getAndIncrement();
    int status = value % 3;
    localTrans.put(msg.getTransactionId(), status);
    try {
        Thread.sleep(20000);
    } catch (InterruptedException e) {
        e.printStackTrace();
    }
    return LocalTransactionState.UNKNOW;
}

@Override
public LocalTransactionState checkLocalTransaction(MessageExt msg) {
    System.out.println("执行回查本地事务接口,msg="+ JSON.toJSONString(msg));
    Integer status = localTrans.get(msg.getTransactionId());
    if (status != null) {
        switch (status) {
            case 0:
                return LocalTransactionState.UNKNOW;
            case 1:
                return LocalTransactionState.COMMIT_MESSAGE;
            case 2:
                return LocalTransactionState.ROLLBACK_MESSAGE;
        }
    }
    return LocalTransactionState.COMMIT_MESSAGE;
}

}

public class TransactionProducer {
public static void main(String[] args) throws Exception{
TransactionMQProducer producer = new TransactionMQProducer("zkt_test_transaction_producer");
producer.setNamesrvAddr("192.168.4.176:9876");
producer.setTransactionListener(new TransactionListenerImpl());
ExecutorService executorService = new ThreadPoolExecutor(2, 5, 100, TimeUnit.SECONDS, new ArrayBlockingQueue(2000), new ThreadFactory() {
@Override
public Thread newThread(Runnable r) {
Thread thread = new Thread(r);
thread.setName("client-transaction-msg-check-thread");
return thread;
}
});
producer.setExecutorService(executorService);
producer.start();

    String[] tags = new String[] {"TagA", "TagB", "TagC", "TagD", "TagE"};
    for(int i=0;i<10;i++){
        Message msg = new Message("TopicTest1",tags[i%tags.length],"Key"+i,("Hello RocketMq "+i).getBytes(StandardCharsets.UTF_8));
        TransactionSendResult sendResult = producer.sendMessageInTransaction(msg, null);
        System.out.printf("%s%n", sendResult);
        Thread.sleep(10);
    }

    for (int i = 0; i < 100000; i++) {
        Thread.sleep(1000);
    }

    producer.shutdown();
}

}

Most helpful comment

我测试过,是可以回调的,一样的4.8版本

All 3 comments

我测试过,是可以回调的,一样的4.8版本

按照官方的步骤一步一步走下来,才发现rocketmq-client的版本是4.3.0,希望更新一下吧。
image

ping @ShannonDing

Was this page helpful?
0 / 5 - 0 ratings