RocketMQ发送消息的三种策略


使用 RocketMQ 发送三种类型的消息:同步消息、异步消息和单向消息,其中前两种消息是可靠的,因为会有发送是否成功的应答

这种可靠性同步地发送方式使用的比较广泛,比如:重要的消息通知,短信通知

同步发送

原理

同步发送是指消息发送方发出数据后,同步等待,直到收到接收方(Broker)发回响应之后才发下一个请求。生产者发送消息到Broker,并等待Broker的确认响应。当Broker成功接收并存储消息后,会返回确认响应给生产者。生产者只有在收到确认响应后,才会认为消息发送成功。

特点:

  • 可靠性高:同步发送确保消息被成功存储在Broker中,适用于对消息可靠性要求较高的场景。
  • 延迟较高:由于需要等待Broker的确认响应,同步发送的延迟相对较高。
  • 吞吐量较低:每次发送都需要等待响应,因此吞吐量相对较低。

案例分析:

  • 场景:重要的消息通知,如短信通知、支付确认等。
  • 实现:生产者发送消息后,等待Broker的确认响应,如果收到确认响应,则消息发送成功;如果未收到确认响应,则可以选择重试或记录日志。
public class SyncProducer {
	public static void main(String[] args) throws Exception {
    	// 实例化消息生产者Producer
        DefaultMQProducer producer = new DefaultMQProducer("please_rename_unique_group_name");
    	// 设置NameServer的地址
    	producer.setNamesrvAddr("localhost:9876");
    	// 启动Producer实例
        producer.start();
    	for (int i = 0; i < 100; i++) {
    	    // 创建消息,并指定Topic,Tag和消息体
    	    Message msg = new Message(
                "TopicTest" /* Topic */,
                "TagA" /* Tag */,
        		("Hello RocketMQ " + i).getBytes(RemotingHelper.DEFAULT_CHARSET) /* Message body */);
            
        	// 发送消息到一个Broker
            SendResult sendResult = producer.send(msg);
            // 通过sendResult返回消息是否成功送达
            System.out.printf("%s%n", sendResult);
    	}
    	// 如果不再发送消息,关闭Producer实例。
    	producer.shutdown();
    }
}

异步发送

原理

异步发送是RocketMQ提供的另一种发送策略,它允许生产者在发送消息后立即返回,而不等待Broker的确认响应。生产者发送消息到Broker,并立即返回,不等待Broker的确认响应。生产者可以注册回调函数,当Broker成功接收并存储消息后,会调用回调函数通知生产者。

特点:

  • 延迟较低:异步发送不需要等待Broker的确认响应,因此延迟较低。
  • 吞吐量较高:由于不需要等待响应,异步发送的吞吐量相对较高。
  • 可靠性适中:虽然异步发送不等待确认响应,但通过回调函数可以确保消息被成功存储在Broker中。

案例分析:

  • 场景:对响应时间敏感的业务场景,如在线支付、实时物流更新等。
  • 实现:生产者发送消息后,立即返回并继续处理其他任务。同时,注册回调函数以处理Broker的确认响应或异常通知。
public class AsyncProducer {
	public static void main(String[] args) throws Exception {
    	// 实例化消息生产者Producer
        DefaultMQProducer producer = new DefaultMQProducer("please_rename_unique_group_name");
    	// 设置NameServer的地址
        producer.setNamesrvAddr("localhost:9876");
    	// 启动Producer实例
        producer.start();
        producer.setRetryTimesWhenSendAsyncFailed(0);
	
        int messageCount = 100;
		// 根据消息数量实例化倒计时计算器
        final CountDownLatch2 countDownLatch = new CountDownLatch2(messageCount);
        for (int i = 0; i < messageCount; i++) {
            final int index = i;
            // 创建消息,并指定Topic,Tag和消息体
            Message msg = new Message("TopicTest", "TagA", "OrderID188",
                                      "Hello world".getBytes(RemotingHelper.DEFAULT_CHARSET));

            // SendCallback接收异步返回结果的回调
            producer.send(msg, new SendCallback() {
                // 发送成功回调函数
                @Override
                public void onSuccess(SendResult sendResult) {
                    countDownLatch.countDown();
                    System.out.printf("%-10d OK %s %n", index, sendResult.getMsgId());
                }
                
                @Override
                public void onException(Throwable e) {
                    countDownLatch.countDown();
                    System.out.printf("%-10d Exception %s %n", index, e);
                    e.printStackTrace();
                }
            });
        }
        // 等待5s
        countDownLatch.await(5, TimeUnit.SECONDS);
        // 如果不再发送消息,关闭Producer实例。
        producer.shutdown();
    }
}

单向发送

原理

单向发送是RocketMQ提供的一种简单快速的发送策略,它允许生产者在发送消息后立即返回,不关心Broker是否成功接收和存储消息。生产者发送消息到Broker,并立即返回,不等待Broker的确认响应。生产者不关心消息是否被成功接收和存储,也不注册任何回调函数。

特点:

  • 延迟最低:单向发送不需要等待任何响应,因此延迟最低。
  • 吞吐量最高:由于不需要等待响应,单向发送的吞吐量最高。
  • 可靠性最低:单向发送不关心消息是否被成功接收和存储,适用于对消息可靠性要求不高的场景。

案例分析:

  • 场景:日志收集、非关键性数据上报等。
  • 实现:生产者发送消息后,立即返回并继续处理其他任务。不关注消息是否成功发送到Broker,也不处理任何确认响应或异常通知。
public class OnewayProducer {
	public static void main(String[] args) throws Exception{
    	// 实例化消息生产者Producer
        DefaultMQProducer producer = new DefaultMQProducer("please_rename_unique_group_name");
    	// 设置NameServer的地址
        producer.setNamesrvAddr("localhost:9876");
    	// 启动Producer实例
        producer.start();
    	for (int i = 0; i < 100; i++) {
        	// 创建消息,并指定Topic,Tag和消息体
        	Message msg = new Message("TopicTest","TagA",
                          ("Hello RocketMQ " + i).getBytes(RemotingHelper.DEFAULT_CHARSET));
        	// 发送单向消息,没有任何返回结果
        	producer.sendOneway(msg);
    	}
    	// 如果不再发送消息,关闭Producer实例。
    	producer.shutdown();
    }
}

文章作者: Fuchanglai
版权声明: 本博客所有文章除特別声明外,均采用 CC BY 4.0 许可协议。转载请注明来源 Fuchanglai !
赏
  目录