可暂停与恢复的无限任务型生产者-消费者模式设计与实现

星明吖_9662

星明吖_9662

2026-08-06

825人浏览

原创

可暂停与恢复的无限任务型生产者-消费者模式设计与实现

本文介绍如何扩展经典生产者-消费者模型,支持无限长度但可主动暂停/恢复的任务(如流式字符串处理),通过状态化任务封装、协作式调度和双角色线程池,实现高效、公平、可中断的并发任务分发与续执行。

本文介绍如何扩展经典生产者-消费者模型,支持无限长度但可主动暂停/恢复的任务(如流式字符串处理),通过状态化任务封装、协作式调度和双角色线程池,实现高效、公平、可中断的并发任务分发与续执行。

在实际系统中,许多任务并非“一次性完成”的原子操作——例如实时日志关键词统计、长连接数据流解析、或分布式爬虫的URL队列处理。这类任务具有无限性(数据源持续到达)、可暂停性(需让出CPU以响应更高优先级任务或负载均衡)和可恢复性(从中断点精确续算)。此时,传统基于 BlockingQueue 的单向生产者→消费者模型不再适用:消费者处理中途需将未完成任务“退回”队列,自身又临时充当生产者,形成双向任务流转。

核心设计原则

  1. 任务状态封装:避免直接传递原始 Queue,而是用 JobStatus 包装任务上下文,包含:

    • 待处理的数据源(如 Iterator 或 Stream)
    • 运行时状态(已统计词频 Map、当前偏移量、处理时间戳等)
    • 控制参数(如 maxItemsPerSlice = 10_000,触发暂停的阈值)
  2. 协作式暂停机制:每个工作线程在处理单个任务时,按预设策略主动让渡控制权,而非依赖外部中断(避免破坏状态一致性)。典型策略包括:

    • 按处理项数暂停(如每处理 1 万条后暂停)
    • 按耗时暂停(如单次 slice 超过 50ms)
    • 按队列水位动态调整(若待处理任务数 > 线程数 × 2,则加速切片)
  3. 统一任务队列 + 终止信号:使用 BlockingQueue 作为共享中枢,配合全局 END_MARKER 对象实现优雅关闭。

实现示例(Java)

// 任务状态封装类
public class JobStatus {
    private final Iterator<string> stream;
    private final Map<string integer> wordCount = new HashMap();
    private final long startTime;
    private final int maxItemsPerSlice;

    public JobStatus(Iterator<string> stream, int maxItemsPerSlice) {
        this.stream = stream;
        this.maxItemsPerSlice = maxItemsPerSlice;
        this.startTime = System.nanoTime();
    }

