Java RocketMQ案例怎么写

wen java案例 3

本文目录导读:

Java RocketMQ案例怎么写

  1. 环境准备
  2. 生产者示例
  3. 消费者示例
  4. 有序消息生产者
  5. 事务消息示例
  6. 批量消息示例
  7. 运行说明
  8. 注意事项

我来给你一个完整的Java RocketMQ案例,包括生产者和消费者的实现。

环境准备

Maven依赖

<dependency>
    <groupId>org.apache.rocketmq</groupId>
    <artifactId>rocketmq-client</artifactId>
    <version>4.9.4</version>
</dependency>

生产者示例

import org.apache.rocketmq.client.producer.DefaultMQProducer;
import org.apache.rocketmq.client.producer.SendResult;
import org.apache.rocketmq.common.message.Message;
import org.apache.rocketmq.remoting.common.RemotingHelper;
public class SimpleProducer {
    public static void main(String[] args) throws Exception {
        // 1. 创建生产者实例,指定生产者组名
        DefaultMQProducer producer = new DefaultMQProducer("producer_group");
        // 2. 指定NameServer地址
        producer.setNamesrvAddr("localhost:9876");
        // 3. 启动生产者
        producer.start();
        try {
            // 4. 发送消息
            for (int i = 0; i < 10; i++) {
                // 创建消息,指定Topic、Tag和消息体
                Message msg = new Message(
                    "TopicTest",  // Topic
                    "TagA",       // Tag
                    ("Hello RocketMQ " + i).getBytes(RemotingHelper.DEFAULT_CHARSET)  // 消息内容
                );
                // 5. 发送消息并获取结果
                SendResult sendResult = producer.send(msg);
                System.out.printf("发送结果: %s%n", sendResult);
            }
        } catch (Exception e) {
            e.printStackTrace();
        } finally {
            // 6. 关闭生产者
            producer.shutdown();
        }
    }
}

消费者示例

1 简单消费者(Push模式)

import org.apache.rocketmq.client.consumer.DefaultMQPushConsumer;
import org.apache.rocketmq.client.consumer.listener.ConsumeConcurrentlyContext;
import org.apache.rocketmq.client.consumer.listener.ConsumeConcurrentlyStatus;
import org.apache.rocketmq.client.consumer.listener.MessageListenerConcurrently;
import org.apache.rocketmq.common.consumer.ConsumeFromWhere;
import org.apache.rocketmq.common.message.MessageExt;
import java.util.List;
public class SimpleConsumer {
    public static void main(String[] args) throws Exception {
        // 1. 创建消费者实例,指定消费者组名
        DefaultMQPushConsumer consumer = new DefaultMQPushConsumer("consumer_group");
        // 2. 指定NameServer地址
        consumer.setNamesrvAddr("localhost:9876");
        // 3. 设置消费位置(从最新消费还是从头消费)
        consumer.setConsumeFromWhere(ConsumeFromWhere.CONSUME_FROM_FIRST_OFFSET);
        // 4. 订阅Topic和Tag
        consumer.subscribe("TopicTest", "*");  // "*"表示订阅所有Tag
        // 5. 注册消息监听器
        consumer.registerMessageListener(new MessageListenerConcurrently() {
            @Override
            public ConsumeConcurrentlyStatus consumeMessage(List<MessageExt> msgs,
                   ConsumeConcurrentlyContext context) {
                try {
                    for (MessageExt msg : msgs) {
                        String body = new String(msg.getBody(), "UTF-8");
                        System.out.printf("接收到消息: Topic=%s, Tag=%s, 消息内容=%s%n",
                            msg.getTopic(), msg.getTags(), body);
                    }
                    // 消费成功
                    return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
                } catch (Exception e) {
                    e.printStackTrace();
                    // 消费失败,稍后重试
                    return ConsumeConcurrentlyStatus.RECONSUME_LATER;
                }
            }
        });
        // 6. 启动消费者
        consumer.start();
        System.out.println("消费者启动成功");
    }
}

2 有序消息消费者

import org.apache.rocketmq.client.consumer.DefaultMQPushConsumer;
import org.apache.rocketmq.client.consumer.listener.ConsumeOrderlyContext;
import org.apache.rocketmq.client.consumer.listener.ConsumeOrderlyStatus;
import org.apache.rocketmq.client.consumer.listener.MessageListenerOrderly;
import org.apache.rocketmq.common.message.MessageExt;
import java.util.List;
public class OrderConsumer {
    public static void main(String[] args) throws Exception {
        DefaultMQPushConsumer consumer = new DefaultMQPushConsumer("order_consumer_group");
        consumer.setNamesrvAddr("localhost:9876");
        // 订阅订单Topic
        consumer.subscribe("OrderTopic", "*");
        // 注册有序消息监听器
        consumer.registerMessageListener(new MessageListenerOrderly() {
            @Override
            public ConsumeOrderlyStatus consumeMessage(List<MessageExt> msgs,
                   ConsumeOrderlyContext context) {
                context.setAutoCommit(true);
                for (MessageExt msg : msgs) {
                    String orderId = msg.getKeys();  // 获取消息Key(订单ID)
                    String body = new String(msg.getBody());
                    System.out.printf("消费订单消息: 订单ID=%s, 内容=%s%n", orderId, body);
                }
                return ConsumeOrderlyStatus.SUCCESS;
            }
        });
        consumer.start();
        System.out.println("有序消费者启动成功");
    }
}

有序消息生产者

