并發(fā)設(shè)計模式 生產(chǎn)者消費者

使用BlockingQueue實現(xiàn)

import java.util.Random;
import java.util.concurrent.BlockingQueue;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicInteger;

/**生產(chǎn)者*/
public class Producer implements Runnable {

    /**內(nèi)存緩沖區(qū)*/
    private BlockingQueue<PCData> queue;
    /**總數(shù)跨晴,原子操作*/
    private static AtomicInteger count = new AtomicInteger();
    /**停止線程*/
    private volatile boolean isRunning = true;

    public Producer(BlockingQueue<PCData> queue) {
        this.queue = queue;
    }


    @Override
    public void run() {
        System.out.println("start producer id=" + Thread.currentThread().getId());
        try {
            while (isRunning) {
                //構(gòu)造任務(wù)數(shù)據(jù)
                PCData data = new PCData(count.incrementAndGet());
                //提交數(shù)據(jù)到緩沖區(qū)中
                System.out.println(data+" is put into queue");
                if (!queue.offer(data, 2, TimeUnit.SECONDS)) {
                    System.out.println("failed to put data:" + data);
                }
                Thread.sleep(new Random().nextInt(1000));
            }
        } catch (InterruptedException e) {
            e.printStackTrace();
            Thread.currentThread().interrupt();
        }
    }

    /**停止線程*/
    public void stop() {
        isRunning = false;
    }
}
/**任務(wù)相關(guān)的數(shù)據(jù)*/
public final class PCData {
    private final int data;

    public PCData(int data) {
        this.data = data;
    }

    public int getData() {
        return data;
    }

    @Override
    public String toString() {
        return "PCData{" +
                "data=" + data +
                '}';
    }
}
import java.text.MessageFormat;
import java.util.Random;
import java.util.concurrent.BlockingQueue;

public class Consumer implements Runnable {

    /**內(nèi)存緩沖區(qū)*/
    private BlockingQueue<PCData> queue;

    public Consumer(BlockingQueue<PCData> queue) {
        this.queue = queue;
    }

    @Override
    public void run() {
        System.out.println("start Consumer id=" + Thread.currentThread().getId());
        try {
            while (true) {
                //提取任務(wù)
                PCData data = queue.take();
                if (null != data) {
                    int getData = data.getData();
                    System.out.println(MessageFormat.format("{0}*{1}={2}",
                            getData, getData, getData * getData));
                }
                Thread.sleep(new Random().nextInt(1000));
            }
        } catch (InterruptedException e) {
            e.printStackTrace();
            Thread.currentThread().interrupt();
        }
    }
}
import java.util.concurrent.BlockingQueue;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
import java.util.concurrent.LinkedBlockingDeque;

public class Client {
    public static void main(String[] args) throws InterruptedException {
        //創(chuàng)建對象
        BlockingQueue<PCData> queue = new LinkedBlockingDeque<>(10);
        ExecutorService service = Executors.newCachedThreadPool();
        Producer producer1 = new Producer(queue);
        Producer producer2 = new Producer(queue);
        Producer producer3 = new Producer(queue);
        Consumer consumer1 = new Consumer(queue);
        Consumer consumer2 = new Consumer(queue);
        Consumer consumer3 = new Consumer(queue);

        //運行
        service.execute(producer1);
        service.execute(producer2);
        service.execute(producer3);
        service.execute(consumer1);
        service.execute(consumer2);
        service.execute(consumer3);

        //停止生產(chǎn)者
        Thread.sleep(10 * 1000);
        producer1.stop();
        producer2.stop();
        producer3.stop();
        Thread.sleep(3000);
        service.shutdown();

        //消費者還在等,處于WAITING狀態(tài)
    }
}

使用Disruptor實現(xiàn)

<dependency>
   <groupId>com.lmax</groupId>
   <artifactId>disruptor</artifactId>
   <version>3.4.2</version>
</dependency>
public class PCData {

    private long value;

    public void set(long value) {
        this.value = value;
    }

