NOTE

Producer–Consumer Pattern

What the producer/consumer pattern is, why it is needed, and Go implementations using channels, timeouts, object-oriented encapsulation, and error handling.

Software Architecture & EngineeringCreated Updated 1 min readhistorical

This is a historical learning note and may contain outdated or incomplete understanding.

1. What Is the Producer/Consumer Pattern?

  • There are three roles: producer, consumer, and queue.
    • After producing data, the producer puts it into the queue, and the consumer takes data from the queue and consumes it.

2. Why Is the Producer/Consumer Pattern Needed?

  • Decoupling
    • The producer only focuses on producing data and does not care how consumers consume it.
    • The consumer only focuses on consuming data and does not care how producers produce it.
  • Asynchrony

3. Producer/Consumer Pattern Implementation

3.1. Golang

  • Imitating Java BlockingQueue
    • If n is 1, it is suitable for scenarios where the producer and consumer speeds match.
    • If n > 1, it is suitable for scenarios where the producer is fast and the consumer is slow.
type Config struct {
	ID int
}

func (c *Config) String() string {
	return fmt.Sprintf("Config[ID=%v]", c.ID)

}

type BlockingQueue struct {
	pipe chan *Config
}

func NewBlockingQueue(n int) *BlockingQueue {
	return &BlockingQueue{make(chan *Config, n)}
}

func (b *BlockingQueue) Offer(data *Config) {
	b.pipe <- data
}

func (b *BlockingQueue) Take() (*Config, bool) {
	saasConfig, ok := <-b.pipe
	return saasConfig, ok
}

func (b *BlockingQueue) Close() {
	close(b.pipe)
}

func TestMain3(t *testing.T) {
	queue := NewBlockingQueue(10)

	go produce(queue)
	go consume(queue)

	time.Sleep(time.Hour)
}

func consume(queue *BlockingQueue) {
	for {
		saasConfig, ok := queue.Take()
		if !ok {
			break
		}
		fmt.Println(saasConfig)
	}
}

func produce(queue *BlockingQueue) {
	defer queue.Close()
	for i := 0; i < 100; i++ {
		time.Sleep(time.Microsecond * 200)
		queue.Offer(&Config{ID: i})
	}
}
  • Add a timeout.
type Config struct {
	ID int
}

func (c *Config) String() string {
	return fmt.Sprintf("Config[ID=%v]", c.ID)

}

func TestMain3(t *testing.T) {
	//Set a context that times out after one second
	d := time.Now().Add(1 * time.Second)
	ctx, cancel := context.WithDeadline(context.Background(), d)
	defer cancel()

	channel := make(chan *Config, 10)
	go produce(ctx, channel)
	go consume(ctx, channel)

	time.Sleep(time.Hour)
}

func consume(ctx context.Context, queue chan *Config) {
	defer fmt.Println("consumer exits")
	for {
		select {
		case cfg, ok := <-queue:
			if !ok {
				fmt.Println("consumer received shutdown signal")
				return
			}
			handleCfg(cfg)
		case <-ctx.Done():
			fmt.Println("consumer timed out: err=", ctx.Err())
			return
		}
	}
}

func handleCfg(cfg *Config) {
	fmt.Println(cfg)
}

func produce(ctx context.Context, queue chan *Config) {
	defer fmt.Println("producer exits")
	defer close(queue)
	defer fmt.Println("producer sends shutdown signal")
	for i := 0; i < 100; i++ {
		time.Sleep(time.Millisecond * 200)
		c := fetchCfg(i)
		select {
		case queue <- c:
		case <-ctx.Done():
			fmt.Println("producer timed out: err=", ctx.Err())
			return
		}
	}
}

func fetchCfg(i int) *Config {
	return &Config{ID: i}
}
  • Object-oriented style.
type Config struct {
	ID int
}

func (c *Config) String() string {
	return fmt.Sprintf("Config[ID=%v]", c.ID)
}

type ConfigHandler struct {
	queue chan *Config
}

