From 72df28dbafb24d3f9c2bb65f07c0abdf6ddf089b Mon Sep 17 00:00:00 2001 From: Song Gao Date: Tue, 25 Jun 2024 11:31:54 +0800 Subject: [PATCH] feat: change kafka sink default balancer (#2946) Signed-off-by: Song Gao --- extensions/sinks/kafka/ext/kafka.go | 7 ++++--- 1 file changed, 4 insertions(+), 3 deletions(-) diff --git a/extensions/sinks/kafka/ext/kafka.go b/extensions/sinks/kafka/ext/kafka.go index 22dca4f05b..655b97a8e4 100644 --- a/extensions/sinks/kafka/ext/kafka.go +++ b/extensions/sinks/kafka/ext/kafka.go @@ -114,9 +114,10 @@ func (m *kafkaSink) buildKafkaWriter() error { } brokers := strings.Split(m.c.Brokers, ",") w := &kafkago.Writer{ - Addr: kafkago.TCP(brokers...), - Topic: m.c.Topic, - Balancer: &kafkago.LeastBytes{}, + Addr: kafkago.TCP(brokers...), + Topic: m.c.Topic, + // kafka java-client default balancer + Balancer: &kafkago.Murmur2Balancer{}, Async: false, AllowAutoTopicCreation: true, MaxAttempts: m.kc.MaxAttempts,