Java源码示例:org.apache.rocketmq.client.producer.TransactionListener

示例1
public static void main(String[] args) throws MQClientException, InterruptedException {
    TransactionListener transactionListener = new TransactionListenerImpl();
    TransactionMQProducer producer = new TransactionMQProducer("please_rename_unique_group_name");
    ExecutorService executorService = new ThreadPoolExecutor(2, 5, 100, TimeUnit.SECONDS, new ArrayBlockingQueue<Runnable>(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.setTransactionListener(transactionListener);
    producer.start();

    String[] tags = new String[] {"TagA", "TagB", "TagC", "TagD", "TagE"};
    for (int i = 0; i < 10; i++) {
        try {
            Message msg =
                new Message("TopicTest1234", tags[i % tags.length], "KEY" + i,
                    ("Hello RocketMQ " + i).getBytes(RemotingHelper.DEFAULT_CHARSET));
            SendResult sendResult = producer.sendMessageInTransaction(msg, null);
            System.out.printf("%s%n", sendResult);

            Thread.sleep(10);
        } catch (MQClientException | UnsupportedEncodingException e) {
            e.printStackTrace();
        }
    }

    for (int i = 0; i < 100000; i++) {
        Thread.sleep(1000);
    }
    producer.shutdown();
}
 
示例2
@Override
public TransactionListener checkListener() {
    if (this.defaultMQProducer instanceof TransactionMQProducer) {
        TransactionMQProducer producer = (TransactionMQProducer) defaultMQProducer;
        return producer.getTransactionListener();
    }

    return null;
}
 
示例3
public static void main(String[] args) throws MQClientException, InterruptedException {
    TransactionListener transactionListener = new TransactionListenerImpl();
    TransactionMQProducer producer = new TransactionMQProducer("please_rename_unique_group_name");
    ExecutorService executorService = new ThreadPoolExecutor(2, 5, 100, TimeUnit.SECONDS, new ArrayBlockingQueue<Runnable>(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.setTransactionListener(transactionListener);
    producer.start();

    String[] tags = new String[] {"TagA", "TagB", "TagC", "TagD", "TagE"};
    for (int i = 0; i < 10; i++) {
        try {
            Message msg =
                new Message("TopicTest1234", tags[i % tags.length], "KEY" + i,
                    ("Hello RocketMQ " + i).getBytes(RemotingHelper.DEFAULT_CHARSET));
            SendResult sendResult = producer.sendMessageInTransaction(msg, null);
            System.out.printf("%s%n", sendResult);

            Thread.sleep(10);
        } catch (MQClientException | UnsupportedEncodingException e) {
            e.printStackTrace();
        }
    }

    for (int i = 0; i < 100000; i++) {
        Thread.sleep(1000);
    }
    producer.shutdown();
}
 
示例4
/**
 * 获取事务回查checklistener
 * @return ;
 */
@Override
public TransactionListener getCheckListener() {
    if (this.defaultMQProducer instanceof TransactionMQProducer) {
        TransactionMQProducer producer = (TransactionMQProducer) defaultMQProducer;
        return producer.getTransactionListener();
    }
    return null;
}
 
示例5
public static void main(String[] args) throws MQClientException, InterruptedException {
	TransactionListener transactionListener = new TransactionListenerImpl();
	TransactionMQProducer producer = new TransactionMQProducer("transactionProducerGroupName");
	producer.setNamesrvAddr("192.168.237.128:9876");
	ExecutorService executorService = new ThreadPoolExecutor(2, 5, 100, TimeUnit.SECONDS,
			new ArrayBlockingQueue<Runnable>(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.setTransactionListener(transactionListener);
	producer.start();

	String[] tags = new String[] { "TagA", "TagB", "TagC", "TagD", "TagE" };
	for (int i = 0; i < 1; i++) {
		try {
			Message msg = new Message("TopicTest1234", tags[i % tags.length], "KEY" + i,
					("Hello RocketMQ " + i).getBytes(RemotingHelper.DEFAULT_CHARSET));
			System.out.println("start send message " + msg);
			SendResult sendResult = producer.sendMessageInTransaction(msg, null);
			System.out.printf("%s%n", sendResult);

			Thread.sleep(10);
		} catch (MQClientException | UnsupportedEncodingException e) {
			e.printStackTrace();
		}
	}

	for (int i = 0; i < 100000; i++) {
		Thread.sleep(1000);
	}
	producer.shutdown();
}
 
示例6
public static void main(String[] args) throws MQClientException, InterruptedException {
    TransactionListener transactionListener = new TransactionListenerImpl();
    TransactionMQProducer producer = new TransactionMQProducer("please_rename_unique_group_name");
    ExecutorService executorService = new ThreadPoolExecutor(2, 5, 100, TimeUnit.SECONDS, new ArrayBlockingQueue<Runnable>(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.setTransactionListener(transactionListener);
    producer.start();

    String[] tags = new String[] {"TagA", "TagB", "TagC", "TagD", "TagE"};
    for (int i = 0; i < 10; i++) {
        try {
            Message msg =
                new Message("TopicTest1234", tags[i % tags.length], "KEY" + i,
                    ("Hello RocketMQ " + i).getBytes(RemotingHelper.DEFAULT_CHARSET));
            SendResult sendResult = producer.sendMessageInTransaction(msg, null);
            System.out.printf("%s%n", sendResult);

            Thread.sleep(10);
        } catch (MQClientException | UnsupportedEncodingException e) {
            e.printStackTrace();
        }
    }

    for (int i = 0; i < 100000; i++) {
        Thread.sleep(1000);
    }
    producer.shutdown();
}
 
示例7
@Override
public TransactionListener getCheckListener() {
    if (this.defaultMQProducer instanceof TransactionMQProducer) {
        TransactionMQProducer producer = (TransactionMQProducer) defaultMQProducer;
        return producer.getTransactionListener();
    }
    return null;
}
 
示例8
protected void create(boolean useTLS, TransactionListener transactionListener) {
    producer = new TransactionMQProducer();
    producer.setProducerGroup(getProducerGroupName());
    producer.setInstanceName(getProducerInstanceName());
    producer.setTransactionListener(transactionListener);
    producer.setUseTLS(useTLS);

    if (nsAddr != null) {
        producer.setNamesrvAddr(nsAddr);
    }
}
 
示例9
public static RMQTransactionalProducer getTransactionalProducer(String nsAddr, String topic, TransactionListener transactionListener) {
    RMQTransactionalProducer producer = new RMQTransactionalProducer(nsAddr, topic, false, transactionListener);
    if (debug) {
        producer.setDebug();
    }
    mqClients.add(producer);
    return producer;
}
 
示例10
public static void main(String[] args) throws MQClientException, InterruptedException {

        TransactionListener transactionListener = new TransactionListenerImpl();
        TransactionMQProducer producer = new TransactionMQProducer("transaction_Producer");
        producer.setNamesrvAddr(NAMESRVADDR);

        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.setTransactionListener(transactionListener);
        producer.start();

        String[] tags = new String[]{"TagA", "TagB", "TagC", "TagD", "TagE"};
        for (int i = 0; i < 10; i++) {
            try {
                Message message = new Message("TopicTransactionTest", tags[i % tags.length],
                        "KEY" + i, ("Hello RocketMQ " + i).getBytes(RemotingHelper.DEFAULT_CHARSET));
                SendResult sendResult = producer.sendMessageInTransaction(message, null);
                System.out.printf("%s%n", sendResult);

                Thread.sleep(10);
            } catch (MQClientException | UnsupportedEncodingException e) {
                e.printStackTrace();
            }
        }

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

        producer.shutdown();

    }
 
示例11
@Override
protected TransactionListener newTxListener() {
    return new RocketMQTransactionListener();
}
 
示例12
public RMQTransactionalProducer(String nsAddr, String topic, TransactionListener transactionListener) {
    this(nsAddr, topic, false, transactionListener);
}
 
示例13
public RMQTransactionalProducer(String nsAddr, String topic, boolean useTLS, TransactionListener transactionListener) {
    super(topic);
    this.nsAddr = nsAddr;
    create(useTLS, transactionListener);
    start();
}
 
示例14
/**
 * 获取 transactionListener
 * @return ;
 */
TransactionListener getCheckListener();
 
示例15
TransactionListener checkListener(); 
示例16
TransactionListener getCheckListener();