import org.apache.rocketmq.client.producer.DefaultMQProducer;
import org.apache.rocketmq.client.producer.MessageQueueSelector;
import org.apache.rocketmq.client.producer.SendResult;
import org.apache.rocketmq.common.message.Message;
import org.apache.rocketmq.common.message.MessageQueue;
import java.util.List;
public class OrderProducer {
    public static void main(String[] args) throws Exception {
        DefaultMQProducer producer = new DefaultMQProducer("order_producer_group");
        producer.setNamesrvAddr("localhost:9876");
        producer.start();
        try {
            // 模拟发送同一个订单的消息到同一个队列
            for (int i = 0; i < 10; i++) {
                // 假设订单ID为 order_001
                String orderId = "order_001";
                Message msg = new Message(
                    "OrderTopic",
                    "TagA",
                    orderId,  // 设置消息Key为订单ID
                    ("订单操作步骤 " + i).getBytes()
                );
                // 使用MessageQueueSelector确保同一个订单的消息发送到同一个队列
                SendResult result = producer.send(msg, new MessageQueueSelector() {
                    @Override
                    public MessageQueue select(List<MessageQueue> mqs, Message msg, Object arg) {
                        // 根据订单ID选择队列
                        String orderId = (String) arg;
                        long index = Math.abs(orderId.hashCode()) % mqs.size();
                        return mqs.get((int) index);
                    }
                }, orderId);
                System.out.printf("发送有序消息: %s%n", result);
            }
        } catch (Exception e) {
            e.printStackTrace();
        } finally {
            producer.shutdown();
        }
    }
}

事务消息示例

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.client.producer.SendResult;
import org.apache.rocketmq.common.message.Message;
import org.apache.rocketmq.common.message.MessageExt;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.TimeUnit;
public class TransactionProducer {
    public static void main(String[] args) throws Exception {
        TransactionMQProducer producer = new TransactionMQProducer("transaction_producer_group");
        producer.setNamesrvAddr("localhost:9876");
        // 设置事务监听器
        producer.setTransactionListener(new TransactionListener() {
            private ConcurrentHashMap<String, Boolean> localTrans = new ConcurrentHashMap<>();
            @Override
            public LocalTransactionState executeLocalTransaction(Message msg, Object arg) {
                // 执行本地事务
                try {
                    System.out.println("执行本地事务");
                    // 模拟本地事务操作(如数据库操作)
                    Thread.sleep(1000);
                    // 本地事务成功
                    localTrans.put(msg.getTransactionId(), true);
                    return LocalTransactionState.COMMIT_MESSAGE;
                } catch (Exception e) {
                    // 本地事务失败,回滚
                    localTrans.put(msg.getTransactionId(), false);
                    return LocalTransactionState.ROLLBACK_MESSAGE;
                }
            }
            @Override
            public LocalTransactionState checkLocalTransaction(MessageExt msg) {
                // 检查本地事务状态
                Boolean success = localTrans.get(msg.getTransactionId());
                if (success != null && success) {
                    return LocalTransactionState.COMMIT_MESSAGE;
                }
                return LocalTransactionState.UNKNOW;
            }
        });
        producer.start();
        try {
            Message msg = new Message("TransactionTopic", "TagA", 
                "事务消息测试".getBytes());
            // 发送事务消息
            SendResult result = producer.sendMessageInTransaction(msg, null);
            System.out.printf("发送结果: %s%n", result);
        } catch (Exception e) {
            e.printStackTrace();
        }
        producer.shutdown();
    }
}

批量消息示例

import org.apache.rocketmq.client.producer.DefaultMQProducer;
import org.apache.rocketmq.client.producer.SendResult;
import org.apache.rocketmq.common.message.Message;
import java.util.ArrayList;
import java.util.List;
public class BatchProducer {
    public static void main(String[] args) throws Exception {
        DefaultMQProducer producer = new DefaultMQProducer("batch_producer_group");
        producer.setNamesrvAddr("localhost:9876");
        producer.start();
        // 创建批量消息
        List<Message> messages = new ArrayList<>();
        for (int i = 0; i < 10; i++) {
            Message msg = new Message("BatchTopic", "TagA", 
                ("批量消息 " + i).getBytes());
            messages.add(msg);
        }
        try {
            // 发送批量消息
            SendResult result = producer.send(messages);
            System.out.printf("批量发送结果: %s%n", result);
        } catch (Exception e) {
            e.printStackTrace();
        } finally {
            producer.shutdown();
        }
    }
}

运行说明

启动RocketMQ(需要先安装Docker)

# 使用Docker启动RocketMQ
docker run -d -p 9876:9876 --name rocketmq-namesrv rocketmq-namesrv
docker run -d -p 10911:10911 --name rocketmq-broker rocketmq-broker

运行步骤

  1. 先启动消费者程序
  2. 再启动生产者程序
  3. 观察消费者控制台输出的消息

注意事项

  1. 消息顺序性:需要有序消息时,必须使用MessageQueueSelector确保消息路由到同一队列
  2. 消费重试:消费失败时返回RECONSUME_LATER,消息会重试
  3. 消息去重:Producer发送消息时建议设置唯一Key,便于Consumer去重
  4. 异常处理:生产者和消费者都要妥善处理异常,避免资源泄漏
  5. 配置优化:生产环境需要根据实际需求调整生产者和消费者的配置参数

这个案例涵盖了RocketMQ的主要使用场景,你可以根据实际需求进行相应的调整和扩展。

抱歉,评论功能暂时关闭!