首页 > 编程语言 >实现按属性分组的线程安全串行执行与全局并发控制

实现按属性分组的线程安全串行执行与全局并发控制

来源:互联网 2026-07-24 08:23:20

针对按属性分组的事件处理需求,提出一种分组队列与共享工作者池相结合的解耦架构。中央调度器将事件按颜色路由至对应队列,工作线程轮询非空队列,通过加锁标记实现同组串行、跨组并发,并利用固定大小线程池控制全局并发数,避免锁竞争与饥饿问题。

今天聊一个实际生产中经常遇到的并发调度问题:如何让同一类事件(比如根据“color”字段分组)严格按顺序一个一个执行,而不同类的事件又能同时并行处理,同时还要控制整个系统的并发线程数不超限。这其实是一个兼顾顺序性、隔离性与资源利用率的经典难题。

在高吞吐事件处理场景中,常常需要同时满足“同组串行、跨组并发、全局限流”这三个要求。举个例子,假设事件按 color 字段分组:所有绿色事件必须严格按照接收顺序依次执行,黄色事件也是如此;但绿色和黄色之间没有执行顺序的约束,可以同时跑;与此同时,整个系统里所有颜色的执行线程总数不能超过一个预设上限,比如 8 个。

直接去改造 ThreadPoolExecutor 的任务队列,比如自定义一个 BlockingQueue,试图在里面实现“跳过同色正在运行的任务”这种动态出队逻辑,不仅会破坏线程池原有的设计契约,还特别容易引发竞态条件、死锁甚至饥饿问题——比如某个颜色的事件持续积压,其他颜色长期得不到调度。因此,行业里更推荐的做法是采用一种解耦架构:分组队列 + 共享工作者池

长期稳定更新的攒劲资源: >>>点此立即查看<<<

核心设计:分组队列 + 统一调度器

  • 每个 color 对应一个线程安全队列(比如 ConcurrentLinkedQueue 或 LinkedBlockingQueue),保证该颜色内部的事件严格遵循 FIFO(先进先出)顺序;
  • 一个中央调度器线程(Distributor)持续从原始事件源(比如 Kafka、消息队列或生产者队列)读取事件,然后根据 event.color() 将其路由到对应的颜色队列里;
  • 一个固定大小的共享线程池(比如 Executors.newFixedThreadPool(N))负责消费所有颜色队列。每个工作线程会循环尝试从任意非空队列中取任务(优先取队列头部的旧事件),执行前加锁标记“该 color 正在运行”,执行完成后立即释放锁。

示例实现(Java)

// 1. 分组队列容器
private final ConcurrentMap> colorQueues = new ConcurrentHashMap<>();
private final ReentrantLock lock = new ReentrantLock();
private final Set runningColors = ConcurrentHashMap.newKeySet();

// 2. 工作线程任务(提交至共享线程池)
Runnable workerTask = () -> {
    while (!Thread.currentThread().isInterrupted()) {
        Event event = null;
        String color = null;
        // 轮询所有队列,找到首个可执行的 oldest 事件(避免饿死)
        for (Queue queue : colorQueues.values()) {
            if (!queue.isEmpty()) {
                event = queue.peek(); // 先看一眼,不移除
                if (event != null && !runningColors.contains(event.color())) {
                    color = event.color();
                    event = queue.poll(); // 确认后出队
                    break;
                }
            }
        }
        if (event == null) {
            Thread.sleep(10); // 短暂让出 CPU
            continue;
        }
        // 标记 color 正在运行
        runningColors.add(color);
        try {
            event.execute(); // 执行业务逻辑
        } finally {
            runningColors.remove(color); // 必须确保释放
        }
    }
};

// 启动 N 个 worker 线程
ExecutorService workers = Executors.newFixedThreadPool(8);
for (int i = 0; i < 8; i++) {
    workers.submit(workerTask);
}

关键注意事项

  • 避免锁竞争:runningColors 使用 ConcurrentHashMap.newKeySet() 替代 synchronized 块,能显著提升并发读写性能;
  • 防止任务丢失:peek() + poll() 这个组合要确保原子性。如果 poll() 返回 null(被其他线程抢先了),需要重试;
  • 公平性保障:轮询所有队列(而不是固定顺序)能缓解某些颜色长期积压的问题。进阶方案可以引入优先级队列,按队列头的时间戳排序;
  • 资源清理:空队列可以定期清理,比如 colorQueues.entrySet().removeIf(e -> e.getValue().isEmpty() && !runningColors.contains(e.getKey())),防止内存泄漏;
  • 扩展性:支持动态增加或销毁 color,无需重启服务。

这套方案天然就能满足所有原始需求:同色严格 FIFO、跨色完全并发、全局线程数可控,而且代码结构清晰,便于监控和调试。相比侵入式地修改线程池的队列,它更符合面向对象和关注点分离的原则,是生产环境里值得推荐的稳健做法。

侠游戏发布此文仅为了传递信息,不代表侠游戏网站认同其观点或证实其描述

热游推荐

更多
湘ICP备14008430号-1 湘公网安备 43070302000280号
All Rights Reserved
本站为非盈利网站,不接受任何广告。本站所有软件,都由网友
上传,如有侵犯你的版权,请发邮件给xiayx666@163.com
抵制不良色情、反动、暴力游戏。注意自我保护,谨防受骗上当。
适度游戏益脑,沉迷游戏伤身。合理安排时间,享受健康生活。