【问题标题】:producer / consumer task. Problem with correct writing to shared buffer生产者/消费者任务。正确写入共享缓冲区的问题
【发布时间】:2020-07-30 08:06:06
【问题描述】:

我正在开展一个解决生产者/消费者调度的经典问题的项目。 Linux Open Suse 42.3 Leep,API System V,C语言

该项目由三个程序组成:生产者、消费者和调度程序。 调度程序的目的是创建 3 个信号量,共享内存,其中有一个缓冲区(数组),其中写入(生产者)和读取(消费者)并运行 n 个生产者和 m 个消费者进程。

每个生产者必须对缓冲区执行 k 个写入周期,消费者必须执行 k 个读取周期。

使用了 3 个信号量:互斥量、空和满。完整信号量的值在程序中用作数组中的索引。

问题是:比如缓冲区大小为3时,生产者写入4部分数据,缓冲区大小为4-5部分数据时(虽然应该有4个)...

消费者阅读正常。

此外,调用 get_semval 函数时,程序的行为无法预测。

请帮忙,我将非常非常感谢您的回答。

制片人

#define BUFFER_SIZE 3
#define MY_RAND_MAX 99      // Highest integer for random number generator
#define LOOP 3              //the number of write / read cycles for each process
#define DATA_DIMENSION 4    // size of portion of data for 1 iteration 
struct Data {
int buf[DATA_DIMENSION];
};
typedef struct Data buffer_item;

buffer_item buffer[BUFFER_SIZE];

void P(int semid)
{
struct sembuf op;
op.sem_num = 0;
op.sem_op = -1;
op.sem_flg = 0;
semop(semid,&op,1);
}

void V(int semid)
{
struct sembuf op;
op.sem_num = 0;
op.sem_op = +1;
op.sem_flg = 0;
semop(semid,&op,1);
}

void Init(int semid,int index,int value)
{
    semctl(semid,index,SETVAL,value);
}

int get_semVal(int sem_id)
{
    int value = semctl(sem_id,0,GETVAL,0);
    return value;
}

int main()
{
    sem_mutex = semget(KEY_MUTEX,1,0);
    sem_empty = semget(KEY_EMPTY,1,0);
    sem_full = semget(KEY_FULL,1,0);

    srand(time(NULL));

const int SIZE = sizeof(buffer[BUFFER_SIZE]);

shm_id = shmget(KEY_SHARED_MEMORY,SIZE, 0);

int i=0;
buffer_item *adr;
do {
    buffer_item nextProduced;
   
    P(sem_empty);

    P(sem_mutex);

    //prepare portion of data
    for(int j=0;j<DATA_DIMENSION;j++)
    {
        nextProduced.buf[j]=rand()%5;
    }

    adr = (buffer_item*)shmat(shm_id,NULL,0);

    int full_value = get_semVal(sem_full);//get index of array
    printf("-----%d------\n",full_value-1);//it’s for test the index of array in buffer

   // write the generated portion of data by index full_value-1
    adr[full_value-1].buf[0] = nextProduced.buf[0];
    adr[full_value-1].buf[1] = nextProduced.buf[1];
    adr[full_value-1].buf[2] = nextProduced.buf[2];
    adr[full_value-1].buf[3] = nextProduced.buf[3];

    shmdt(adr);


    printf("producer %d produced %d %d %d %d\n", getpid(), nextProduced.buf[0],nextProduced.buf[1],nextProduced.buf[2],nextProduced.buf[3]);

    V(sem_mutex);
    V(sem_full);

    i++;
    } while (i<LOOP);

V(sem_empty); 
sleep(1); 
 }

消费者

 …
