mirror of https://github.com/milvus-io/milvus.git
fix: apply custom producer config for kafkaHealthCheck (#39283)
issue: https://github.com/milvus-io/milvus/issues/39287 KafkaHealthCheck init without ProducerExtraConfig. This PR fix that. --------- Signed-off-by: DLT1412 <tuduc93@gmail.com> Co-authored-by: DucLT <duc.le1@be.com.vn>pull/39366/head
parent
172051b050
commit
2a962ad1ec
|
@ -109,6 +109,11 @@ func PulsarHealthCheck(clusterStatus *pcommon.MQClusterStatus) {
|
|||
// KafkaHealthCheck Perform a health check by retrieving cluster metadata
|
||||
func KafkaHealthCheck(clusterStatus *pcommon.MQClusterStatus) {
|
||||
config := kafkamqwrapper.GetBasicConfig(¶mtable.Get().KafkaCfg)
|
||||
// Set extra config for producer
|
||||
pConfig := (¶mtable.Get().KafkaCfg).ProducerExtraConfig.GetValue()
|
||||
for k, v := range pConfig {
|
||||
config.SetKey(k, v)
|
||||
}
|
||||
producer, err := kafka.NewProducer(&config)
|
||||
if err != nil {
|
||||
clusterStatus.Reason = fmt.Sprintf("failed to create Kafka producer: %v", err)
|
||||
|
|
Loading…
Reference in New Issue