讲师中心 微信公众号
AI工具推荐 视频效率加速

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

星明吖_9662

星明吖_9662

发布时间:2026-08-06 21:46:02

|

825人浏览过

|

来源于php中文网

原创

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

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

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

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

核心设计原则

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

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

    • 按处理项数暂停(如每处理 1 万条后暂停)
    • 按耗时暂停(如单次 slice 超过 50ms)
    • 按队列水位动态调整(若待处理任务数 > 线程数 × 2,则加速切片)
  3. 统一任务队列 + 终止信号:使用 BlockingQueue<JobStatus> 作为共享中枢,配合全局 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 < maxItemsPerSlice) {
            String line = stream.next();
            // 示例:统计关键词 "error" 和 "warning"
            if (line.contains("error")) wordCount.merge("error", 1, Integer::sum);
            if (line.contains("warning")) wordCount.merge("warning", 1, Integer::sum);
            processed++;
        }
        return !stream.hasNext(); // true 表示任务彻底完成
    }

    // 获取当前状态快照(用于调试或监控)
    public Map<String, Integer> 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<String>> streams = generateTestStreams();
        for (Iterator<String> stream : streams) {
            queue.offer(new JobStatus(stream, 10_000));
        }

        // 启动工作线程
        for (int i = 0; i < 4; i++) {
            pool.submit(new WorkerThread(queue));
        }

        // 添加终止标记
        queue.offer(END_MARKER);

        pool.shutdown();
        pool.awaitTermination(1, TimeUnit.MINUTES);
    }
}

关键注意事项

  • ✅ 状态一致性:JobStatus 必须是线程安全的(本例中仅由单一线程修改,故无需同步;若需跨线程读取快照,应加 synchronized 或使用 ConcurrentHashMap)。
  • ⚠️ 避免虚假唤醒:BlockingQueue.take() 已处理中断,但需在 catch (InterruptedException) 中恢复中断状态(Thread.currentThread().interrupt())。
  • ? 禁止共享可变集合:不要将 ArrayList 或普通 HashMap 直接暴露给多线程,否则会导致 ConcurrentModificationException 或数据丢失。
  • ? 暂停位置选择:应在自然边界暂停(如处理完一条完整日志行),而非在循环中间强制打断,防止状态残缺。
  • ? 监控与调优:建议记录每个 JobStatus 的处理时长、切片次数、最终结果大小,用于动态调整 maxItemsPerSlice 参数。

该模式本质上是一种轻量级协程调度思想在 JVM 线程模型中的落地——它不依赖语言级协程(如 Kotlin suspend),而是通过任务状态显式保存 + 队列重入,达成近似协作式多任务的效果。适用于中高吞吐、低延迟敏感的流式数据处理场景,是传统生产者-消费者模式面向真实业务复杂性的必要演进。

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

热门AI工具

更多
Lovart
Lovart Hot

一款面向视觉设计创作的AI设计平台,可通过智能体和画布工作流辅助制作海报、Logo、网页、PPT及其他视觉内容。

咔片AIPPT

一款在线AI演示文稿制作工具,可根据主题和内容需求辅助生成PPT结构与页面,提高演示材料制作效率。

Atoms
Atoms Hot

Atoms是一款AI智能体工具,第一支自动构建真实业务的 AI 团队。

WorkBuddy

一款AI办公效率工具,主要用于腾讯云推出的AI原生桌面智能体工作台,适合需要提升相关任务效率的用户。

DeepSeek

DeepSeek是一款面向对话、写作、编程和推理场景的AI大模型工具。

豆包大模型

豆包大模型是一款由字节跳动推出的企业级大语言模型服务平台。

LibLibAI
LibLibAI Hot

一款AI视频创作工具,主要用于国内领先的AI创意平台,以海量模型、低门槛操作与“创作-分享-商业化”生态,让小白与专业创作者都能高效实现图文乃至视频创意表达,适合需要提升相关任务效率的用户。

立刻MV
立刻MV Hot

立刻MV是一款AI文本写作工具,AI 音乐视频(MV)创作工具。

AionClaw
AionClaw Hot

AionClaw是一款面向办公、创作和编程任务的AI桌面智能体。

相关专题

更多
java
java

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

9377

2023.06.15

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

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

6542

2023.07.05

java自学难吗
java自学难吗

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

5832

2023.07.31

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

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

1024

2023.08.01

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

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

868

2023.08.02

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

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

1216

2023.08.02

java有什么用
java有什么用

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

2469

2023.08.02

java在线网站
java在线网站

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

19811

2023.08.03

LLVM自定义Pass怎么写
LLVM自定义Pass怎么写

本专题聚焦LLVM自定义Pass开发,整理Pass类结构、run()方法、PreservedAnalyses、CMake构建、插件注册、-load-pass-plugin加载和测试用例编写流程。

0

2026.09.30

热门下载

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

精品课程

更多
热门推荐
/
最新课程
关于我们 免责申明 举报中心 意见反馈 讲师合作 广告合作 最新更新
php中文网:公益在线php培训,帮助PHP学习者快速成长!
关注服务号
PHP中文网订阅号
每天精选资源文章推送

Copyright 2014-2026 https://www.php.cn/ All Rights Reserved | php.cn