int main()
{

sem_mutex = semget(KEY_MUTEX,1,0);
sem_empty = semget(KEY_EMPTY,1,0);
sem_full = semget(KEY_FULL,1,0);

srand(time(NULL));

const int SIZE = sizeof(buffer[BUFFER_SIZE]);

shm_id = shmget(KEY_SHARED_MEMORY,SIZE,0);

int i=0;
buffer_item *adr;
do
{
    buffer_item nextConsumed;

    P(sem_full);
    P(sem_mutex);

    int full_value = get_semVal(sem_full);

    adr = (buffer_item*)shmat(shm_id,NULL,0);

    for(int i=0;i<BUFFER_SIZE;i++)
    {
        printf("--%d %d %d %d\n",adr[i].buf[0],adr[i].buf[1],adr[i].buf[2],adr[i].buf[3]);
    }

    for(int i=0;i<BUFFER_SIZE;i++)
    {
        buffer[i].buf[0] = adr[i].buf[0];
        buffer[i].buf[1] = adr[i].buf[1];
        buffer[i].buf[2] = adr[i].buf[2];
        buffer[i].buf[3] = adr[i].buf[3];
    }

    tab(nextConsumed);

        nextConsumed.buf[0]=buffer[full_value-1].buf[0];
        nextConsumed.buf[1]=buffer[full_value-1].buf[1];
        nextConsumed.buf[2]=buffer[full_value-1].buf[2];
        nextConsumed.buf[3]=buffer[full_value-1].buf[3];



    // Set buffer to 0 since we consumed that item
    for(int j=0;j<DATA_DIMENSION;j++)
    {
        buffer[full_value-1].buf[j]=0;
    }

    for(int i=0;i<BUFFER_SIZE;i++)
    {
        adr[i].buf[0]=buffer[i].buf[0];
        adr[i].buf[1]=buffer[i].buf[1];
        adr[i].buf[2]=buffer[i].buf[2];
        adr[i].buf[3]=buffer[i].buf[3];

    }

    shmdt(adr);

    printf("consumer %d consumed %d %d %d %d\n", getpid() ,nextConsumed.buf[0],nextConsumed.buf[1],nextConsumed.buf[2],nextConsumed.buf[3]);

    V(sem_mutex);
    // increase empty
    V(sem_empty);

    i++;
    } while (i<LOOP);

V(sem_full);
sleep(1);
}

调度器

…
struct Data {
   int buf[DATA_DIMENSION];
 };
typedef struct Data buffer_item;

buffer_item buffer[BUFFER_SIZE];

struct TProcList
{
  pid_t processPid;
};
typedef struct TProcList ProcList;
 …

ProcList createProcess(char *name)
{
   pid_t pid;
   ProcList a;

   pid = fork();
   if (!pid){
      kill(getpid(),SIGSTOP);
      execl(name,name,NULL);
      exit(0);
   }
   else if(pid){
       a.processPid=pid;
   }
   else
       cout<<"error forking"<<endl;
   return a;
   }

   int main()
   {
      sem_mutex = semget(KEY_MUTEX,1,IPC_CREAT|0600);
      sem_empty = semget(KEY_EMPTY,1,IPC_CREAT|0600);
      sem_full = semget(KEY_FULL,1,IPC_CREAT|0600);

Init(sem_mutex,0,1);//unlock mutex
Init(sem_empty,0,BUFFER_SIZE); 
Init(sem_full,0,0);//unlock empty

const int SIZE = sizeof(buffer[BUFFER_SIZE]);

shm_id = shmget(KEY_SHARED_MEMORY,SIZE,IPC_CREAT|0600);

buffer_item *adr;
adr = (buffer_item*)shmat(shm_id,NULL,0);

for(int i=0;i<BUFFER_SIZE;i++)
{
    buffer[i].buf[0]=0;
    buffer[i].buf[1]=0;
    buffer[i].buf[2]=0;
    buffer[i].buf[3]=0;
}

for(int i=0;i<BUFFER_SIZE;i++)
{
    adr[i].buf[0] = buffer[i].buf[0];
    adr[i].buf[1] = buffer[i].buf[1];
    adr[i].buf[2] = buffer[i].buf[2];
    adr[i].buf[3] = buffer[i].buf[3];

}

int consumerNumber = 2;
int produserNumber = 2;

ProcList producer_pids[produserNumber];
ProcList consumer_pids[consumerNumber];

for(int i=0;i<produserNumber;i++)
{
    producer_pids[i]=createProcess("/home/andrey/build-c-unknown-Debug/c");//create sleeping processes 
}

for(int i=0;i<consumerNumber;i++)
{
    consumer_pids[i]=createProcess("/home/andrey/build-p-unknown-Debug/p");
}

    sleep(3);

for(int i=0;i<produserNumber;i++)
{
   kill(producer_pids[i].processPid,SIGCONT);//continue processes
   sleep(1);
}
for(int i=0;i<consumerNumber;i++)
{
   kill(consumer_pids[i].processPid,SIGCONT);
   sleep(1);
}


for(int i=0;i<produserNumber;i++)
    {
        waitpid(producer_pids[i].processPid,&stat,WNOHANG);//wait
    }
for(int i=0;i<consumerNumber;i++)
    {
        waitpid(consumer_pids[i].processPid,&stat,WNOHANG);
    }

shmdt(adr);

semctl(sem_mutex,0,IPC_RMID);
semctl(sem_full,0,IPC_RMID);
semctl(sem_empty,0,IPC_RMID);

}

