系统优化

This commit is contained in:
hahwu 2025-03-19 18:20:16 +08:00
parent 903559ca2a
commit daf7152619
2 changed files with 2 additions and 1 deletions

View File

@ -32,6 +32,6 @@
"RemoteAddr":"host.docker.internal:9001",
"Partition":3,
"KafkaHost":"kafka-server",
"KafkaHost":"kafka-server-2",
"KafkaPort":"9092"
}

View File

@ -60,6 +60,7 @@ func reconnectKafka() {
for _, p := range Partitions {
if p.Leader.Host == conf.Server.KafkaHost {
log.Debug("Kafka connected partition :%v", p)
conn, err := kafka.DialLeader(context.Background(), "tcp", fmt.Sprintf("%s:%s", conf.Server.KafkaHost, conf.Server.KafkaPort), topic, p.ID)
if err != nil {
log.Fatal("kafka.DialLeader err: %v", err)