You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在生产者消费者模型中正确创建多线程以避免任务重叠

解决生产者/消费者线程任务重复问题

我来帮你梳理下问题所在,再给你几个实用的解决方案~

首先,你的代码里出现重复的生产/消费输出,核心原因不是线程创建的方式有问题,而是任务生成的唯一性没保证,再加上你可能对生产者消费者模型的任务执行逻辑有些误解:

问题分析

  1. 任务内容重复:你应该是用System.currentTimeMillis()作为生产的任务内容,但这个方法的精度是毫秒级的——当多个生产者线程在同一毫秒内启动并执行生产操作时,就会生成完全相同的时间戳,导致输出重复。这不是“同一个任务被多个线程执行”,而是多个线程生产了内容相同的任务。
  2. 线程职责混淆:生产者消费者模型里,任务是存在阻塞队列中的,队列本身(比如你用的ArrayBlockingQueue)是线程安全的,put()和take()都是原子操作,同一个任务绝对不会被多个消费者取走执行。你看到的重复消费输出,本质是队列里有多个相同内容的任务,而不是同一个任务被重复处理。

解决方案

要实现“每个线程处理不同的任务”,我们只需要保证生产的任务唯一,同时合理设计线程的生产/消费逻辑即可。

方案1:给生产者分配唯一ID,生成带标识的任务

给每个生产者线程分配唯一ID,生产的任务带上生产者ID和递增序号,确保任务绝对唯一。

先修正Producer和Consumer类:

class Producer implements Runnable {
    private final BlockingQueue<String> queue;
    private final int producerId;
    private int taskSeq;

    public Producer(BlockingQueue<String> queue, int producerId) {
        this.queue = queue;
        this.producerId = producerId;
        this.taskSeq = 0;
    }

    @Override
    public void run() {
        try {
            // 每个生产者生产5个任务,可根据需求调整为无限循环
            for (int i = 0; i < 5; i++) {
                taskSeq++;
                String uniqueTask = "任务-生产者" + producerId + "-序号" + taskSeq;
                queue.put(uniqueTask);
                System.out.println("生产者" + producerId + "生产:" + uniqueTask);
                // 模拟生产耗时,避免线程集中执行
                Thread.sleep(100);
            }
        } catch (InterruptedException e) {
            Thread.currentThread().interrupt();
            System.out.println("生产者" + producerId + "被中断");
        }
    }
}

class Consumer implements Runnable {
    private final BlockingQueue<String> queue;
    private final int consumerId;

    public Consumer(BlockingQueue<String> queue, int consumerId) {
        this.queue = queue;
        this.consumerId = consumerId;
    }

    @Override
    public void run() {
        try {
            // 每个消费者消费5个任务,对应生产者的总任务数
            for (int i = 0; i < 5; i++) {
                String task = queue.take();
                System.out.println("消费者" + consumerId + "消费:" + task);
                // 模拟消费耗时
                Thread.sleep(150);
            }
        } catch (InterruptedException e) {
            Thread.currentThread().interrupt();
            System.out.println("消费者" + consumerId + "被中断");
        }
    }
}

然后修改main方法,给每个线程分配唯一ID:

import java.util.concurrent.ArrayBlockingQueue;
import java.util.concurrent.BlockingQueue;

public class ProducerConsumerExample {
    public static void main(String[] args) {
        BlockingQueue<String> blockingQueue = new ArrayBlockingQueue<>(10);
        final int PRODUCER_NUM = 10;
        final int CONSUMER_NUM = 10;

        // 批量创建生产者线程
        for (int i = 0; i < PRODUCER_NUM; i++) {
            new Thread(new Producer(blockingQueue, i)).start();
        }

        // 批量创建消费者线程
        for (int i = 0; i < CONSUMER_NUM; i++) {
            new Thread(new Consumer(blockingQueue, i)).start();
        }
    }
}

方案2:全局唯一任务ID(适合无限循环生产场景)

如果你的生产者需要持续无限生产任务,可以用AtomicInteger生成全局唯一的任务ID,彻底避免重复:

import java.util.concurrent.atomic.AtomicInteger;

class Producer implements Runnable {
    private final BlockingQueue<String> queue;
    private final int producerId;
    // 全局原子类,保证多线程下任务ID递增唯一
    private static final AtomicInteger GLOBAL_TASK_ID = new AtomicInteger(0);

    public Producer(BlockingQueue<String> queue, int producerId) {
        this.queue = queue;
        this.producerId = producerId;
    }

    @Override
    public void run() {
        try {
            // 无限循环生产,直到线程被中断
            while (!Thread.currentThread().isInterrupted()) {
                int taskId = GLOBAL_TASK_ID.incrementAndGet();
                String uniqueTask = "全局任务" + taskId + "-来自生产者" + producerId;
                queue.put(uniqueTask);
                System.out.println("生产者" + producerId + "生产:" + uniqueTask);
                Thread.sleep(100);
            }
        } catch (InterruptedException e) {
            Thread.currentThread().interrupt();
            System.out.println("生产者" + producerId + "停止生产");
        }
    }
}

关键说明

  • ArrayBlockingQueue本身是线程安全的,put()和take()方法已经帮我们处理了线程同步,所以同一个任务绝对只会被一个消费者取走执行,不用担心任务重复执行的问题。
  • 加入Thread.sleep()是为了模拟真实场景的生产/消费耗时,避免所有线程同时操作导致输出混乱,你可以根据实际需求调整或移除。

内容的提问来源于stack exchange,提问作者WIIcky

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.04.29 15:27:43