【问题讨论】:

  • 为什么是信号量而不是互斥量和一些条件变量?另外,cout&lt;&lt; 是 C++,而不是 C。它们是两种不同的编程语言,因此有单独的 c 和 c++ 标签。
  • 我需要准确地使用信号量和进程来执行此操作。关于 с 和 с++ 我同意。
  • 每个消费者是只消耗缓冲区中的一个元素,还是消耗整个缓冲区?每个生产者是只为缓冲区生成一个元素,还是整个缓冲区?换句话说,什么是“作业单元”:一个元素、一些元素还是整个缓冲区?这很重要,当使用信号量时,你会看到。 (并将cout &lt;&lt; "error forking" &lt;&lt; endl; 替换为同样描述错误的C printf("Cannot fork: %s.\n", strerror(errno));。fork() 错误情况是pid == -1、0 仅在子级中返回,而所有其他值在父级中返回为子进程 ID。)
  • 在我的任务中,一部分数据是用于对函数(开始、结束、步骤和函数)进行制表的数据。这部分数据以结构数据的形式呈现。同时有几部分数据可以在缓冲区(结构数组)中。在一个周期内,生产者写入一部分数据(开始、结束、步骤和函数),消费者也读取一部分(开始、结束、步骤和函数)。屏幕截图显示(见 ----- 1 ------)信号量值已满:1 0 1 2. 缓冲区大小 3. 询问为什么缓冲区中出现了额外的数据。
  • 我使用 POSIX,包括 Linux(和其他 POSIXy 系统)上的 POSIX semaphores,而您使用的是 System V;恐怕我帮不了你。

标签: c linux process parent-child semaphore


【解决方案1】:

尝试解开其他人编写的未注释代码并不好玩,因此,我将解释一个经过验证的工作方案。

(请注意,cmets 应该始终解释程序员的意图或想法,而不是代码做什么;我们可以阅读代码以了解它的作用。问题是,我们需要先了解代码首先是程序员的想法/意图,然后我们才能将其与实现进行比较。没有 cmets,我需要先阅读代码以尝试猜测意图,然后将其与代码本身进行比较;这就像工作的两倍。)

(我怀疑 OP 的潜在问题是尝试使用信号量值作为缓冲区索引,但并没有仔细研究所有代码以 100% 确定。)

假设共享内存结构如下所示:

struct shared {
    sem_t       lock;            /* Initialized to value 1 */
    sem_t       more;            /* Initialized to 0 */
    sem_t       room;            /* Initialized to MAX_ITEMS */
    size_t      num_items;       /* Initialized to 0 */
    size_t      next_item;       /* Initialized to 0 */
    item_type   item[MAX_ITEMS];
};

我们有struct shared *mem 指向共享内存区域。

请注意,您应该在运行时包含&lt;limits.h&gt;,并验证MAX_ITEMS &lt;= SEM_VALUE_MAX。否则MAX_ITEMS 太大,这个信号量方案可能会失败。 (Linux上的SEM_VALUE_MAX通常是INT_MAX,足够大,但可能会有所不同。而且,如果你在编译时使用-O进行优化,检查将被完全优化掉。所以它是一个非常便宜和合理的检查有。)

mem-&gt;lock 信号量用作互斥体。也就是说,要锁定结构以进行独占访问,进程会等待它。完成后,它会在上面发布。

