当前位置:首页 > 技术 > 正文内容

Java消息队列的并发消费策略

访客 技术 2026年9月7日 6

在现代分布式系统中,消息队列扮演着至关重要的角色,它们能够解耦服务、削峰填谷、实现最终一致性等。然而,随着业务量的增长,如何高效、快速地消费消息成为一个关键挑战。当消息生产速度远超单个线程的消费能力时,采用多线程并发消费机制就显得尤为必要。本文将探讨一种基于Java的多线程消息消费模型,旨在提高消息处理的吞吐量,并提供一个灵活的框架来管理不同类型的消息处理任务。

核心设计理念

为了构建一个高效的消息消费者,我们遵循以下设计原则:

  1. 持续拉取: 消费者应不间断地从消息中间件(例如 Apache RocketMQ、Kafka 或 RabbitMQ)拉取新消息。
  2. 批量处理与分片: 一次性拉取多条消息,并将这些消息进一步分成更小的批次(分片)。
  3. 并发执行: 利用专门的线程池,为每个消息分片分配一个工作线程进行并行处理。
  4. 同步等待: 等待当前批次的所有消息分片处理完成后,再进行下一轮的消息拉取,确保消息的有序处理或批次完整性。
  5. 任务隔离: 针对不同类型的业务消息(例如短信通知、邮件发送),使用独立的处理器和线程池,避免相互影响,提高系统的健壮性和可扩展性。
  6. 优雅停机: 在系统关闭时,确保正在处理的消息批次能够完整执行完毕,然后平稳关闭相关线程池。

代码实现

我们将通过一系列Java类来构建这个并发消息消费框架。为了简化演示,消息源将通过模拟方式提供,而不是直接连接到真实的消息队列。

消息模型:ConsumedMessage

首先,定义一个简单的消息数据结构。

package com.example.mqconsumer.model;

import java.time.Instant;
import java.util.UUID;

/**
 * 表示从消息队列中消费到的消息实体。
 */
public class ConsumedMessage {
    private String messageId;
    private String payload; // 消息内容
    private String messageType; // 消息类型,用于区分不同的业务处理逻辑
    private Instant creationTimestamp;

    public ConsumedMessage(String payload, String messageType) {
        this.messageId = UUID.randomUUID().toString();
        this.payload = payload;
        this.messageType = messageType;
        this.creationTimestamp = Instant.now();
    }

    public String getMessageId() { return messageId; }
    public String getPayload() { return payload; }
    public String getMessageType() { return messageType; }
    public Instant getCreationTimestamp() { return creationTimestamp; }

    @Override
    public String toString() {
        return "ConsumedMessage{" +
               "messageId='" + messageId + '\'' +
               ", payload='" + payload + '\'' +
               ", messageType='" + messageType + '\'' +
               ", creationTimestamp=" + creationTimestamp +
               '}';
    }
}

模拟消息源:MockMessageBroker

这个类模拟了从消息队列中拉取消息和发送确认(ACK)的逻辑。

package com.example.mqconsumer.mock;

import com.example.mqconsumer.model.ConsumedMessage;
import java.util.ArrayList;
import java.util.Collections;
import java.util.List;
import java.util.concurrent.ThreadLocalRandom;
import java.util.concurrent.atomic.AtomicInteger;

/**
 * 模拟消息队列代理,用于生成和确认消息。
 */
public class MockMessageBroker {
    private String topicName;
    private String messageTag;
    private AtomicInteger messageCounter = new AtomicInteger(0);

    public MockMessageBroker(String topicName, String messageTag) {
        this.topicName = topicName;
        this.messageTag = messageTag;
    }