    // 执行一个处理切片,返回是否已完成
    public boolean processSlice() {
        int processed = 0;
        while (stream.hasNext() && processed  getSnapshot() {
        return new HashMap(wordCount);
    }
}

// 工作线程实现
public class WorkerThread implements Runnable {
    private final BlockingQueue<jobstatus> workQueue;
    private static final JobStatus END_MARKER = new JobStatus(
        Collections.emptyIterator(), 0);

    public WorkerThread(BlockingQueue<jobstatus> workQueue) {
        this.workQueue = workQueue;
    }

    @Override
    public void run() {
        try {
            while (true) {
                JobStatus job = workQueue.take();
                if (job == END_MARKER) {
                    workQueue.put(END_MARKER); // 广播终止信号
                    break;
                }

                boolean completed = job.processSlice();
                if (!completed) {
                    // 未完成 → 放回队尾,实现轮转调度
                    workQueue.put(job);
                }
                // 可选:添加延迟避免忙等待(如队列空闲时)
                if (workQueue.isEmpty()) {
                    Thread.sleep(1);
                }
            }
        } catch (InterruptedException e) {
            Thread.currentThread().interrupt();
        }
    }
}

// 启动与使用
public class StreamProcessor {
    public static void main(String[] args) throws InterruptedException {
        BlockingQueue<jobstatus> queue = new LinkedBlockingQueue();
        ExecutorService pool = Executors.newFixedThreadPool(4);

        // 提交多个流任务(模拟不同数据源)
        List<iterator>> streams = generateTestStreams();
        for (Iterator<string> stream : streams) {
            queue.offer(new JobStatus(stream, 10_000));
        }

        // 启动工作线程
        for (int i = 0; i <h3>关键注意事项</h3>
<ul>
<li>✅ <strong>状态一致性</strong>:JobStatus 必须是线程安全的(本例中仅由单一线程修改,故无需同步;若需跨线程读取快照,应加 synchronized 或使用 ConcurrentHashMap)。</li>
<li>⚠️ <strong>避免虚假唤醒</strong>:BlockingQueue.take() 已处理中断,但需在 catch (InterruptedException) 中恢复中断状态(Thread.currentThread().interrupt())。</li>
<li>? <strong>禁止共享可变集合</strong>:不要将 ArrayList 或普通 HashMap 直接暴露给多线程,否则会导致 ConcurrentModificationException 或数据丢失。</li>
<li>? <strong>暂停位置选择</strong>:应在自然边界暂停(如处理完一条完整日志行),而非在循环中间强制打断,防止状态残缺。</li>
<li>? <strong>监控与调优</strong>:建议记录每个 JobStatus 的处理时长、切片次数、最终结果大小,用于动态调整 maxItemsPerSlice 参数。</li>
</ul>
<p>该模式本质上是一种轻量级协程调度思想在 JVM 线程模型中的落地——它不依赖语言级协程(如 Kotlin suspend),而是通过任务状态显式保存 + 队列重入,达成近似协作式多任务的效果。适用于中高吞吐、低延迟敏感的流式数据处理场景,是传统生产者-消费者模式面向真实业务复杂性的必要演进。</p></string></iterator></jobstatus></jobstatus></jobstatus></string></string></string>
PHP速学视频免费教程(入门到精通)
PHP速学视频免费教程(入门到精通)

PHP怎么学习?PHP怎么入门?PHP在哪学?PHP怎么学才快?不用担心,这里为大家提供了PHP速学教程(入门到精通),有需要的小伙伴保存下载就能学习啦!

下载

相关标签:

本站声明:本文内容由网友自发贡献,版权归原作者所有,本站不承担相应法律责任。如您发现有涉嫌抄袭侵权的内容,请联系admin@php.cn

相关专题

更多
java
java

Java是一个通用术语,用于表示Java软件及其组件,包括“Java运行时环境 (JRE)”、“Java虚拟机 (JVM)”以及“插件”。php中文网还为大家带了Java相关下载资源、相关课程以及相关文章等内容,供大家免费下载使用。

2023.06.15

9417

6

java正则表达式语法
java正则表达式语法

java正则表达式语法是一种模式匹配工具,它非常有用,可以在处理文本和字符串时快速地查找、替换、验证和提取特定的模式和数据。本专题提供java正则表达式语法的相关文章、下载和专题,供大家免费下载体验。

2023.07.05

6582

9

java自学难吗
java自学难吗

Java自学并不难。Java语言相对于其他一些编程语言而言,有着较为简洁和易读的语法,本专题为大家提供java自学难吗相关的文章,大家可以免费体验。

2023.07.31

5852

8

java配置jdk环境变量
java配置jdk环境变量

Java是一种广泛使用的高级编程语言,用于开发各种类型的应用程序。为了能够在计算机上正确运行和编译Java代码,需要正确配置Java Development Kit(JDK)环境变量。php中文网给大家带来了相关的教程以及文章,欢迎大家前来阅读学习。

2023.08.01

1024

3

java保留两位小数
java保留两位小数

Java是一种广泛应用于编程领域的高级编程语言。在Java中,保留两位小数是指在进行数值计算或输出时,限制小数部分只有两位有效数字,并将多余的位数进行四舍五入或截取。php中文网给大家带来了相关的教程以及文章,欢迎大家前来阅读学习。

2023.08.02

868

3

java基本数据类型
java基本数据类型

java基本数据类型有:1、byte;2、short;3、int;4、long;5、float;6、double;7、char;8、boolean。本专题为大家提供java基本数据类型的相关的文章、下载、课程内容,供大家免费下载体验。

2023.08.02

1236

5

java有什么用
java有什么用

java可以开发应用程序、移动应用、Web应用、企业级应用、嵌入式系统等方面。本专题为大家提供java有什么用的相关的文章、下载、课程内容,供大家免费下载体验。

2023.08.02

2469

5

java在线网站
java在线网站

Java在线网站是指提供Java编程学习、实践和交流平台的网络服务。近年来,随着Java语言在软件开发领域的广泛应用,越来越多的人对Java编程感兴趣,并希望能够通过在线网站来学习和提高自己的Java编程技能。php中文网给大家带来了相关的视频、教程以及文章,欢迎大家前来学习阅读和下载。

2023.08.03

19811

3

配置java环境变量
配置java环境变量

配置Java环境变量是为了让操作系统能够识别和使用Java的相关命令和功能。本专题为大家提供配置java环境变量相关文章,帮助大家解决问题。

2023.08.03

1115

8

热门下载

更多
网站特效
/
网站源码
/
网站素材
/
前端模板

精品课程

更多
热门推荐
/
最新课程
phpStudy极速入门视频教程
phpStudy极速入门视频教程

共6课时 | 54.6万人学习

独孤九贱(4)_PHP视频教程
独孤九贱(4)_PHP视频教程

共89课时 | 133.4万人学习