请注意,虽然sem_post(&amp;(mem-&gt;lock)) 总是会成功(忽略诸如 mem 为 NULL 或指向未初始化内存或已被垃圾覆盖的错误),但从技术上讲,sem_wait() 可以被传递到用户空间处理程序的信号中断安装时没有 SA_RESTART 标志。这就是为什么我建议使用静态内联帮助函数而不是 sem_wait():

static inline int  semaphore_wait(sem_t *const s)
{
    int  result;
    do {
        result = sem_wait(s);
    } while (result == -1 && errno == EINTR);
    return result;
}

static inline int  semaphore_post(sem_t *const s)
{
    return sem_post(s);
}

在信号传递不应中断信号量等待的情况下,您可以使用semaphore_wait()。如果您确实希望信号传递中断信号量的等待,请使用sem_wait();如果它返回-1 和errno == EINTR,则操作由于信号传递而中断,并且信号量实际上并没有递减。 (许多其他低级函数,如read()、write()、send()、recv(),可以以完全相同的方式中断;它们也可以只返回一个短计数,以防中断发生在中途。)

semaphore_post() 只是一个包装器,因此您可以使用“匹配”的 post 和 wait 操作。执行那种“无用”的包装器确实有助于理解代码,你看。

item[] 数组用作循环队列。 num_items 表示其中的项目数。如果num_items &gt; 0,下一个要消费的项目是item[next_item]。如果num_items &lt; MAX_ITEMS,下一个要生产的项目是item[(next_item + num_items) % MAX_ITEMS]。

% 是模运算符。在这里,因为next_item 和num_items 总是正数,所以(next_item + num_items) % MAX_ITEMS 总是介于0 和MAX_ITEMS - 1 之间,包括在内。这就是缓冲区循环的原因。

当生产者构建了一个新项目,比如item_type newitem;,并希望将其添加到共享内存中时,它基本上会执行以下操作:

    /* Omitted: Initialize and fill in 'newitem' members */

    /* Wait until there is room in the buffer */
    semaphore_wait(&(mem->room));

    /* Get exclusive access to the structure members */
    semaphore_wait(&(mem->lock));

    mem->item[(mem->next_item + mem->num_items) % MAX_ITEMS] = newitem;
    mem->num_items++;
    sem_post(&(mem->more));

    semaphore_post(&(mem->lock));

上面通常被称为enqueue,因为它将一个项目附加到一个队列中(这恰好是通过一个循环缓冲区来实现的)。

当消费者想要使用共享缓冲区中的项目 (item_type nextitem;) 时,它会执行以下操作:

    /* Wait until there are items in the buffer */
    semaphore_wait(&(mem->more));

    /* Get exclusive access to the structure members */
    semaphore_wait(&(mem->lock));

    nextitem = mem->item[mem->next_item];
    mem->next_item = (mem->next_item + 1) % MAX_ITEMS;
    mem->num_items = mem->num_items - 1;

    semaphore_post(&(mem->room));

    mem->item[(mem->next_item + mem->num_items) % MAX_ITEMS] = newitem;
    mem->num_items++;
    sem_post(&(mem->more));

    semaphore_post(&(mem->lock));

    /* Omitted: Do work on 'nextitem' here. */

这通常称为 dequeue,因为它从队列中获取下一项。

我建议您首先编写一个单进程测试用例,将MAX_ITEMS 加入队列,然后将它们出列,并验证信号量值是否恢复为初始值。这并不能保证正确性,但它可以解决最典型的错误。

在实践中,我会亲自将排队函数编写为描述共享内存结构的同一头文件中的静态内联帮助程序。差不多

static inline int  shared_get(struct shared *const mem, item_type *const into)
{
    int  err;

    if (!mem || !into)
        return errno = EINVAL; /* Set errno = EINVAL, and return EINVAL. */

    /* Wait for the next item in the buffer. */
    do {
        err = sem_wait(&(mem->more));
    } while (err == -1 && errno == EINTR);
    if (err)
        return errno;

    /* Exclusive access to the structure. */
    do {
        err = sem_wait(&(mem->lock));
    } while (err == -1 && errno == EINTR);

    /* Copy item to caller storage. */
    *into = mem->item[mem->next_item];

    /* Update queue state. */
    mem->next_item = (mem->next_item + 1) % MAX_ITEMS;
    mem->num_items--;

    /* Account for the newly freed slot. */
    sem_post(&(mem->room));

    /* Done. */
    sem_post(&(mem->lock));
    return 0;
}       

