架构师修行录

MQ组件重磅更新!灵活切换Rocket/Redis/Kafka/Rabbit多种实现

哈喽,各位代码战士们,我是Jensen,一个梦想着和大家一起在代码的海洋里遨游,顺便捡起那些散落的知识点的程序员小伙伴。

笔者进公司一年多,一直想上MQ(消息队列)很久了,但这边的技术相对保守,不是因为MQ不好,要看团队适不适合,团队里总有一些声音不想上MQ:

  • MQ占用资源太大了,客户的服务器资源不足,这么多微服务,16G的内存已经扛不住了

  • 上了MQ谁来维护?运维也是有成本的,别到时开发用着很爽,运维就一堆负担

  • 不用MQ,系统调用是同步的,被调用方挂了,调用方直接能感知到,业务做重试就行了,不会有乱七八糟的问题

  • 还有,怎么处理分布式事务问题?

确实是一些比较现实绕不开的话题,但上MQ也有非常多的好处,这里我说两点:
  • 系统解耦:微服务各个系统是有不同定位的,让一个底层系统去往业务系统调就很不合理。比如我们有个租赁业务的微服务,调用链路就不应该由基础平台服务去调用租赁服务,而应该由租赁去监听基础平台的事件,依赖反转嘛,SpringIOC也是这个道理

  • 职责划分:通过MQ解耦后,往往开发人员的职责边界也更清晰。只需要把MQ事件发布出去就行了,我管你消费者是哪个接口呢?我就非常受不了这边的IOT服务,想把当前的IOT设备的在线状态同步出去业务系统,既调用基础平台的同步设备状态接口,又调用租赁设备的同步设备状态接口,万一后续又多了个业务要做设备的,是不是还得IOT再调一次?还有商城的店铺同步功能也一样,需要在Nacos配置哪个租户需要调哪个系统的同步店铺接口,多开一个租户还得配置一下,不然都不知道为什么店铺没有同步出去。

关于上不上MQ,团队内赞同者还是占多数的,只要能解决现有的问题,让程序猿开发起来爽一些不好么?最后我们讨论多次后采用了相对权衡的方案——用Redis的发布订阅模型去实现MQ。
那怎么去实现这个MQ呢?实现Redis的发布订阅模型本身是非常简单的,AI分分钟能帮你写好一个Redis工具类,但MQ的范畴更大,需要一定的设计模式,例如SpringStream是能做到不更改代码而换MQ实现的,我们自己去实现,也要参考这个思想。

之前在我的D3Boot框架内也写过一个base-mq组件,当时只有kafka实现,现在看来有点鸡肋了,得升级一下,所以我花了几天功夫,把主流的几个MQ的共通点抽了出来,再全部实现了一遍,话不多说,来看看MQ组件升级后是怎么用的。

0x1

新MQ组件的打开方式

我们来举一个例子:商城服务在创建店铺后,需要发MQ事件,由IOT服务和租赁服务处理后续的逻辑,那么MQ的生产者是商城,MQ的消费者分别是IOT、租赁。

首先定义一个店铺同步事件ShopSyncEvent:

@Data@Builder@AllArgsConstructor@NoArgsConstructorpublic class ShopSyncEvent extends MQEvent {    // 店铺ID    private String shopId;    // 店铺名称    private String shopName;    // 店铺详情    private JSONObject shopInfo;
@Override public String getTopic() { return "mall"; }
@Override public String getTag() { return "shop_sync"; }}

MQ生产者(商城服务)发布MQ事件:

// 在创建店铺成功后,发布店铺同步MQ事件ShopSyncEvent.builder().shopId(shopInfo.getId()).shopName(shopInfo.getName())  .shopInfo(new JSONObject(shopInfo)).build().tenantId(shopInfo.getTenantId())  .publish();

MQ消费者(IOT服务、租赁服务)分别监听MQ事件:

