#583·watermill

can watermill kafka client support sarama.NewAsyncProducer() ?

Author: JaylenwaCreated Jul 4, 2025Updated Jul 17, 2026
Labelsenhancementhelp wantedgood first issue

Feature request

Description

From the source code, Kafka's Publisher uses sarama-NewSyncProducer and does not support Asynchronous, which can cause performance bottlenecks.

sourcecode

// NewPublisher creates a new Kafka Publisher.
func NewPublisher(
	config PublisherConfig,
	logger watermill.LoggerAdapter,
) (*Publisher, error) {
	config.setDefaults()

	if err := config.Validate(); err != nil {
		return nil, err
	}

	if logger == nil {
		logger = watermill.NopLogger{}
	}

	producer, err := sarama.NewSyncProducer(config.Brokers, config.OverwriteSaramaConfig)
	if err != nil {
		return nil, errors.Wrap(err, "cannot create Kafka producer")
	}

	if config.OTELEnabled && config.Tracer == nil {
		config.Tracer = NewOTELSaramaTracer()
	}

	if config.Tracer != nil {
		producer = config.Tracer.WrapSyncProducer(config.OverwriteSaramaConfig, producer)
	}

	return &Publisher{
		config:   config,
		producer: producer,
		logger:   logger,
	}, nil
}