    /**
     * 模拟从消息队列拉取一批消息。
     * @param requestedCount 请求拉取的消息数量。
     * @return 消息列表,可能为空。
     */
    public List<ConsumedMessage> fetchMessages(int requestedCount) {
        // 模拟偶尔拉取不到消息的情况
        if (ThreadLocalRandom.current().nextInt(10) < 2) {
            System.out.printf("[%s] MockBroker: 暂时没有消息可拉取。\n", topicName);
            return Collections.emptyList();
        }

        List<ConsumedMessage> messages = new ArrayList<>(requestedCount);
        for (int i = 0; i < requestedCount; i++) {
            String content = String.format("%s-%s-Msg%d-%d", topicName, messageTag, messageCounter.getAndIncrement(), System.currentTimeMillis());
            String type = topicName.contains("sms") ? "SMS" : (topicName.contains("email") ? "EMAIL" : "UNKNOWN");
            messages.add(new ConsumedMessage(content, type));
        }
        System.out.printf("[%s] MockBroker: 成功拉取 %d 条消息。\n", topicName, messages.size());
        return messages;
    }

    /**
     * 模拟向消息队列发送消息确认。
     * @param message 已处理的消息。
     */
    public void acknowledgeMessage(ConsumedMessage message) {
        // 实际场景中,这里会调用消息队列客户端的ack方法
        System.out.printf("[%s] MockBroker: 确认消息ID: %s\n", topicName, message.getMessageId());
    }
}

线程池管理:DynamicThreadPoolManager

这个工具类负责根据任务类型创建和管理专属的线程池。它确保了不同任务之间线程资源的隔离,并且支持线程池的优雅关闭。

package com.example.mqconsumer.thread;

import com.google.common.util.concurrent.ThreadFactoryBuilder;
import java.util.Map;
import java.util.concurrent.*;

/**
 * 动态线程池管理器,用于为不同类型的任务提供和管理专属线程池。
 */
public class DynamicThreadPoolManager {

    private static final Map<String, ExecutorService> MANAGED_THREAD_POOLS = new ConcurrentHashMap<>();

    private DynamicThreadPoolManager() {
        // 防止实例化
    }

    /**
     * 获取或创建一个指定名称和核心线程数的线程池。
     * @param poolName 线程池的唯一名称。
     * @param corePoolSize 线程池的核心线程数。
     * @return 对应的 ExecutorService 实例。
     */
    public static ExecutorService getOrCreateThreadPool(String poolName, int corePoolSize) {
        // 使用 computeIfAbsent 确保线程池只被创建一次
        return MANAGED_THREAD_POOLS.computeIfAbsent(poolName, k -> {
            System.out.println("创建新的线程池: " + poolName + ",核心线程数: " + corePoolSize);
            return new ThreadPoolExecutor(
                corePoolSize,
                corePoolSize, // 最大线程数与核心线程数相同,避免动态扩容
                60L, TimeUnit.SECONDS, // 线程空闲时间
                new LinkedBlockingQueue<>(), // 无界任务队列
                new ThreadFactoryBuilder().setNameFormat(poolName + "-worker-%d").build(), // 自定义线程命名
                new ThreadPoolExecutor.CallerRunsPolicy() // 拒绝策略:由调用者线程执行任务
            );
        });
    }

    /**
     * 优雅地关闭指定名称的线程池。
     * @param poolName 要关闭的线程池名称。
     */
    public static void shutdownThreadPool(String poolName) {
        ExecutorService executor = MANAGED_THREAD_POOLS.remove(poolName);
        if (executor != null) {
            System.out.println("尝试关闭线程池: " + poolName);
            executor.shutdown(); // 启动有序关闭
            try {
                // 等待正在执行的任务完成
                if (!executor.awaitTermination(30, TimeUnit.SECONDS)) {
                    executor.shutdownNow(); // 强制关闭
                    if (!executor.awaitTermination(10, TimeUnit.SECONDS))
                        System.err.println("线程池 " + poolName + " 未完全终止。");
                }
            } catch (InterruptedException ie) {
                executor.shutdownNow();
                Thread.currentThread().interrupt(); // 保留中断状态
            }
        }
    }
}

抽象消息处理器:AbstractMessageConsumer

这是所有具体消息处理器的基类,它定义了消息拉取、分片、并发处理和优雅停机的通用逻辑。