    public long get() {
        return value;
    }
}
import com.lmax.disruptor.EventFactory;

/**
 * 工廠類在Disruptor系統(tǒng)初始化時螃壤,構(gòu)造所有緩沖區(qū)中的對象實例
 * (Disruptor會預(yù)先分配空間)
 */
public class PCDataFactory implements EventFactory<PCData> {

    @Override
    public PCData newInstance() {
        return new PCData();
    }
}
import com.lmax.disruptor.WorkHandler;

import java.text.MessageFormat;

public class Consumer implements WorkHandler<PCData> {

    /**
     * 這里只需要數(shù)據(jù)的處理就可以了
     * @param pcData
     * @throws Exception
     */
    @Override
    public void onEvent(PCData pcData) throws Exception {
        String msg = MessageFormat.format(
                "{0}:Event: --{1}--",
                Thread.currentThread().getId(),
                pcData.get()*pcData.get());
        System.out.println(msg);
    }
}
import com.lmax.disruptor.RingBuffer;

import java.nio.ByteBuffer;

public class Producer   {

    private final RingBuffer<PCData> ringBuffer;

    public Producer(RingBuffer<PCData> ringBuffer) {
        this.ringBuffer = ringBuffer;
    }

    public void pushData(ByteBuffer bb) {
        //獲取序列
        long sequence = ringBuffer.next();
        try {
            //設(shè)置數(shù)據(jù)
            PCData pcData = ringBuffer.get(sequence);
            pcData.set(bb.getLong(0));
        } finally {
            //標記可用
            ringBuffer.publish(sequence);
        }
    }
}
import com.lmax.disruptor.BlockingWaitStrategy;
import com.lmax.disruptor.dsl.Disruptor;
import com.lmax.disruptor.dsl.ProducerType;

import java.nio.ByteBuffer;
import java.util.concurrent.Executors;