和

static inline int  shared_put(struct shared *const mem, const item_type *const from)
    int  err;

    if (!mem || !into)
        return errno = EINVAL; /* Set errno = EINVAL, and return EINVAL. */

    /* Wait for room in the buffer. */
    do {
        err = sem_wait(&(mem->room));
    } while (err == -1 && errno == EINTR);
    if (err)
        return errno;

    /* Exclusive access to the structure. */
    do {
        err = sem_wait(&(mem->lock));
    } while (err == -1 && errno == EINTR);

    /* Copy item to queue. */
    mem->item[(mem->next_item + mem->num_items) % MAX_ITEMS] = *from;

    /* Update queue state. */
    mem->num_items++;

    /* Account for the newly filled slot. */
    sem_post(&(mem->more));

    /* Done. */
    sem_post(&(mem->lock));
    return 0;
}       

但请注意,这些是我凭记忆编写的,而不是从我的测试程序中复制粘贴的,因为我希望您学习而不是在不理解(和怀疑)的情况下复制粘贴其他人的代码。

为什么我们需要单独的计数器(first_item,num_items)当我们有信号量时,有相应的值?

因为我们无法在sem_wait() 成功/继续/停止阻塞时捕获信号量值。

例如,最初room 信号量被初始化为MAX_ITEMS,所以最多可以有多个生产者并行运行。在sem_wait() 之后立即运行sem_getvalue() 的任何一个都将获得一些稍后 值,而不是导致sem_wait() 返回的值或转换。 (即使使用 SysV 信号量,您也无法获得导致此进程等待返回的信号量值。)

因此,我们认为more 信号量的值不是缓冲区的索引或计数器,而是具有多少次可以从缓冲区中出列而不会阻塞的值,而room 的值是多少次次可以在不阻塞的情况下排队到缓冲区。 lock 信号量授予独占访问权限,这样我们就可以原子地修改共享内存结构(嗯,next_item 和 num_items),而不需要不同的进程同时尝试更改值。

我不能 100% 确定这是最好或最佳的模式,这是最常用的模式之一。它不像我想要的那么健壮:对于num_items 中的每个增量(一个),必须在more 上发布一次;对于num_items 中的每个(一)减量,必须将next_item 增加一并在room 上发布一次,否则该方案将崩溃​​。

不过,还有最后一个问题:

生产者如何表明他们已经完成了? 调度程序如何告诉生产者和/或消费者停止?

我首选的解决方案是在共享内存结构中添加一个标志,例如unsigned int status;,并使用特定的位掩码告诉生产者和消费者该做什么,在等待lock 后立即检查:

#define  STOP_PRODUCING  (1 << 0)
#define  STOP_CONSUMING  (1 << 1)

static inline int  shared_get(struct shared *const mem, item_type *const into)
{
    int  err;

    if (!mem || !into)
        return errno = EINVAL; /* Set errno = EINVAL, and return EINVAL. */

    /* Wait for the next item in the buffer. */
    do {
        err = sem_wait(&(mem->more));
    } while (err == -1 && errno == EINTR);
    if (err)
        return errno;

    /* Exclusive access to the structure. */
    do {
        err = sem_wait(&(mem->lock));
    } while (err == -1 && errno == EINTR);

    /* Need to stop consuming? */
    if (mem->state & STOP_CONSUMING) {
        /* Ensure all consumers see the state immediately */
        sem_post(&(mem->more));
        sem_post(&(mem->lock));
        /* ENOMSG == please stop. */
        return errno = ENOMSG;
    }

    /* Copy item to caller storage. */
    *into = mem->item[mem->next_item];

    /* Update queue state. */
    mem->next_item = (mem->next_item + 1) % MAX_ITEMS;
    mem->num_items--;

    /* Account for the newly freed slot. */
    sem_post(&(mem->room));

    /* Done. */
    sem_post(&(mem->lock));
    return 0;
}       