func NewConfigHandler(num int) *ConfigHandler {
	return &ConfigHandler{queue: make(chan *Config, num)}
}

func (c *ConfigHandler) produce(ctx context.Context) {
	defer fmt.Println("producer exits")
	defer close(c.queue)
	defer fmt.Println("producer sends shutdown signal")
	for i := 0; i < 100; i++ {
		time.Sleep(time.Millisecond * 200)
		cfg := fetchCfg(i)
		select {
		case c.queue <- cfg:
		case <-ctx.Done():
			fmt.Println("producer timed out: err=", ctx.Err())
			return
		}
	}
}

func (c *ConfigHandler) consume(ctx context.Context) {
	defer fmt.Println("consumer exits")
	for {
		select {
		case cfg, ok := <-c.queue:
			if !ok {
				fmt.Println("consumer received shutdown signal")
				return
			}
			handleCfg(cfg)
		case <-ctx.Done():
			fmt.Println("consumer timed out: err=", ctx.Err())
			return
		}
	}
}

func TestMain3(t *testing.T) {
	//Set a context that times out after one second
	d := time.Now().Add(1 * time.Second)
	ctx, cancel := context.WithDeadline(context.Background(), d)
	defer cancel()

	handler := NewConfigHandler(10)

	go handler.produce(ctx)
	go handler.consume(ctx)

	time.Sleep(time.Hour)
}

func handleCfg(cfg *Config) {
	fmt.Println(cfg)
}

func fetchCfg(i int) *Config {
	return &Config{ID: i}
}
  • Add error handling.
type Config struct {
	ID int
}

func (c *Config) String() string {
	return fmt.Sprintf("Config[ID=%v]", c.ID)
}

type ConfigHandler struct {
	queue chan *Config
}

func NewConfigHandler(num int) *ConfigHandler {
	return &ConfigHandler{queue: make(chan *Config, num)}
}

func (c *ConfigHandler) produce(ctx context.Context) error {
	defer fmt.Println("producer exits")
	defer close(c.queue)
	defer fmt.Println("producer sends shutdown signal")
	for i := 0; i < 100; i++ {
		time.Sleep(time.Millisecond * 200)
		cfg, err := fetchCfg(i)
		if err != nil {
			continue
		}
		select {
		case c.queue <- cfg:
		case <-ctx.Done():
			fmt.Println("producer timed out: err=", ctx.Err())
			return ctx.Err()
		}
	}
	return nil
}

func (c *ConfigHandler) consume(ctx context.Context) error {
	defer fmt.Println("consumer exits")
	for {
		select {
		case cfg, ok := <-c.queue:
			if !ok {
				fmt.Println("consumer received shutdown signal")
				return nil
			}
			err := handleCfg(cfg)
			if err != nil {
				return err
			}
		case <-ctx.Done():
			fmt.Println("consumer timed out: err=", ctx.Err())
			return ctx.Err()
		}
	}
}

func TestMain3(t *testing.T) {
	//Set a context that times out after one second
	ctx, cancel := context.WithDeadline(context.Background(), time.Now().Add(10*time.Second))
	defer cancel()

	handler := NewConfigHandler(10)

	group, newCtx := errgroup.WithContext(ctx)
	group.Go(func() error {
		return handler.produce(newCtx)
	})
	group.Go(func() error {
		return handler.consume(newCtx)
	})
	group.Go(func() error {
		return handler.consume(newCtx)
	})
	err := group.Wait()
	if err != nil {
		fmt.Println("WaitGroup: err=", err)
	}
	fmt.Println("main done")
}

func handleCfg(cfg *Config) error {
	fmt.Println(cfg)
	if cfg.ID == 20 {
		return fmt.Errorf("consumer processing error")
	}
	return nil
}

func fetchCfg(i int) (*Config, error) {
	return &Config{ID: i}, nil
}

4. References

Discussion

Sign in with GitHub to comment. Discussions are stored as GitHub Issues.View on GitHub