public class TestClient {
    public static void main(String[] args) throws InterruptedException {
        //RingBuffer的長度的必須是2的整數(shù)次冪领斥。
        int ringBufferSize = 4;

        Disruptor<PCData> disruptor = new Disruptor<PCData>(
                new PCDataFactory(),
                ringBufferSize,
                Executors.defaultThreadFactory(),
                ProducerType.MULTI,
                new BlockingWaitStrategy()
        );
        //每一個消費者會映射到一個線程中。
        disruptor.handleEventsWithWorkerPool(
                new Consumer(),
                new Consumer(),
                new Consumer(),
                new Consumer()
        );
        //不要忘了啟動Disruptor
        disruptor.start();

        Producer producer = new Producer(disruptor.getRingBuffer());
        ByteBuffer bb = ByteBuffer.allocate(8);
        for (int i = 0; true ; i++) {
            Thread.sleep(100);
            producer.pushData(bb.putLong(0, i));
        }

    }
}
?著作權(quán)歸作者所有,轉(zhuǎn)載或內(nèi)容合作請聯(lián)系作者
  • 序言:七十年代末哩都,一起剝皮案震驚了整個濱河市,隨后出現(xiàn)的幾起案子婉徘,更是在濱河造成了極大的恐慌漠嵌,老刑警劉巖咐汞,帶你破解...
    沈念sama閱讀 217,657評論 6 505
  • 序言:濱河連續(xù)發(fā)生了三起死亡事件,死亡現(xiàn)場離奇詭異儒鹿,居然都是意外死亡化撕,警方通過查閱死者的電腦和手機,發(fā)現(xiàn)死者居然都...
    沈念sama閱讀 92,889評論 3 394
  • 文/潘曉璐 我一進店門约炎,熙熙樓的掌柜王于貴愁眉苦臉地迎上來植阴,“玉大人,你說我怎么就攤上這事圾浅÷邮郑” “怎么了?”我有些...
    開封第一講書人閱讀 164,057評論 0 354
  • 文/不壞的土叔 我叫張陵狸捕,是天一觀的道長喷鸽。 經(jīng)常有香客問我,道長灸拍,這世上最難降的妖魔是什么做祝? 我笑而不...
    開封第一講書人閱讀 58,509評論 1 293
  • 正文 為了忘掉前任,我火速辦了婚禮株搔,結(jié)果婚禮上剖淀,老公的妹妹穿的比我還像新娘。我一直安慰自己纤房,他們只是感情好纵隔,可當我...
    茶點故事閱讀 67,562評論 6 392
  • 文/花漫 我一把揭開白布。 她就那樣靜靜地躺著炮姨,像睡著了一般捌刮。 火紅的嫁衣襯著肌膚如雪。 梳的紋絲不亂的頭發(fā)上舒岸,一...
    開封第一講書人閱讀 51,443評論 1 302
  • 那天绅作,我揣著相機與錄音,去河邊找鬼蛾派。 笑死俄认,一個胖子當著我的面吹牛,可吹牛的內(nèi)容都是我干的洪乍。 我是一名探鬼主播眯杏,決...
    沈念sama閱讀 40,251評論 3 418
  • 文/蒼蘭香墨 我猛地睜開眼,長吁一口氣:“原來是場噩夢啊……” “哼壳澳!你這毒婦竟也來了岂贩?” 一聲冷哼從身側(cè)響起,我...
    開封第一講書人閱讀 39,129評論 0 276
  • 序言:老撾萬榮一對情侶失蹤巷波,失蹤者是張志新(化名)和其女友劉穎萎津,沒想到半個月后卸伞,有當?shù)厝嗽跇淞掷锇l(fā)現(xiàn)了一具尸體,經(jīng)...
    沈念sama閱讀 45,561評論 1 314
  • 正文 獨居荒郊野嶺守林人離奇死亡锉屈,尸身上長有42處帶血的膿包…… 初始之章·張勛 以下內(nèi)容為張勛視角 年9月15日...
    茶點故事閱讀 37,779評論 3 335
  • 正文 我和宋清朗相戀三年荤傲,在試婚紗的時候發(fā)現(xiàn)自己被綠了。 大學(xué)時的朋友給我發(fā)了我未婚夫和他白月光在一起吃飯的照片部念。...
    茶點故事閱讀 39,902評論 1 348
  • 序言:一個原本活蹦亂跳的男人離奇死亡弃酌,死狀恐怖,靈堂內(nèi)的尸體忽然破棺而出儡炼,到底是詐尸還是另有隱情妓湘,我是刑警寧澤,帶...
    沈念sama閱讀 35,621評論 5 345
  • 正文 年R本政府宣布乌询,位于F島的核電站榜贴,受9級特大地震影響,放射性物質(zhì)發(fā)生泄漏妹田。R本人自食惡果不足惜唬党,卻給世界環(huán)境...
    茶點故事閱讀 41,220評論 3 328
  • 文/蒙蒙 一、第九天 我趴在偏房一處隱蔽的房頂上張望鬼佣。 院中可真熱鬧驶拱,春花似錦、人聲如沸晶衷。這莊子的主人今日做“春日...
    開封第一講書人閱讀 31,838評論 0 22
  • 文/蒼蘭香墨 我抬頭看了看天上的太陽晌纫。三九已至税迷,卻和暖如春,著一層夾襖步出監(jiān)牢的瞬間锹漱,已是汗流浹背箭养。 一陣腳步聲響...
    開封第一講書人閱讀 32,971評論 1 269
  • 我被黑心中介騙來泰國打工, 沒想到剛下飛機就差點兒被人妖公主榨干…… 1. 我叫王不留哥牍,地道東北人毕泌。 一個月前我還...
    沈念sama閱讀 48,025評論 2 370
  • 正文 我出身青樓,卻偏偏與公主長得像嗅辣,于是被迫代替她去往敵國和親懈词。 傳聞我的和親對象是個殘疾皇子,可洞房花燭夜當晚...
    茶點故事閱讀 44,843評論 2 354

推薦閱讀更多精彩內(nèi)容