package com.example.mqconsumer.processor;

import com.example.mqconsumer.mock.MockMessageBroker;
import com.example.mqconsumer.model.ConsumedMessage;
import com.example.mqconsumer.thread.DynamicThreadPoolManager;
import java.util.ArrayList;
import java.util.Collections;
import java.util.List;
import java.util.concurrent.CountDownLatch;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.TimeUnit;

/**
 * 抽象消息消费者基类,提供消息拉取、分片、并发处理和终止的通用框架。
 */
public abstract class AbstractMessageConsumer implements Runnable {
    protected String consumerIdentifier; // 消费者实例的唯一标识
    protected int workerThreadCount; // 处理消息的子线程数
    protected volatile boolean shutdownRequested = false; // 终止信号
    protected int fetchBatchSize = 10; // 每次从消息源拉取的消息数量
    protected int partitionChunkSize = 2; // 每个工作线程处理的消息分片大小

    protected MockMessageBroker messageBroker; // 注入的消息源

    public AbstractMessageConsumer(String identifier, int threadCount, MockMessageBroker broker) {
        this.consumerIdentifier = identifier;
        this.workerThreadCount = threadCount;
        this.messageBroker = broker;
    }

    /**
     * 抽象方法:由具体的消费者实现消息处理逻辑。
     * @param message 待处理的单条消息。
     */
    protected abstract void processSingleMessage(ConsumedMessage message);

    @Override
    public void run() {
        // 获取或创建用于处理具体消息的线程池
        ExecutorService workerExecutor = DynamicThreadPoolManager.getOrCreateThreadPool(consumerIdentifier + "-workers", workerThreadCount);

        while (!shutdownRequested) {
            List<ConsumedMessage> currentBatchMessages = messageBroker.fetchMessages(fetchBatchSize);

            if (currentBatchMessages.isEmpty()) {
                try {
                    System.out.printf("[%s] 暂无消息,等待 5 秒后重试...\n", consumerIdentifier);
                    TimeUnit.SECONDS.sleep(5); // 如果没有消息,等待一段时间再拉取
                } catch (InterruptedException e) {
                    Thread.currentThread().interrupt();
                    System.err.printf("[%s] 主拉取线程在等待时被中断:%s\n", consumerIdentifier, e.getMessage());
                    shutdownRequested = true; // 收到中断信号,请求终止
                }
                continue;
            }

            // 将拉取到的消息进行分片处理
            List<List<ConsumedMessage>> partitionedMessageLists = partitionList(currentBatchMessages, partitionChunkSize);
            CountDownLatch batchCompletionLatch = new CountDownLatch(partitionedMessageLists.size());

            for (List<ConsumedMessage> subList : partitionedMessageLists) {
                workerExecutor.execute(() -> {
                    try {
                        for (ConsumedMessage msg : subList) {
                            processSingleMessage(msg); // 处理单条消息
                            messageBroker.acknowledgeMessage(msg); // 确认消息
                        }
                    } catch (Exception e) {
                        System.err.printf("[%s] 处理消息分片时发生错误:%s - %s\n", consumerIdentifier, subList, e.getMessage());
                        // 实际系统中,这里可能需要NACK消息或记录错误日志以便后续重试
                    } finally {
                        batchCompletionLatch.countDown(); // 当前分片处理完成,计数器减一
                    }
                });
            }

            try {
                batchCompletionLatch.await(); // 等待所有分片处理完成
                System.out.printf("[%s] 成功处理完一批 %d 条消息。\n", consumerIdentifier, currentBatchMessages.size());
            } catch (InterruptedException e) {
                Thread.currentThread().interrupt();
                System.err.printf("[%s] 主拉取线程在等待批处理完成时被中断:%s\n", consumerIdentifier, e.getMessage());
                shutdownRequested = true; // 收到中断信号,请求终止
            }
        }
        System.out.printf("[%s] 收到终止请求。正在关闭工作线程池...\n", consumerIdentifier);
        DynamicThreadPoolManager.shutdownThreadPool(consumerIdentifier + "-workers");
        System.out.printf("[%s] 消费者已停止。\n", consumerIdentifier);
    }

