fix: disable reset kafka connection timeout (#28681)

pr: https://github.com/milvus-io/milvus/pull/28642
issue https://github.com/milvus-io/milvus/issues/28588

Signed-off-by: Enwei Jiao <enwei.jiao@zilliz.com>
This commit is contained in:
Enwei Jiao 2023-11-23 19:42:30 +08:00 committed by GitHub
parent 33bbdf6c88
commit c73bb26782
No known key found for this signature in database
GPG Key ID: 4AEE18F83AFDEB23

View File

@ -59,13 +59,13 @@ func NewKafkaClientInstanceWithConfigMap(config kafka.ConfigMap, extraConsumerCo
func NewKafkaClientInstanceWithConfig(ctx context.Context, config *paramtable.KafkaConfig) (*kafkaClient, error) {
kafkaConfig := getBasicConfig(config.Address.GetValue())
// connection setup timeout, default as 30000ms
// connection setup timeout, default as 30000ms, available range is [1000, 2147483647]
if deadline, ok := ctx.Deadline(); ok {
if deadline.Before(time.Now()) {
return nil, errors.New("context timeout when new kafka client")
}
timeout := time.Until(deadline).Milliseconds()
kafkaConfig.SetKey("socket.connection.setup.timeout.ms", timeout)
// timeout := time.Until(deadline).Milliseconds()
// kafkaConfig.SetKey("socket.connection.setup.timeout.ms", timeout)
}
if (config.SaslUsername.GetValue() == "" && config.SaslPassword.GetValue() != "") ||