How to properly terminate multiple consumers after all producers have ended?

Viewed 390

I've been searching for a solution to this after a while. I've been working in a minimal example where I have a single producer and multiple consumers, and it seems to work properly for the most part (the producer and the consumers don't deadlock, all the values are consumed once, etc...)

However, I'm not sure how to signal the consumers when the producer ends to finish consuming all the produced data left and then end gracefully. Here's a minimal working example I've prepared.

#include <stdio.h>
#include <stdlib.h>
#include <string.h>
#include <semaphore.h>
#include <pthread.h>
#include <stdbool.h>
#define bufferSize 6
#define consumerAmount 3

int buffer[bufferSize];
int producerCell = 0, consumerCell = 0;
sem_t mutexFill, mutexEmpty;
sem_t bufferHasSpace, bufferHasData;
bool producerAlive = true;

void *producer_routine()
{
    int producedCell;
    for (int i = 0; i < 20; i++)
    {
        sem_wait(&bufferHasSpace);
        sem_wait(&mutexFill);
        producedCell = producerCell;
        producerCell = (producerCell + 1) % bufferSize;
        sem_post(&mutexFill);
        memcpy(&buffer[producedCell], &i, sizeof(i));
        printf("Producer generated %i\n", i);
        sem_post(&bufferHasData);
    }
    printf("Producer ended\n");
    producerAlive = false;
}

void *consumer_routine(void *id)
{
    int consumedCell, semaphoreValue;
    do
    {
        sem_wait(&bufferHasData);
        sem_wait(&mutexEmpty);
        consumedCell = consumerCell;
        consumerCell = (consumerCell + 1) % bufferSize;
        sem_getvalue(&bufferHasData, &semaphoreValue);
        sem_post(&mutexEmpty);
        printf("Consumer %i processed %i\n", *(int *)id, buffer[consumedCell]);
        sem_post(&bufferHasSpace);
    }while (producerAlive || semaphoreValue > 0);
    printf("Consumer %i ended\n", *(int *)id);

}

int main()
{
    sem_init(&mutexFill, 1, 1);
    sem_init(&mutexEmpty, 1, 1);
    sem_init(&bufferHasSpace, 1, bufferSize);
    sem_init(&bufferHasData, 1, 0);
    
    pthread_t consumers[consumerAmount];
    int consumerIDs[consumerAmount];
    for (int i = 0; i < consumerAmount; i++)
    {
        consumerIDs[i] = i;
        pthread_create(&consumers[i], NULL, &consumer_routine, &consumerIDs[i]);
    }
    
    pthread_t producer;
    pthread_create(&producer, NULL, &producer_routine, NULL);

    for (int i = 0; i < consumerAmount; i++)
    {
        pthread_join(consumers[i], NULL);
    }
    return 0;
}

Example output from the code snippet:

Producer generated 0
Producer generated 1
Consumer 0 processed 0
Consumer 1 processed 1
Producer generated 2
Producer generated 3
Consumer 2 processed 2
Consumer 2 processed 3
Producer generated 4
Producer generated 5
Producer generated 6
Consumer 1 processed 4
Consumer 1 processed 6
Producer generated 7
Consumer 0 processed 5
Producer generated 8
Consumer 1 processed 7
Producer generated 9
Producer generated 10
Consumer 0 processed 9
Consumer 0 processed 10
Consumer 2 processed 8
Producer generated 11
Producer generated 12
Producer generated 13
Producer generated 14
Producer generated 15
Producer generated 16
Consumer 2 processed 11
Consumer 2 processed 12
Consumer 2 processed 13
Consumer 2 processed 14
Consumer 2 processed 16
Consumer 0 processed 15
Producer generated 17
Producer generated 18
Producer generated 19
Consumer 1 processed 17
Consumer 2 processed 18
Producer ended
Consumer 0 processed 19
Consumer 0 ended
(It gets stuck here indefinitely until the process is killed)

I've tried the following things trying to get the consumers to read the leftover data from the buffer after the producer ends and then exit, but none of them seem to give me an expected result.

  • Not using a boolean flag and relying on the bufferHasData semaphore: Each consumer consumes a single time. After that, since the buffer is empty it exits.
  • Using a boolean flag in conjunction with the bufferHasData semaphore (the example above implements this): Now it works properly, but only the first consumer exits gracefully. The rest of the consumers get stuck in a semaphore waiting for more data, which won't happen because the producer has terminated.
  • Using only a boolean flag: As soon as the producer ends, the consumers will just terminate without processing the rest of the buffer.
  • Not using anything and have the consumers run indefinetly: It consumes all the data, but now I have no way to tell the consumers to stop consuming and exit gracefully.

So, what is the proper way to tell the consumers that the producer(s) have finished, and how to consume the remainder of the buffer without the multiple consumers getting locked in semaphores?

2 Answers

I write producer-consumer programs using Ada, which greatly simplifies the logic. Perhaps the logic I use may be applicable to your C implementation.

Ada has an entity called a task which is analogous to a thread. Ada also has an entity called a protected object which is used to create a buffer shared by tasks. Protected objects implement a high level monitor structure with automatic handling of locking and unlocking the shared data.

I have implemented this example using an Ada package to encapsulate the producer-consumer code. Ada packages have two parts. The package specification defines an API sort of like a C header file. The package body defines the implementation of the functions, procedures, tasks or protected objects exposed in the specification. The package body can also contain type definitions, procedures, functions, tasks and protected objects used by the program but not exposed to the API.

My package specification is

package Stopping_Consumers is
   task type Consumer;
   task Producer;
end Stopping_Consumers;

As you can see, only two things are exposed in the specification. A task type called Consumer allows the programmer to make as many instances of Consumer as needed. The task Producer is not a task type. There is only one possible instance of Producer in this API.

The package body contains all the interesting stuff for this answer.

with Ada.Text_IO; use Ada.Text_IO;

package body Stopping_Consumers is
   -------------------
   -- Shared_Buffer --
   -------------------

   type Buf_Index is mod 10;
   type Buf_Array is array (Buf_Index) of Integer;

   protected Shared_Buffer is
      entry Write (Value : Integer);
      function Producer_Is_Done return Boolean;
      procedure Signal_Done;
      entry Read (Value : out Integer);
      function Is_Empty return Boolean;
      procedure Get_Id (Id : out Natural);
   private
      Buf          : Buf_Array;
      Write_Index  : Buf_Index := 0;
      Read_Index   : Buf_Index := 0;
      Count        : Natural   := 0;
      Is_Done      : Boolean   := False;
      Consumer_Num : Natural   := 0;
   end Shared_Buffer;

   protected body Shared_Buffer is
      entry Write (Value : Integer) when Count < Buf_Index'Modulus is
      begin
         Buf (Write_Index) := Value;
         Write_Index       := Write_Index + 1;
         Count             := Count + 1;
      end Write;

      function Producer_Is_Done return Boolean is
      begin
         return Is_Done;
      end Producer_Is_Done;

      procedure Signal_Done is
      begin
         Is_Done := True;
      end Signal_Done;

      entry Read (Value : out Integer) when Count > 0 is
      begin
         Value      := Buf (Read_Index);
         Read_Index := Read_Index + 1;
         Count      := Count - 1;
      end Read;

      function Is_Empty return Boolean is
      begin
         return Count = 0;
      end Is_Empty;

      procedure Get_Id (Id : out Natural) is
      begin
         Id           := Consumer_Num;
         Consumer_Num := Consumer_Num + 1;
      end Get_Id;

   end Shared_Buffer;

   --------------
   -- Consumer --
   --------------

   task body Consumer is
      Id    : Natural;
      Value : Integer;
   begin
      Shared_Buffer.Get_Id (Id);
      loop
         exit when Shared_Buffer.Producer_Is_Done
           and then Shared_Buffer.Is_Empty;

         select
            Shared_Buffer.Read (Value);
            Put_Line ("Consumer" & Id'Image & " read" & Value'Image);
         or
            delay 0.001;
         end select;
      end loop;

   end Consumer;

   --------------
   -- Producer --
   --------------

   task body Producer is
   begin
      for I in 1 .. 60 loop
         Shared_Buffer.Write (I);
         Put_Line ("Producer wrote" & I'Image);
         delay 0.01;
      end loop;
      Shared_Buffer.Signal_Done;
   end Producer;

end Stopping_Consumers;

The package body contains implementations of the Consumer task type and the Producer task. It also contains the specification and implementation of a protected object named Shared_Buffer. The really interesting part of this example lies in the definition of Shared_Buffer.

Shared_Buffer contains several operations listed as procedures, entries, and functions. It also contains several data items. The Shared_Buffer is first defined by its API and then the implementation is defined. The Shared_Buffer API is

   protected Shared_Buffer is
      entry Write (Value : Integer);
      function Producer_Is_Done return Boolean;
      procedure Signal_Done;
      entry Read (Value : out Integer);
      function Is_Empty return Boolean;
      procedure Get_Id (Id : out Natural);
   private
      Buf          : Buf_Array;
      Write_Index  : Buf_Index := 0;
      Read_Index   : Buf_Index := 0;
      Count        : Natural   := 0;
      Is_Done      : Boolean   := False;
      Consumer_Num : Natural   := 0;
   end Shared_Buffer;

The data element named Buf is defined as an instance of the type Buf_Array. Buf_Array is an array of Integer values indexed by a modular type named Buf_Index. Ada modular types are unsigned integer types exhibiting modular arithmetic. In this example Buf_Index is defined as

type Buf_Index is mod 10;

The range of values of this type is 0 through 9. All arithmetic no a modular type is itself modular. For instance, incrementing a value of Buf_Index with an initial value of 9 results in 0. All the values of such an addition are equivalent to the C code

num = (num + 1) % 10;

Use of a modular type for the index type of the array creates a circular array, which is exactly what is needed in a producer-consumer buffer.

The producer must write data before it can be read. The producer therefore writes using the Write_Index data member and the consumers read using the Read_Index. The Count data member is used to determine both when the buffer is full and when it is empty.

The Is_Done member is a flag set by the producer indicating when the producer is done writing to the buffer.

Consumer_Num assigns an Id number to each consumer when it registers to read data from the buffer.

Ada protected objects can have three kinds of operations; procedures, entries, and functions. Procedures unconditionally modify data in the protected object, automatically manipulating an exclusive read/write lock on the protected object. Entries conditionally modify data in the protected object, automatically manipulating an exclusive read/write lock on the protected object. Functions are read-only and manipulate a shared read lock on the protected object. The compiler writes the lock manipulation code for the programmer.

The implementation of the write entry expresses the bounding condition associated with the execution of the entry.

  entry Write (Value : Integer) when Count < Buf_Index'Modulus is
  begin
     Buf (Write_Index) := Value;
     Write_Index       := Write_Index + 1;
     Count             := Count + 1;
  end Write;

The controlling condition is expressed as "when Count < Buf_Index'Modulus". In this example Buf_Index'Modulus evaluates to 10. If I were to change the definition of Buf_Index to be "type Buf_Index is mod 30" this expression would evaluate to 30. The entry will only execute when the controlling condition evaluates to True. When the controlling condition evaluates to False the calling task is suspended in an entry queue until the condition evaluates to True. Once the condition evaluates to True the tasks suspended in the entry queue are allowed to execute the entry in FIFO order. This entry is called by Producer, forcing Producer to suspend when the buffer is full and execute when the buffer is not full.

The Read entry only allows the consumer to read data from the buffer when the buffer is not empty.

  entry Read (Value : out Integer) when Count > 0 is
  begin
     Value      := Buf (Read_Index);
     Read_Index := Read_Index + 1;
     Count      := Count - 1;
  end Read;

The Producer simply writes a series of numbers to the protected object then calls the Signal_Done procedure to indicate that the producer is done.

   task body Producer is
   begin
      for I in 1 .. 60 loop
         Shared_Buffer.Write (I);
         Put_Line ("Producer wrote" & I'Image);
      end loop;
      Shared_Buffer.Signal_Done;
   end Producer;

The producer is automatically suspended when the Shared_Buffer is full and is automatically un-suspended when the Shared_Buffer is no longer full.

The consumer reads the Shared_Buffer until the Producer is done and the Shared_Buffer is empty. The consumer will continue to read from the shared buffer until both conditions are true.

   task body Consumer is
      Id    : Natural;
      Value : Integer;
   begin
      Shared_Buffer.Get_Id (Id);
      loop
         exit when Shared_Buffer.Producer_Is_Done
           and then Shared_Buffer.Is_Empty;

         select
            Shared_Buffer.Read (Value);
            Put_Line ("Consumer" & Id'Image & " read" & Value'Image);
         or
            delay 0.001;
         end select;
      end loop;

   end Consumer;

The "select" clause within the Consumer loop allows the Consumer to poll the buffer so that the Consumer is not perpetually suspended in the Read entry queue while the Producer has terminated. The statement "delay 0.001;" causes the Consumer to wait 1 millisecond before coming out of the Read entry queue and trying again.

This approach allows any number of Consumer tasks to read the data avoiding duplicate reads and consuming all the produced data without requiring direct communication between the Producer and the Consumers.

EDIT: My experience tells me that elegant and coordinated thread termination is a persistent and difficult problem.

I have modified my example above to more closely mirror the author's example.

Package Specification:

package Stopping_Consumers is
   task type Consumer;
   task Producer;
end Stopping_Consumers;

Package Body:

with Ada.Text_IO; use Ada.Text_IO;

package body Stopping_Consumers is
   -------------------
   -- Shared_Buffer --
   -------------------

   type Buf_Index is mod 6;
   type Buf_Array is array (Buf_Index) of Integer;

   protected Shared_Buffer is
      entry Write (Value : Integer);
      function Producer_Is_Done return Boolean;
      procedure Signal_Done;
      entry Read (Value : out Integer);
      function Is_Empty return Boolean;
      procedure Get_Id (Id : out Natural);
   private
      Buf          : Buf_Array;
      Write_Index  : Buf_Index := 0;
      Read_Index   : Buf_Index := 0;
      Count        : Natural   := 0;
      Is_Done      : Boolean   := False;
      Consumer_Num : Natural   := 0;
   end Shared_Buffer;

   protected body Shared_Buffer is
      entry Write (Value : Integer) when Count < Buf_Index'Modulus is
      begin
         Buf (Write_Index) := Value;
         Write_Index       := Write_Index + 1;
         Count             := Count + 1;
      end Write;

      function Producer_Is_Done return Boolean is
      begin
         return Is_Done;
      end Producer_Is_Done;

      procedure Signal_Done is
      begin
         Is_Done := True;
      end Signal_Done;

      entry Read (Value : out Integer) when Count > 0 is
      begin
         Value      := Buf (Read_Index);
         Read_Index := Read_Index + 1;
         Count      := Count - 1;
      end Read;

      function Is_Empty return Boolean is
      begin
         return Count = 0;
      end Is_Empty;

      procedure Get_Id (Id : out Natural) is
      begin
         Id           := Consumer_Num;
         Consumer_Num := Consumer_Num + 1;
      end Get_Id;

   end Shared_Buffer;

   --------------
   -- Consumer --
   --------------

   task body Consumer is
      Id    : Natural;
      Value : Integer;
   begin
      Shared_Buffer.Get_Id (Id);
      loop
         exit when Shared_Buffer.Producer_Is_Done
           and then Shared_Buffer.Is_Empty;

         select
            Shared_Buffer.Read (Value);
            Put_Line ("Consumer" & Id'Image & " read" & Value'Image);
         or
            delay 0.009;
         end select;
      end loop;
      Put_line("Consumer" & Id'Image & " ended.");
   end Consumer;

   --------------
   -- Producer --
   --------------

   task body Producer is
   begin
      delay 0.001;
      for I in 1 .. 20 loop
         Shared_Buffer.Write (I);
         Put_Line ("Producer wrote" & I'Image);
      end loop;
      Shared_Buffer.Signal_Done;
      Put_Line("Producer ended.");
   end Producer;

end Stopping_Consumers;

Program main procedure:

with Stopping_Consumers; use Stopping_Consumers;
procedure Main is
    Consumers : array(0..2) of Consumer;
begin
   null;
end Main;

Sample Output:

Producer wrote 1
Consumer 0 read 1
Consumer 1 read 2
Producer wrote 2
Producer wrote 3
Consumer 2 read 3
Producer wrote 4
Consumer 0 read 4
Consumer 1 read 5
Producer wrote 5
Producer wrote 6
Consumer 2 read 6
Consumer 0 read 7
Producer wrote 7
Consumer 1 read 8
Producer wrote 8
Producer wrote 9
Consumer 1 read 9
Producer wrote 10
Consumer 0 read 11
Consumer 2 read 10
Producer wrote 11
Producer wrote 12
Consumer 1 read 12
Consumer 2 read 13
Producer wrote 13
Producer wrote 14
Consumer 0 read 14
Consumer 1 read 15
Producer wrote 15
Producer wrote 16
Consumer 2 read 16
Consumer 0 read 17
Producer wrote 17
Producer wrote 18
Consumer 1 read 18
Consumer 0 read 19
Producer wrote 19
Producer wrote 20
Producer ended.
Consumer 2 read 20
Consumer 2 ended.
Consumer 0 ended.
Consumer 1 ended.

How to properly terminate multiple consumers after all producers have ended?

The simple solution is for the producer to produce a special "terminate yourself" value (one for each consumer); where the consumers recognize this value and terminate themselves instead of processing it like normal.

For example (using INT_MAX is the "terminate yourself" value):

void *producer_routine()
{
    int producedCell;
    int terminator = INT_MAX;

    for (int i = 0; i < 20; i++)
    {
        sem_wait(&bufferHasSpace);
        sem_wait(&mutexFill);
        producedCell = producerCell;
        producerCell = (producerCell + 1) % bufferSize;
        sem_post(&mutexFill);
        memcpy(&buffer[producedCell], &i, sizeof(i));
        printf("Producer generated %i\n", i);
        sem_post(&bufferHasData);
    }

    printf("Producer ending\n");
    for (int i = 0; i < consumerAmount; i++)
    {
        sem_wait(&bufferHasSpace);
        sem_wait(&mutexFill);
        producedCell = producerCell;
        producerCell = (producerCell + 1) % bufferSize;
        sem_post(&mutexFill);
        memcpy(&buffer[producedCell], &terminator, sizeof(terminator));
        sem_post(&bufferHasData);
    }
    printf("Producer ended\n");
}

void *consumer_routine(void *id)
{
    int consumedCell, semaphoreValue, value;
    do
    {
        sem_wait(&bufferHasData);
        sem_wait(&mutexEmpty);
        consumedCell = consumerCell;
        consumerCell = (consumerCell + 1) % bufferSize;
        sem_getvalue(&bufferHasData, &semaphoreValue);
        sem_post(&mutexEmpty);
        value = buffer[consumedCell];
        sem_post(&bufferHasSpace);
        if(value != INT_MAX) {
            printf("Consumer %i processed %i\n", *(int *)id, buffer[consumedCell]);
        } else {
            printf("Consumer %i told to terminate\n", *(int *)id);
            break;
        }
    }while (semaphoreValue > 0);
    printf("Consumer %i ended\n", *(int *)id);
}
Related