    /**
     * 发送终止信号,请求消费者优雅停机。
     */
    public void requestShutdown() {
        this.shutdownRequested = true;
        System.out.printf("[%s] 收到终止请求信号。\n", consumerIdentifier);
    }

    public String getConsumerIdentifier() {
        return consumerIdentifier;
    }

    /**
     * 辅助方法:将列表分割成多个小块。
     */
    private <T> List<List<T>> partitionList(List<T> list, int chunkSize) {
        if (list == null || list.isEmpty() || chunkSize <= 0) {
            return Collections.emptyList();
        }
        int numPartitions = (int) Math.ceil((double) list.size() / chunkSize);
        List<List<T>> partitions = new ArrayList<>(numPartitions);
        for (int i = 0; i < numPartitions; i++) {
            int fromIndex = i * chunkSize;
            int toIndex = Math.min((i + 1) * chunkSize, list.size());
            partitions.add(new ArrayList<>(list.subList(fromIndex, toIndex)));
        }
        return partitions;
    }
}

具体消息处理器:SmsNotificationConsumerEmailNotificationConsumer

这些类继承自 AbstractMessageConsumer,并实现各自的消息处理逻辑。

package com.example.mqconsumer.processor;

import com.example.mqconsumer.mock.MockMessageBroker;
import com.example.mqconsumer.model.ConsumedMessage;
import java.util.concurrent.TimeUnit;

/**
 * 短信通知消息的消费者实现。
 */
public class SmsNotificationConsumer extends AbstractMessageConsumer {
    public SmsNotificationConsumer(String identifier, int threadCount, MockMessageBroker broker) {
        super(identifier, threadCount, broker);
    }

    @Override
    protected void processSingleMessage(ConsumedMessage message) {
        // 模拟短信发送业务逻辑
        System.out.printf("[%s] %s [SMS] 处理消息: %s\n", Thread.currentThread().getName(), consumerIdentifier, message.getPayload());
        try {
            TimeUnit.MILLISECONDS.sleep(35); // 模拟耗时操作
        } catch (InterruptedException e) {
            Thread.currentThread().interrupt();
            System.err.printf("[%s] [SMS] 消息处理中断:%s\n", Thread.currentThread().getName(), e.getMessage());
        }
    }
}
package com.example.mqconsumer.processor;

import com.example.mqconsumer.mock.MockMessageBroker;
import com.example.mqconsumer.model.ConsumedMessage;
import java.util.concurrent.TimeUnit;

/**
 * 邮件通知消息的消费者实现。
 */
public class EmailNotificationConsumer extends AbstractMessageConsumer {
    public EmailNotificationConsumer(String identifier, int threadCount, MockMessageBroker broker) {
        super(identifier, threadCount, broker);
    }

    @Override
    protected void processSingleMessage(ConsumedMessage message) {
        // 模拟邮件发送业务逻辑
        System.out.printf("[%s] %s [EMAIL] 处理消息: %s\n", Thread.currentThread().getName(), consumerIdentifier, message.getPayload());
        try {
            TimeUnit.MILLISECONDS.sleep(20); // 模拟耗时操作,比短信快一些
        } catch (InterruptedException e) {
            Thread.currentThread().interrupt();
            System.err.printf("[%s] [EMAIL] 消息处理中断:%s\n", Thread.currentThread().getName(), e.getMessage());
        }
    }
}

启动器:ConsumerApplication

这是整个应用的入口点,负责初始化并启动不同类型的消息消费者,并演示如何优雅地关闭它们。

package com.example.mqconsumer;

import com.example.mqconsumer.mock.MockMessageBroker;
import com.example.mqconsumer.processor.AbstractMessageConsumer;
import com.example.mqconsumer.processor.EmailNotificationConsumer;
import com.example.mqconsumer.processor.SmsNotificationConsumer;

