【问题标题】:Multiple Producer single consumer多生产者单消费者
【发布时间】:2017-02-24 05:44:12
【问题描述】:

我无法理解多个生产者和单个消费者的问题。我正在处理一项任务,但我不确定创建两个生产者是如何工作的。我了解单个生产者/消费者问题的工作原理,但我不明白如何处理多个生产者是否需要为每个生产者创建两个单独的线程,如果是这种情况,如何用他们的“生产的数据”填充队列,其中一个生产者需要睡觉,而另一个生产者填充一个数据项,然后他们来回切换直到队列缓冲区已满?

只是在寻找解释,因为我不完全了解这将如何工作(在有人提出建议之前,我正在寻找某人来做我的家庭作业,但事实并非如此,只是寻找有用的见解来澄清我对此的想法,所以我可以自己实现)

我在这个网站和其他各种网站上查看了许多其他问题/主题,但仍然无法就我的答案得出结论。

谢谢!

【问题讨论】:

  • 快速示例.. 就像加油站的洗手间。您(消费者)必须从服务员(生产者)那里获得浴室钥匙(互斥锁/锁)才能使用资源(洗手间)。加油站可能有很多服务员,检查顾客,拖地等等,都在做自己的事情,等待归还钥匙。加油站可能有很多消费者,加油,买食物等,都在做自己的事情,在洗手间等待腾出。如果是这样,另一个消费者可以获取使用洗手间的密钥。
  • 好的,这绝对有帮助,所以我收集到的是我可以定义两个生产者线程,并让它们都“接受订单”,哪个先到它先到它?我不必让他们中的一个人同时“工作”,而只是为了完成任务的能力而竞争?
  • 线程需要“竞争”的唯一时间是它们必须共享资源时。例如,缓冲区。生产者线程将信息插入缓冲区,消费者将其取出。必须保护对缓冲区的访问,以便 2 个生产者线程不会将它们的数据写入同一个地方;否则你最终会得到损坏的数据。消费者线程也是如此……如果没有与共享资源的线程同步,您将永远不会拥有明确定义的资源状态,然后就会出现混乱。当线程不访问共享资源时,您希望它们关闭
  • 他们自己的东西并行,这就是你从线程中获得速度增强的地方。您希望临界区(线程访问共享资源的部分)尽可能小。如果您有 100 个线程,但他们将所有时间都花在尝试访问共享资源上,那么您的情况不会比只有 1 个线程更好,实际上您可能会更糟 b/c启动线程需要时间。

标签: c queue consumer producer


【解决方案1】:

这是我的解决方案,仅使用 pipeselect 系统调用来实现 MPSCQ。下图详细说明了它的工作原理:

<producer-thread-1>  {msg produced in heap}
  \
   \  /* only address of msg objects were sent to pipe[1] */
    \
     pipe[1] >>>(kernel)>>> pipe[0]  <consumer-thread>:
    /                                1. polling from pipe[0]
   /                                 2. restore msg object via address ptr
  /                                  3. process then delete the msg object
<producer-thread-2>  {msg produced in heap}

演示代码是用 C++ 编写的,用于将队列封装到一个类中,没有使用 C++11/14/17 特性。首先是队列类,模板形式:

// mpscq.hpp
#include <sys/select.h>
#include <unistd.h>
#include <stdio.h>
#include <errno.h>
#include <string.h>

#define PTR_SIZE (sizeof(void*))

template<class T> class MPSCQ { // Multi Producer Single Consumer Queue
public:
        MPSCQ() {
                int pipe_fd_set[2];
                pipe(pipe_fd_set); // err-handler omitted for this demo
                _fdProducer = pipe_fd_set[1];
                _fdConsumer = pipe_fd_set[0];
        }
        ~MPSCQ() { /* pipe close omitted for this demo */ }
        int producerPush(const T* t) {
                // will be blocked when pipe is full, should always return PTR_SIZE
                return t == NULL ? 0 : write(_fdProducer, &t, PTR_SIZE);
        }
        T* consumerPoll(int timeout = 1);
private:
        int _selectFdConsumer(int timeout);
private:
        int _fdProducer; // pipe_fd_set[1]
        int _fdConsumer; // pipe_fd_set[0]
};

template<class T> T* MPSCQ<T>::consumerPoll(int timeout) {
        if (_selectFdConsumer(timeout) <= 0) {  // timeout or error
                return NULL;
        }
        char ptr_buff[PTR_SIZE];
        ssize_t r = read(_fdConsumer, ptr_buff, PTR_SIZE);
        if (r <= 0) {
                fprintf(stderr, "consumer read EOF or error, r=%d, errno=%d\n", r, errno);
                return NULL;
        }
        T* t;
        memcpy(&t, ptr_buff, PTR_SIZE); // cast received bytes to T*
        return t;
}

