
本文介绍如何扩展经典生产者-消费者模型,支持无限长度但可主动暂停/恢复的任务(如流式字符串处理),通过状态化任务封装、协作式调度和双角色线程池,实现高效、公平、可中断的并发任务分发与续执行。
本文介绍如何扩展经典生产者-消费者模型,支持无限长度但可主动暂停/恢复的任务(如流式字符串处理),通过状态化任务封装、协作式调度和双角色线程池,实现高效、公平、可中断的并发任务分发与续执行。
在实际系统中,许多任务并非“一次性完成”的原子操作——例如实时日志关键词统计、长连接数据流解析、或分布式爬虫的URL队列处理。这类任务具有无限性(数据源持续到达)、可暂停性(需让出CPU以响应更高优先级任务或负载均衡)和可恢复性(从中断点精确续算)。此时,传统基于 BlockingQueue<T> 的单向生产者→消费者模型不再适用:消费者处理中途需将未完成任务“退回”队列,自身又临时充当生产者,形成双向任务流转。
核心设计原则
-
任务状态封装:避免直接传递原始 Queue<String>,而是用 JobStatus 包装任务上下文,包含:
- 待处理的数据源(如 Iterator<String> 或 Stream<String>)
- 运行时状态(已统计词频 Map<String, Integer>、当前偏移量、处理时间戳等)
- 控制参数(如 maxItemsPerSlice = 10_000,触发暂停的阈值)
-
协作式暂停机制:每个工作线程在处理单个任务时,按预设策略主动让渡控制权,而非依赖外部中断(避免破坏状态一致性)。典型策略包括:
- 按处理项数暂停(如每处理 1 万条后暂停)
- 按耗时暂停(如单次 slice 超过 50ms)
- 按队列水位动态调整(若待处理任务数 > 线程数 × 2,则加速切片)
统一任务队列 + 终止信号:使用 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),而是通过任务状态显式保存 + 队列重入,达成近似协作式多任务的效果。适用于中高吞吐、低延迟敏感的流式数据处理场景,是传统生产者-消费者模式面向真实业务复杂性的必要演进。

