import java.util.ArrayList;
import java.util.List;
import java.util.concurrent.TimeUnit;

/**
 * 消费者应用启动类,用于初始化并启动多种类型的消息消费者。
 */
public class ConsumerApplication {

    public static void main(String[] args) throws InterruptedException {
        System.out.println("====== 启动多线程消息消费应用 ======");

        List<AbstractMessageConsumer> activeConsumers = new ArrayList<>();

        // 1. 初始化短信通知消息消费者
        MockMessageBroker smsBroker = new MockMessageBroker("smsTopic", "marketing");
        SmsNotificationConsumer smsConsumer = new SmsNotificationConsumer("SMS_Processor_01", 2, smsBroker);
        activeConsumers.add(smsConsumer);

        // 2. 初始化邮件通知消息消费者
        MockMessageBroker emailBroker = new MockMessageBroker("emailTopic", "transactional");
        EmailNotificationConsumer emailConsumer = new EmailNotificationConsumer("EMAIL_Processor_01", 3, emailBroker);
        activeConsumers.add(emailConsumer);

        // 为每个消费者创建一个独立的线程并启动
        for (AbstractMessageConsumer consumer : activeConsumers) {
            Thread consumerMainThread = new Thread(consumer, consumer.getConsumerIdentifier() + "-main-loop");
            consumerMainThread.start();
        }

        System.out.println("\n--- 消费者已启动,运行中... ---\n");

        // 模拟应用运行一段时间
        TimeUnit.SECONDS.sleep(5); // 增加运行时间以观察更多输出

        System.out.println("\n--- 收到停止应用的请求,开始优雅关闭所有消费者 ---\n");
        // 发送终止信号给所有消费者
        for (AbstractMessageConsumer consumer : activeConsumers) {
            consumer.requestShutdown();
        }

        // 等待一段时间,让消费者完成当前批次的消息处理并关闭线程池
        // 在实际应用中,可能会使用 Thread.join() 来等待这些线程真正结束
        TimeUnit.SECONDS.sleep(8);
        System.out.println("====== 应用关闭完成 ======");
    }
}

运行效果演示

执行 ConsumerApplicationmain 方法后,你将看到类似以下的控制台输出。输出会根据消息拉取和处理速度有所不同,但核心流程保持一致。

====== 启动多线程消息消费应用 ======
创建新的线程池: SMS_Processor_01-workers,核心线程数: 2
创建新的线程池: EMAIL_Processor_01-workers,核心线程数: 3

--- 消费者已启动,运行中... ---

[smsTopic] MockBroker: 成功拉取 10 条消息。
[SMS_Processor_01] SMS_Processor_01-main-loop [SMS] 处理消息: smsTopic-marketing-Msg0-1701000000000
[emailTopic] MockBroker: 成功拉取 10 条消息。
[EMAIL_Processor_01] EMAIL_Processor_01-main-loop [EMAIL] 处理消息: emailTopic-transactional-Msg10-1701000000050
[EMAIL_Processor_01-workers-worker-0] EMAIL_Processor_01 [EMAIL] 处理消息: emailTopic-transactional-Msg12-1701000000050
[SMS_Processor_01-workers-worker-0] SMS_Processor_01 [SMS] 处理消息: smsTopic-marketing-Msg1-1701000000000
[EMAIL_Processor_01-workers-worker-1] EMAIL_Processor_01 [EMAIL] 处理消息: emailTopic-transactional-Msg10-1701000000050
[EMAIL_Processor_01-workers-worker-2] EMAIL_Processor_01 [EMAIL] 处理消息: emailTopic-transactional-Msg14-1701000000050
[smsTopic] MockBroker: 确认消息ID: 4d2e...
[SMS_Processor_01-workers-worker-0] SMS_Processor_01 [SMS] 处理消息: smsTopic-marketing-Msg2-1701000000000
[emailTopic] MockBroker: 确认消息ID: bc7f...
[EMAIL_Processor_01-workers-worker-0] EMAIL_Processor_01 [EMAIL] 处理消息: emailTopic-transactional-Msg11-1701000000050
... (更多消息处理日志) ...
[SMS_Processor_01] 成功处理完一批 10 条消息。
[smsTopic] MockBroker: 成功拉取 10 条消息。
[EMAIL_Processor_01] 成功处理完一批 10 条消息。
[emailTopic] MockBroker: 成功拉取 10 条消息。
...