template<class T> int MPSCQ<T>::_selectFdConsumer(int timeout) {
        int nfds = _fdConsumer + 1;
        fd_set readfds;
        struct timeval tv;
        while (true) {
                tv.tv_sec = timeout;
                tv.tv_usec = 0;
                FD_ZERO(&readfds);
                FD_SET(_fdConsumer, &readfds);
                int r = select(nfds, &readfds, NULL, NULL, &tv);
                if (r < 0 && errno == EINTR) {
                        continue;
                }
                return r;
        }
}

然后是测试用例:4 个生产者线程发出 1..100000,1 个消费者线程总结它。

// g++ -o mpscq mpscq.cpp -lpthread
#include "mpscq.hpp"
#include <sys/types.h>
#include <pthread.h>

#define PER_THREAD_LOOPS        25000
#define SAMPLE_INTERVAL         10000
#define PRODUCER_THREAD_NUM     4

struct TestMsg {
        int _msgId;     // a dummy demo member
        int64_t _val;   // _val < 0 is an end flag
        TestMsg(int msg_id, int64_t val) :
                _msgId(msg_id),
                _val(val) { };
};

static MPSCQ<TestMsg> TEST_QUEUE;

void* functor_producer(void* arg) {
        int* task_seg = (int*) arg;
        TestMsg* msg;
        for (int i = 0; i <= PER_THREAD_LOOPS; ++ i) {
                int64_t id = PER_THREAD_LOOPS * (*task_seg) + i;
                msg = new TestMsg(id, i >= PER_THREAD_LOOPS ? -1 : id + 1);
                TEST_QUEUE.producerPush(msg);
        }
        delete task_seg;
        return NULL;
}

void* functor_consumer(void* arg) {
        int64_t* sum = (int64_t*)arg;
        int msg_cnt = 0;
        int stop_cnt = 0; // for shutdown gracefully
        TestMsg* msg;
        while (true) {
                if ((msg = TEST_QUEUE.consumerPoll()) == NULL) {
                        continue;
                }
                int64_t val = msg->_val;
                delete msg; // this delete is essential to prevent memory leak
                if (val <= 0) {
                        if ((++ stop_cnt) >= PRODUCER_THREAD_NUM) {
                                printf("all done, sum=%ld\n", *sum);
                                break;
                        }
                } else {
                        *sum += val;
                        if ((++ msg_cnt) % SAMPLE_INTERVAL == 0) {
                                printf("msg_cnt=%d, sum=%ld\n", msg_cnt, *sum);
                        }
                }
        }
        return NULL;
}

int main(int argc, char* const* argv) {
        int64_t sum = 0;
        printf("PTR_SIZE: %d, target: sum(1..%d)\n", PTR_SIZE, PRODUCER_THREAD_NUM * PER_THREAD_LOOPS);
        pthread_t consumer;
        pthread_create(&consumer, NULL, functor_consumer, &sum);
        pthread_t producers[PRODUCER_THREAD_NUM];
        for (int i = 0; i < PRODUCER_THREAD_NUM; ++ i) {
                pthread_create(&producers[i], NULL, functor_producer, new int(i));
        }
        for (int i = 0; i < PRODUCER_THREAD_NUM; ++ i) {
                pthread_join(producers[i], NULL);
        }
        pthread_join(consumer, NULL);
        return 0;
}

样本测试结果:

$ ./mpscq 
PTR_SIZE: 8, target: sum(1..100000)
msg_cnt=10000, sum=490096931
msg_cnt=20000, sum=888646187
msg_cnt=30000, sum=1282852073
msg_cnt=40000, sum=1606611602
msg_cnt=50000, sum=2088863858
msg_cnt=60000, sum=2573791058
msg_cnt=70000, sum=3180398370
msg_cnt=80000, sum=3768718659
msg_cnt=90000, sum=4336431164
msg_cnt=100000, sum=5000050000
all done, sum=5000050000

这里实现的 MPSCQ 是一种消息传递模式,让内核来处理内部队列操作的复杂性。这个技巧的一个副作用是当工作负载很重时,消费者端会有太多的select 调用,这将显着影响性能。 (在这个演示中,消费者每次只获取 8 个字节。为了缓解这种情况,消费者应该维护一个额外的接收缓冲区。)

【讨论】:

    猜你喜欢
    • 2015-04-05
    • 2011-03-12
    • 2019-05-24
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多