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.
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
}
Discussion
Sign in with GitHub to comment. Discussions are stored as GitHub Issues.View on GitHub