--- 收到停止应用的请求,开始优雅关闭所有消费者 ---

[SMS_Processor_01] 收到终止请求信号。
[EMAIL_Processor_01] 收到终止请求信号。
[SMS_Processor_01] 收到终止请求。正在关闭工作线程池...
尝试关闭线程池: SMS_Processor_01-workers
[SMS_Processor_01] 消费者已停止。
[EMAIL_Processor_01] 收到终止请求。正在关闭工作线程池...
尝试关闭线程池: EMAIL_Processor_01-workers
[EMAIL_Processor_01] 消费者已停止。
====== 应用关闭完成 ======

从输出中可以看出,不同的消息处理器在各自的线程池中并行工作,持续拉取并消费消息。当程序收到终止信号时,它们会等待当前正在处理的批次完成后,再关闭各自的线程池,实现了优雅停机。

标签: Java

相关文章

Linux crontab 详解

1) crontab 是什么cron 是 Linux 的定时任务守护进程;crontab 是用来编辑/查看“按时间周期执行命令”的表(cron table)。常见两类:用户 crontab:每个用户一份(crontab -e 编辑)系统级 crontab / cron.d:可指定执行用户(/etc/crontab、/etc/cron.d/*)2) crontab 时间...

富文本里可以允许的 HTML 属性

一、所有标签默认允许的安全属性(极少)class        (可选)id           (通常建议禁用)title️ 注意:id 容易被滥用做锚点注入,很多系统直接禁用class 允许的话最好只允许固定前缀(如 editor-*)二、a 标签允许属性<a href="" t...

Mac 安装 Node.js 指南

方法一:通过官网安装包(最简单,适合初学者)如果你只是想快速安装并开始使用,这是最直接的方法。访问 Node.js 官网。页面会显示两个版本:LTS (Recommended For Most Users):长期支持版,最稳定。建议选这个。Current:最新特性版,包含最新功能但可能不够稳定。下载 .pkg 安装包并运行。按照安装向导点击“下一步”即可完成。方法二:使用 Homebrew 安装(...

Dom\HTML_NO_DEFAULT_NS 的副作用:自动加闭合标签

在使用Dom\HTMLDocument时,Dom\HTML_NO_DEFAULT_NS 将禁止在解析过程中设置元素的命名空间, 此设置是为了与DOMDocument向后兼容而存在的。当使用它时,已知的一个副作用就是:自动加闭合标签例如 </img> 为什么会这样?当你使用:Dom\HTML_NO_DEFAULT_NS文档会变成 无命名空间模式,此时内部更接近 XML...

Laravel 事件和监听器创建

在 Laravel 中,使用 Artisan 命令创建 Events(事件) 和 Listeners(监听器) 是非常高效的。你可以通过以下几种方式来实现:1. 手动创建单个 Event如果你只想创建一个事件类,可以使用 make:event 命令:Bashphp artisan make:event UserRegistered执行后,文件将生成在 app/Even...

自定义域名解析神器 dnsmasq

什么是 dnsmasq?dnsmasq 是一个轻量级、功能强大的网络服务工具,专为小型和中等规模网络设计。它是一个综合的网络基础设施解决方案[1]。dnsmasq 能做什么?功能说明应用场景DNS 转发与缓存将 DNS 查询转发到上游服务器(ISP、Google DNS 等),并在本地缓存结果加快 DNS 查询速度,减少外部 DNS 流量本地 DNS解析本地网络设备的主机名,无需编辑&n...

发表评论

访客

◎欢迎参与讨论,请在这里发表您的看法和观点。