fix: disable reset kafka connection timeout (#28681) (#29286)

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>
Co-authored-by: Enwei Jiao <enwei.jiao@zilliz.com>
pull/29324/head
Jiquan Long 2023-12-19 10:11:30 +08:00 committed by GitHub
parent d34ada83cb
commit 90847fbcfd
No known key found for this signature in database
GPG Key ID: 4AEE18F83AFDEB23
1 changed files with 3 additions and 3 deletions

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() != "") ||