static inline int  shared_put(struct shared *const mem, const item_type *const from)
    int  err;

    if (!mem || !into)
        return errno = EINVAL; /* Set errno = EINVAL, and return EINVAL. */

    /* Wait for room in the buffer. */
    do {
        err = sem_wait(&(mem->room));
    } while (err == -1 && errno == EINTR);
    if (err)
        return errno;

    /* Exclusive access to the structure. */
    do {
        err = sem_wait(&(mem->lock));
    } while (err == -1 && errno == EINTR);

    /* Time to stop? */
    if (mem->state & STOP_PRODUCING) {
        /* Ensure all producers see the state immediately */
        sem_post(&(mem->lock));
        sem_post(&(mem->room));
        /* ENOMSG == please stop. */
        return errno = ENOMSG;
    }

    /* Copy item to queue. */
    mem->item[(mem->next_item + mem->num_items) % MAX_ITEMS] = *from;

    /* Update queue state. */
    mem->num_items++;

    /* Account for the newly filled slot. */
    sem_post(&(mem->more));

    /* Done. */
    sem_post(&(mem->lock));
    return 0;
}

如果调用者应该停止,它会将ENOMSG 返回给调用者。当状态改变时,当然应该持有lock。添加STOP_PRODUCING 时,还应该在room 信号量上发布(一次)以启动“级联”,以便所有生产者停止;并且在添加STOP_CONSUMING 时,在more 信号量上发布(一次)以启动消费者停止级联。 (他们每个人都会再次发布,以确保每个生产者/消费者尽快看到状态。)

不过,还有其他方案;例如信号(设置volatile sig_atomic_t 标志),但通常很难确保没有竞争窗口:一个进程在标志更改之前检查标志,然后阻塞信号量。

在这个方案中,最好同时验证MAX_ITEMS + NUM_PRODUCERS &lt;= SEM_VALUE_MAX和MAX_ITEMS + NUM_CONSUMERS &lt;= SEM_VALUE_MAX,这样即使在停止级联期间,信号量值也不会溢出。

【讨论】:

  • 美好的一天。我非常感谢您对我的问题的彻底而详细的回答。我仔细阅读了你的每一行,意识到我的错误是概念性的。此外,我的大错误是使用 sem_wait 函数我没有考虑到 sem_wait() 将获得一些稍后的值,而不是导致 sem_wait() 返回的值或转换。您提出的想法确实解决了我的问题。有了你的想法,我实现了我的生产者/消费者问题。
  • 但是,我还有一个小问题。我不完全了解工作标志状态如何分别停止消费者或生产者消费或生产。正如您所写,我在 struct 共享结构中创建了一个字段,并在调度程序中将其初始化为 0(即 mem-> status = 0)。然后在生产者和消费者进程中,我只是在检查条件时读取这个值,即 if (mem-> state & STOP_PRODUCING and if) 和 if (mem-> state & STOP_CONSUMING)。它真的有效!但是我只是不明白它是如何工作的。
  • 例如对于消费者:mem-> state = 0 = 0000 STOP_CONSUMING = (1
  • @Andrii:真的有两种方法。一是生产者和消费者事先同意他们做了多少工作。例如,如果您决定有N * NUM_PRODUCERS * NUM_CONSUMERS 工作单元要完成,那么您需要做的就是让每个生产者都做N * NUM_CONSUMERS,每个消费者都做N * NUM_PRODUCERS 工作单元,然后退出。另一种方法是使用status 标志,让它们都连续工作,直到有人 通过status 标志告诉它们停止。如果您实现第一个,则根本不需要 status 标志,并且可以删除该代码。
  • 谢谢你的一切,我明白了。我真的做了第一个选择
猜你喜欢
  • 2011-03-21
  • 2020-10-24
  • 1970-01-01
  • 1970-01-01
  • 2020-09-21
  • 2018-09-13
  • 2011-02-15
  • 2018-05-05
  • 1970-01-01
相关资源
最近更新 更多