Skip to content

Commit c6fdb95

Browse files
committed
Chore: changed async producer to sync producer
1 parent 3758f28 commit c6fdb95

2 files changed

Lines changed: 14 additions & 9 deletions

File tree

internal/queue/kafka/operator.go

Lines changed: 5 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -304,9 +304,13 @@ func (k *BytesProduceOperator) Produce(message []byte) error {
304304
return errors.New("message is nil")
305305
}
306306

307-
producer.Input() <- &sarama.ProducerMessage{
307+
_, _, err := producer.SendMessage(&sarama.ProducerMessage{
308308
Topic: k.topic,
309309
Value: sarama.ByteEncoder(message),
310+
})
311+
312+
if err != nil {
313+
return err
310314
}
311315

312316
return nil

internal/queue/kafka/producer.go

Lines changed: 9 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -2,22 +2,23 @@ package kafkaqueue
22

33
import (
44
"fmt"
5-
"github.com/IBM/sarama"
65
"sync"
6+
7+
"github.com/IBM/sarama"
78
)
89

910
// ProducerPool is a pool of producers that can be used to produce messages to Kafka for one set of brokers.
1011
// It is not related to Transaction, Transactional Producer implements by configProvider.
1112
type ProducerPool struct {
1213
locker sync.Mutex
13-
producers []sarama.AsyncProducer
14+
producers []sarama.SyncProducer
1415

1516
brokers []string
1617
configProvider func() *sarama.Config
1718
}
1819

1920
// Take returns a producer from pool. If the producer does not exist, it creates a new one.
20-
func (p *ProducerPool) Take() (producer sarama.AsyncProducer) {
21+
func (p *ProducerPool) Take() (producer sarama.SyncProducer) {
2122
p.locker.Lock()
2223
defer p.locker.Unlock()
2324

@@ -33,7 +34,7 @@ func (p *ProducerPool) Take() (producer sarama.AsyncProducer) {
3334
}
3435

3536
// Return returns a producer to the pool.
36-
func (p *ProducerPool) Return(producer sarama.AsyncProducer) {
37+
func (p *ProducerPool) Return(producer sarama.SyncProducer) {
3738
p.locker.Lock()
3839
defer p.locker.Unlock()
3940

@@ -51,7 +52,7 @@ func (p *ProducerPool) Return(producer sarama.AsyncProducer) {
5152
p.producers = append(p.producers, producer)
5253
}
5354

54-
func (p *ProducerPool) Producers() []sarama.AsyncProducer {
55+
func (p *ProducerPool) Producers() []sarama.SyncProducer {
5556
return p.producers
5657
}
5758

@@ -71,16 +72,16 @@ func NewProducerPool(brokers []string, configProvider func() *sarama.Config) *Pr
7172

7273
pool := &ProducerPool{
7374
locker: sync.Mutex{},
74-
producers: []sarama.AsyncProducer{},
75+
producers: []sarama.SyncProducer{},
7576
brokers: brokers,
7677
configProvider: configProvider,
7778
}
7879

7980
return pool
8081
}
8182

82-
func (p *ProducerPool) generateProducer() sarama.AsyncProducer {
83-
producer, err := sarama.NewAsyncProducer(p.brokers, p.configProvider())
83+
func (p *ProducerPool) generateProducer() sarama.SyncProducer {
84+
producer, err := sarama.NewSyncProducer(p.brokers, p.configProvider())
8485
if err != nil {
8586
fmt.Println("Error creating producer", err)
8687
return nil

0 commit comments

Comments
 (0)