@Slf4j@Componentpublic class MallMQListener {
/** * 店铺同步 */ @MQEventListener(topic = "mall", tags = "shop_sync") public void onShopSync(ShopSyncEvent event) { log.info("店铺同步至租赁/IOT:{}", event.getShopInfo()); // TODO 具体的消费逻辑 }
}

那么问题来了,具体的MQ在框架内是如何实现的?好学的你此时选择了继续往下看~

0x2

定义核心组件

首先我们使用门面模式去设计一个自定义注解@MQEventListener,该注解与具体的MQ实现无关,是最轻量的存在:

/** * MQ事件监听器,配合MQEvent使用 */@Documented@Target(ElementType.METHOD)@Retention(RetentionPolicy.RUNTIME)public @interface MQEventListener {
// 消费者组,不设置默认为当前启动的应用名${spring.application.name},一般多实例同一个消费者需要保持一致,避免同一业务重复消费 String group() default "";
// 主题,大分类,消费线程隔离,用于区分不同业务。也可作为小分类使用 String topic() default "DEFAULT";
// 标签列表,小分类,同一主题下共享消费线程,支持:通配符*、或||、非-,支持复合表达式如:A || B -C String tags() default "*";}

再设计一个基类MQEvent,发送与接收的MQ事件统一继承该基类:

/** * MQ事件基类 */@Datapublic class MQEvent implements Serializable {    // 消息ID,默认当前时间戳    protected String msgId;    // 主题,配置base-mq.default-topic后无须每次指定    protected String topic;    // 标签,只支持单个标签,多标签需要分开发送    protected String tag;    // 租户ID,默认从线程上下文获取    protected String tenantId;
// 发布MQ事件 public void publish() { this.publish(getTopic(), getTag(), getTenantId()); }
// 发布MQ事件 public void publish(String topic) { this.publish(topic, getTag(), getTenantId()); }
// 发布MQ事件 public void publish(String topic, String tag) { this.publish(topic, tag, getTenantId()); }
// 发布MQ事件 public void publish(String topic, String tag, String tenantId) { setTopic(topic); if (getTopic() == null) { setTopic(SpringContext.getEnv().getProperty("base-mq.default-topic", "DEFAULT")); } setTag(tag); setTenantId(tenantId != null ? tenantId : ThreadContext.get(ContextConstants.TENANT_ID)); setMsgId(getMsgId() == null ? String.valueOf(System.currentTimeMillis()) : getMsgId()); // 这里的BaseContext是在基础框架全局使用的上下文,用于解耦各个组件具体的技术代码 if (BaseContext.contains("MQEventPublisher")) { BaseContext.<String, Consumer<MQEvent>>get("MQEventPublisher").accept(this); } }
public <T extends MQEvent> T tenantId(String tenantId) { this.tenantId = tenantId; return (T) this; }
}

接下来设计一个用于抽象不同MQ实现的配置文件BaseMQProperties:

@Data@ConfigurationProperties(prefix = "base-mq")public class BaseMQProperties {    // 是否启用    private boolean enable = false;    // MQ实现:目前支持kafka|rocket|redis|rabbit,一个应用只用一套MQ发布和订阅    private String impl = "none";    // 服务地址    private String server = "";    // 命名空间,或作为所有主题的前缀,用于环境隔离等场景,如UAT/生产/租户环境共用一个MQ    private String namespace = "";    // 是否持久化到本地,需要定义实现了MQEventStorer的Bean    private boolean persist = false;    // 序列化器,定义实现了MQEventSerialization的Bean    private String serialization = "JsonMQEventSerialization";    // 是否自动提交    private boolean autoAck = false;    // 发送失败重试次数    private int retries = 0;    // 用户    private String username = "";    // 密码    private String password = "";    // 生产者组    private String producerGroup = "DEFAULT";    // 默认主题    private String defaultTopic = "DEFAULT";    // 交换机    private String exchange = "";
public String namespace(String concat) { if (this.namespace != null && !this.namespace.isEmpty()) { return namespace + concat; } return ""; }
}

最最重要的MQ组件MQClient,设计为一个接口,并在接口内定义default方法:

/** * MQClient接口,MQ实现类实现该接口做差异化实现 */public interface MQClient {    // MQ实现    String impl();
// MQ配置 default BaseMQProperties config() { return SpringContext.getBean(BaseMQProperties.class); }
// 整体初始化 default void init() { if (!config().isEnable() || !Objects.equals(config().getImpl(), impl())) return; Logger log = LoggerFactory.getLogger("### BASE-MQ : " + impl() + "Client ###"); // 初始化MQ事件发布器 if (Objects.equals(config().getImpl(), impl()) && !BaseContext.contains("MQEventPublisher")) { log.info("Initializing MQEventPublisher"); Consumer<MQEvent> producer = initProducer(); if (producer != null) { BaseContext.inject("MQEventPublisher", producer); } } // 初始化MQ事件监听器 log.info("Initializing MQEventListener"); String[] beanNames = SpringContext.getBeanDefinitionNames(); List<MQListener> mqListeners = new ArrayList<>(); for (String name : beanNames) { try { // 找出所有Bean下注解了@MQEventListener的方法 Class<?> targetClass = AopUtils.getTargetClass(SpringContext.getBean(name)); Method[] declaredMethods = targetClass.getDeclaredMethods(); for (Method listenerMethod : declaredMethods) { MQEventListener mqEventListener = AnnotationUtils.findAnnotation(listenerMethod, MQEventListener.class); if (mqEventListener != null) { // 消费者组,不设置默认按Spring工程的应用名隔离,当然也可以自定义 String group = mqEventListener.group(); if (group.isEmpty()) { group = SpringContext.getEnv().getProperty("spring.application.name"); } // 取出首个参数,用于接收MQ事件后自动解析到该类型 Class<MQEvent> eventClass = (Class<MQEvent>) listenerMethod.getParameterTypes()[0]; mqListeners.add(MQListener.builder() .listenerMethod(listenerMethod) .eventClass(eventClass) .consumer(mqEvent -> { try { // 接收并解析为MQEvent后,具体调用该方法的位置 listenerMethod.invoke(SpringContext.getBean(listenerMethod.getDeclaringClass()), mqEvent); } catch (Exception e) { throw new RuntimeException(e); } }) .group(group) .topic(mqEventListener.topic()) .tags(mqEventListener.tags()) .build()); } } } catch (Exception ignore) { } } for (MQListener listener : mqListeners) { try { // 具体创建消费者的逻辑,不同MQ实现不一样 initConsumer(listener); listener.setSuccess(true); } catch (Exception e) { log.error("Listen MQ [{}] failed!", listener.topicAndTags(), e); } } log.info("Listening MQ: {}", mqListeners.stream().filter(MQListener::isSuccess).map(MQListener::topicAndTags).collect(Collectors.joining(","))); }
// 初始化生产者 Consumer<MQEvent> initProducer();
// 初始化消费者 void initConsumer(MQListener mqListener) throws Exception;
// 启动 void start();
// 序列化器,默认取配置了的实现了MQEventSerialization的Bean default MQEventSerialization serialization() { return SpringContext.getBean(config().getSerialization()); }
// 消费MQ事件 default void consume(MQListener listener, MQEvent mqEvent) { Logger log = LoggerFactory.getLogger("### BASE-MQ : " + impl() + "Client ###"); log.info("Consume MQ: {}", serialization().<String>serialize(mqEvent)); // MQ事件持久化 if (config().isPersist()) { try { MQEventStorer mqEventStorer = SpringContext.getBean(MQEventStorer.class); mqEventStorer.store(mqEvent.getTopic(), mqEvent); } catch (Exception e) { log.error("Persist MQ failed: {}", serialization().serialize(mqEvent), e); } } // 再调用方法执行业务 listener.getConsumer().accept(mqEvent); }
}

最后咱们定义好配置类,统一加载那些实现了MQClient的Bean:

/** * 基础MQ配置 */@Order(PriorityOrdered.HIGHEST_PRECEDENCE + 10)@Configuration@ConditionalOnProperty(prefix = "base-mq", name = "enable", havingValue = "true")public class BaseMQConfig {    @Autowired    private List<MQClient> mqClients;
@PostConstruct public void init() { if (mqClients != null && !mqClients.isEmpty()) { // 获取所有实现了MQClient的MQ客户端实现类 for (MQClient mqClient : mqClients) { // 先初始化 mqClient.init(); // 再启动 mqClient.start(); } } }}

核心组件中的其它类比较简单,这里先略过。

0x3

定义不同的MQClient实现

各个MQ之间的特性在这里就不一一介绍了,其中70%左右的代码由阿里的通义千问AI实现(推荐大家平时多用用),直接看实现代码吧,相关作用在注释里有介绍:

一、RocketMQ实现

/** * Rocket客户端实现 */@Slf4j(topic = "### BASE-MQ : rocketClient ###")public final class RocketClient implements MQClient {    private static Map<String, Supplier> LISTENERS = new ConcurrentHashMap<>();
@Override public String impl() { return "rocket"; }
@Override public Consumer<MQEvent> initProducer() { try { // 创建 DefaultMQProducer 实例 DefaultMQProducer rocketProducer; if (!config().getUsername().isEmpty() && !config().getPassword().isEmpty()) { // 需要鉴权 rocketProducer = new DefaultMQProducer(config().getNamespace(), config().getProducerGroup(), new AclClientRPCHook(new SessionCredentials(config().getUsername(), config().getPassword()))); } else { rocketProducer = new DefaultMQProducer(config().getNamespace(), config().getProducerGroup()); } rocketProducer.setNamesrvAddr(config().getServer()); // 启动rocketProducer实例 rocketProducer.start(); DefaultMQProducer finalRocketProducer = rocketProducer; return mqEvent -> { // 定义具体的MQ事件发布逻辑 String message = serialization().serialize(mqEvent); log.info("Publish MQ: {}", message); String tags = mqEvent.getTag() == null ? "" : mqEvent.getTag(); try { finalRocketProducer.send(new Message(mqEvent.getTopic(), tags, message.getBytes(RemotingHelper.DEFAULT_CHARSET))); } catch (Exception e) { log.error("Publish MQ: {} failed!", message, e); } }; } catch (MQClientException e) { log.error("创建Rocket生产者失败", e); } return null; }
@Override public void initConsumer(MQListener mqListener) throws Exception { DefaultMQPushConsumer defaultMQPushConsumer; if (!config().getUsername().isEmpty() && !config().getPassword().isEmpty()) { // 需要鉴权的情况 defaultMQPushConsumer = new DefaultMQPushConsumer(config().getNamespace(), mqListener.getGroup(), new AclClientRPCHook(new SessionCredentials(config().getUsername(), config().getPassword()))); } else { defaultMQPushConsumer = new DefaultMQPushConsumer(config().getNamespace(), mqListener.getGroup()); } defaultMQPushConsumer.setNamesrvAddr(config().getServer()); defaultMQPushConsumer.subscribe(mqListener.getTopic(), mqListener.getTags()); // 创建MessageListenerConcurrently实例 defaultMQPushConsumer.registerMessageListener((MessageListenerConcurrently) (messageExts, context) -> { // 获取方法参数 for (MessageExt messageExt : messageExts) { MQEvent mqEvent = null; try { // 反序列化为MQEvent对象 mqEvent = serialization().deserialize(new String(messageExt.getBody(), StandardCharsets.UTF_8), mqListener.getEventClass()); consume(mqListener, mqEvent); } catch (Exception e) { log.error("Consume MQ failed: {}", mqEvent, e); if (!(e instanceof ServiceException)) { return ConsumeConcurrentlyStatus.RECONSUME_LATER; } } } return ConsumeConcurrentlyStatus.CONSUME_SUCCESS; }); LISTENERS.put(mqListener.getTopic(), () -> defaultMQPushConsumer); }
@Override @SneakyThrows public void start() { if (!config().isEnable() || !Objects.equals(config().getImpl(), impl())) return; // 在此统一启动所有消费者 for (Supplier<DefaultMQPushConsumer> listener : LISTENERS.values()) { listener.get().start(); } }}

二、Kafka实现

之前底层依赖使用spring-data-kafka,用了两年,回过头来发现Spring封装得太死了,后来换成了kafka-clients:

/** * Kafka客户端实现 */@Slf4j(topic = "### BASE-MQ : kafkaClient ###")public final class KafkaClient implements MQClient, DisposableBean {    private static List<Supplier> LISTENERS = new ArrayList<>();
@Override public String impl() { return "kafka"; }
@Override public Consumer<MQEvent> initProducer() { // 创建 Producer 配置 Properties props = new Properties(); props.put("bootstrap.servers", config().getServer()); // Kafka broker 地址 props.put("acks", "all"); props.put("retries", config().getRetries());// 失败重试次数 props.put("batch.size", 16384); props.put("linger.ms", 1); props.put("buffer.memory", 33554432); props.put("key.serializer", "org.apache.kafka.common.serialization.StringSerializer"); props.put("value.serializer", "org.apache.kafka.common.serialization.StringSerializer"); // 创建KafkaProducer实例 Producer<String, String> producer = new KafkaProducer<>(props); return mqEvent -> { // 定义具体的MQ事件发布逻辑 String message = serialization().serialize(mqEvent); log.info("Publish MQ: {}", message); try { // topic使用下划线的方式拼接 producer.send(new ProducerRecord<>(config().namespace("_") + mqEvent.getTopic(), message)); } catch (Exception e) { log.error("Publish MQ: {} failed!", message, e); } }; }
@Override public void initConsumer(MQListener mqListener) throws Exception { Properties props = new Properties(); props.put("bootstrap.servers", config().getServer()); props.put("group.id", mqListener.getGroup()); // 消费者组 props.put("enable.auto.commit", config().isAutoAck()); props.put("auto.offset.reset", "earliest"); // 从最早的偏移量开始消费 props.put("key.deserializer", "org.apache.kafka.common.serialization.StringDeserializer"); props.put("value.deserializer", "org.apache.kafka.common.serialization.StringDeserializer"); KafkaConsumer<String, String> consumer = new KafkaConsumer<>(props); consumer.subscribe(Collections.singletonList(config().namespace("_") + mqListener.getTopic())); LISTENERS.add(() -> { GlobalThreadPool.submit(() -> { while (true) { ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(1000)); for (ConsumerRecord<String, String> record : records) { // 处理消息 // Kafka没有标签过滤机制,所以自己实现一遍 MQEvent mqEvent = serialization().deserialize(record.value(), mqListener.getEventClass()); if (!MQFilter.matchExps(mqEvent.getTag(), mqListener.getTags())) { continue; } try { consume(mqListener, mqEvent); if (!config().isAutoAck()) { consumer.commitSync(); // 手动提交偏移量 } } catch (Exception e) { log.error("Consume MQ failed: {}", mqEvent, e); if (e instanceof ServiceException) { if (!config().isAutoAck()) { consumer.commitSync(); // 手动提交偏移量 } } } } } }); return consumer; }); }
@Override public void start() { if (!config().isEnable() || !Objects.equals(config().getImpl(), impl())) return; LISTENERS.forEach(Supplier::get); }
@Override public void destroy() throws Exception { for (Supplier<KafkaConsumer> listener : LISTENERS) { listener.get().close(); } }}

三、Redis实现

Redis严格来说不算MQ,但因为它是一个独立于应用之外的带发布/订阅功能的中间件,且算它半个MQ吧,一般缓存服务用得多,如果要考虑减少服务器资源开销,可以考虑资源重复利用:

/** * Redis客户端实现 */@Slf4j(topic = "### BASE-MQ : redisClient ###")@Componentpublic final class RedisClient implements MQClient, DisposableBean {    private static JedisPool JEDIS_POOL;    private static List<Supplier> LISTENERS = new ArrayList<>();
@Override public String impl() { return "redis"; }
@Override public Consumer<MQEvent> initProducer() { return mqEvent -> { String message = serialization().serialize(mqEvent); log.info("Publish MQ: {}", message); // channel=namespace:topic:tag或namespace:topic或topic String channel = String.format("%s%s%s", (config().getNamespace().isEmpty() ? "" : config().getNamespace() + ":"), mqEvent.getTopic(), (mqEvent.getTag() == null || mqEvent.getTag().isEmpty()) ? "" : (":" + mqEvent.getTag())); GlobalThreadPool.submit(() -> { try { jedis().publish(channel, message); } catch (Throwable e) { log.error("Publish MQ to [{}] failed: {}", channel, message, e); } }); }; }
@Override public void initConsumer(MQListener mqListener) throws Exception { // channel=namespace:topic:tag或namespace:topic或topic,TODO 暂不支持通配符*和非- List<String> channels = new ArrayList<>(); if (mqListener.getTags() != null && !mqListener.getTags().isEmpty()) { Set<String> tags = MQFilter.findIncludes(mqListener.getTags()); if (!tags.isEmpty()) { for (String tag : tags) { channels.add((config().getNamespace().isEmpty() ? "" : config().getNamespace() + ":") + mqListener.getTopic() + (":" + tag)); } } } else { channels.add((config().getNamespace().isEmpty() ? "" : config().getNamespace() + ":") + mqListener.getTopic()); } // Redis原子加锁脚本 String lockScript = "if redis.call('get', KEYS[1]) == false then redis.call('setex', KEYS[1], tonumber(ARGV[1]), KEYS[1]) return 1 else return 0 end"; LISTENERS.add(() -> { GlobalThreadPool.submit(() -> { try (Jedis jedis = jedis()) { jedis.subscribe(new JedisPubSub() { @Override public void onMessage(String channel, String message) { MQEvent event = null; try { event = serialization().deserialize(message, mqListener.getEventClass()); //按组消费:Redis本身没有消费组的概念,这里我们锁住某条消息的消费资格10秒钟,以实现同消费组(如多实例)只能由一个消费者消费 String lockKey = String.format("ChannelGroupLock:%s:%s:%s", channel, mqListener.getGroup(), event.getMsgId()); if ((long) jedis().eval(lockScript, 1, lockKey, "10000") == 1) { //本条消息在该组未消费 consume(mqListener, event); } } catch (Exception e) { log.error("Consume MQ failed: {}", event, e); } }
@Override public void onSubscribe(String channel, int subscribedChannels) { log.info("Subscribed channel: {}", channel); }
@Override public void onUnsubscribe(String channel, int subscribedChannels) { log.info("Unsubscribed channel: {}", channel); } }, channels.toArray(new String[0])); } }); return null; }); }
@Override public void start() { if (!config().isEnable() || !Objects.equals(config().getImpl(), impl())) return; LISTENERS.forEach(Supplier::get); }
@Override public void destroy() throws Exception { if (JEDIS_POOL != null) { JEDIS_POOL.close(); } }
private Jedis jedis() { if (JEDIS_POOL == null) { String[] hostAndPort = config().getServer().split(":"); JedisPoolConfig poolConfig = new JedisPoolConfig(); poolConfig.setMaxTotal(100); poolConfig.setMaxIdle(50); poolConfig.setMinIdle(10); poolConfig.setTestOnBorrow(true); if (!config().getPassword().isEmpty()) { JEDIS_POOL = new JedisPool(poolConfig, hostAndPort[0], Integer.parseInt(hostAndPort[1]), 5000, config().getPassword()); } else { JEDIS_POOL = new JedisPool(poolConfig, hostAndPort[0], Integer.parseInt(hostAndPort[1])); } } return JEDIS_POOL.getResource(); }
}

之前的底层依赖使用了spring-data-redis,但我发现很多工程中RedisTemplate版本冲突太多了,一气之下我换成更轻量的jedis(spring-data-redis底层也是用jedis)。

四、RabbitMQ实现

RabbitMQ也是常用的MQ之一,我们需要考虑兼容它的AMQP协议,以及标签Tags的实现,这里用了它默认的Direct交换机:

@Component@Slf4j(topic = "### BASE-MQ : rabbitClient ###")public final class RabbitClient implements MQClient, DisposableBean {    private static Connection CONNECTION;    private static List<Supplier> LISTENERS = new ArrayList<>();
@Override public String impl() { return "rabbit"; }
@Override public Consumer<MQEvent> initProducer() { try { com.rabbitmq.client.Channel channel = connection().createChannel(); return mqEvent -> { // 定义具体的MQ事件发布逻辑 String message = serialization().serialize(mqEvent); log.info("Publish MQ: {}", message); try { // 路由键=namespace.topic或namespace.topic.tag String routingKey = config().namespace(".") + mqEvent.getTopic() + (mqEvent.getTag() == null || mqEvent.getTag().isEmpty() ? "" : ("." + mqEvent.getTag())); channel.basicPublish(config().getExchange(), routingKey, null, message.getBytes()); } catch (Exception e) { log.error("Publish MQ: {} failed!", message, e); } }; } catch (IOException e) { throw new RuntimeException(e); } }
@Override public void initConsumer(MQListener mqListener) throws Exception { // 队列名=group.namespace.topic.className.methodName String queue = mqListener.group(".") + config().namespace(".") + mqListener.getListenerMethod().getDeclaringClass().getSimpleName() + "." + mqListener.getListenerMethod().getName(); List<String> routingKeys = new ArrayList<>(); if (mqListener.getTags() != null && !mqListener.getTags().isEmpty()) { // 通过多个tags绑定到对应的路由键 Set<String> tags = MQFilter.findIncludes(mqListener.getTags()); if (!tags.isEmpty()) { // 路由键=namespace.topic.tag1、namespace.topic.tag2,TODO 这里只能用或||的关系,通配符*和非-是不能生效的 for (String tag : tags) { String routingKey = config().namespace(".") + mqListener.getTopic() + "." + tag; routingKeys.add(routingKey); } } } else { // 路由键=namespace.topic String routingKey = config().namespace(".") + mqListener.getTopic(); routingKeys.add(routingKey); } try (com.rabbitmq.client.Channel channel = connection().createChannel()) { channel.queueDeclare(queue, true, false, false, null); // 队列绑定到多个路由器 for (String routingKey : routingKeys) { channel.queueBind(queue, config().getExchange(), routingKey); } DeliverCallback deliverCallback = (consumerTag, delivery) -> { String message = new String(delivery.getBody(), StandardCharsets.UTF_8); MQEvent event = serialization().deserialize(message, mqListener.getEventClass()); if (!MQFilter.matchExps(event.getTag(), mqListener.getTags())) { return; } try { consume(mqListener, event); if (!config().isAutoAck()) { channel.basicAck(delivery.getEnvelope().getDeliveryTag(), false); } } catch (Exception e) { log.error("Consume MQ failed: {}", event, e); if (e instanceof ServiceException) { if (!config().isAutoAck()) { channel.basicAck(delivery.getEnvelope().getDeliveryTag(), false); } } } }; LISTENERS.add(() -> GlobalThreadPool.submit(() -> { try { channel.basicConsume(queue, config().isAutoAck(), deliverCallback, consumerTag -> { }); } catch (Exception e) { } return channel; }) ); } }
@Override public void start() { if (!config().isEnable() || !Objects.equals(config().getImpl(), impl())) return; LISTENERS.forEach(Supplier::get); }
@Override public void destroy() throws Exception { try { for (Supplier<com.rabbitmq.client.Channel> listener : LISTENERS) { listener.get().close(); } } catch (Exception e) { log.error("Error closing RabbitMQ channel", e); } try { connection().close(); } catch (Exception e) { log.error("Error closing RabbitMQ connection", e); } }
private Connection connection() { if (CONNECTION == null) { ConnectionFactory factory = new ConnectionFactory(); String[] hostAndPort = config().getServer().split(":"); factory.setHost(hostAndPort[0]); factory.setPort(Integer.parseInt(hostAndPort[1])); factory.setUsername(config().getUsername()); factory.setPassword(config().getPassword()); try { CONNECTION = factory.newConnection(); } catch (IOException | TimeoutException e) { throw new RuntimeException(e); } } return CONNECTION; }}

0x4

基础MQ组件整体划分

以下是所有代码写完后的整体结构:

Image

  • config包放MQ组件的配置类

  • core包放MQ组件的核心类:包括MQ客户端接口、MQ事件序列化器、MQ事件存储器、MQ过滤器、MQ监听器

  • impl包放MQ组件的实现类:包括默认使用的JsonMQ事件序列化器、不同MQ的客户端实现

此外,还有两个类是不放在MQ组件的,上移到了base-core核心组件(MQ组件依赖核心组件),目的是为了在引入MQ时更轻量,即使后期业务工程要去掉base-mq依赖,业务也不会报错(当然了,MQ事件发布订阅会失效):

Image

  • MQEventListener注解:参考Spring事件机制的@EventListener来实现

  • MQEvent:MQ事件基类,与领域事件基类DomainEvent同级,具体业务的MQ事件承继它即可

MQ组件的全貌就这些,其实非常简单,只需要定义完核心组件,剩下的在impl包去做不同的MQ实现即可。后续想要自己扩展别的MQ实现,也只需要实现MQClient接口并注入到Spring容器中,在配置类里指定该实现就行了。

0x5

写在最后

工欲善其事,必先利其器,咱们总不能一个系统就写一个工具类吧,所以抽一个公共的MQ组件是非常有必要的。有人说,为什么不直接用SpringStream,还要自己费工夫写一个?

我只想说,虽然Spring集成了很多内容,生态很大,但同时也有一些历史包袱,引入过多的设计模式反而是一种“反模式”,封装得多也会变得很重,学习成本大大增加,框架想定制自己的业务代码也变得异常困难,而我对基础框架的看法一直都是:简洁、灵活、包容。

对于文中的完整代码,我把它集成在了我的D3Boot(DDD快速启动)开源基础框架里的BASE-MQ组件,大家需要可以移步Gitee抄作业。

Gitee源码地址(差1个ʷ到10k Stars):
https://gitee.com/jensvn/d3boot

EOF

Image

作者:Jensen

专注分享程序员日常/架构技术/职场干货

Java老兵一枚,深耕电商/IoT/康美/AI等领域产品研发多年。

D3Boot开源框架作者,现任职SaaS软件架构师。

关注回复“DDD”,免费领取DDD学习大礼包。

关注回复“进群”,我拉你进架构技术交流群。

点赞在看ᴸᵉᵗᴹᵉ年薪百万