如何将Knative KafkaSource连接至启用SASL的AWS MSK集群
解决Knative Eventing KafkaSource连接SASL认证AWS MSK失败问题
问题背景
尝试用Knative Eventing构建事件驱动架构,连接AWS MSK集群,明文无认证方式可正常连接,但启用SASL认证后,KafkaSource无法连接MSK Broker,报错kafka: client has run out of available brokers to talk to (Is your cluster reachable?)。
现有配置与错误信息
KafkaSource配置
kind: KafkaSource metadata: name: kafka-source spec: consumerGroup: kntive-groups bootstrapServers: - my-cluster-kafka-bootstrap.kafka:9096 #MSK broker topics: - knative-input-topic sink: ref: apiVersion: serving.knative.dev/v1 kind: Service name: my-app-service uri: /myAppUrl net: sasl: enable: true user: secretKeyRef: name: msk-secret key: user password: secretKeyRef: name: msk-secret key: password type: secretKeyRef: name: msk-secret key: saslType tls: enable: false caCert: secretKeyRef: name: msk-secret-tf key: ca.crt
状态信息
kubectl get kafkasource kafka-source -n my-ns -o yaml status: conditions: - lastTransitionTime: "2023-02-03T06:10:15Z" message: 'kafka: client has run out of available brokers to talk to (Is your cluster reachable?)' reason: ClientCreationFailed status: "False" type: ConnectionEstablished - lastTransitionTime: "2023-02-03T06:10:15Z" status: Unknown type: Deployed - lastTransitionTime: "2023-02-03T06:10:15Z" status: Unknown type: InitialOffsetsCommitted - lastTransitionTime: "2023-02-03T06:10:15Z" message: 'kafka: client has run out of available brokers to talk to (Is your cluster reachable?)' reason: ClientCreationFailed status: "False" type: Ready - lastTransitionTime: "2023-02-03T06:10:15Z" status: "True" type: SinkProvided
控制器日志
kubectl logs deployment.apps/kafka-controller-manager -n knative-eventing { "level": "info", "ts": "2023-02-03T06:21:18.865Z", "logger": "kafka-controller", "caller": "client/config.go:288", "msg": "Built Sarama config: &{Admin:{Retry:{Max:5 Backoff:100ms} Timeout:3s} Net:{MaxOpenRequests:5 DialTimeout:30s ReadTimeout:30s WriteTimeout:30s TLS:{Enable:false Config:<nil>} SASL:{Enable:true Mechanism:SCRAM-SHA-512 Version:0 Handshake:true AuthIdentity: User:maas-ml-testuser2 Password: SCRAMAuthzID: SCRAMClientGeneratorFunc:0x1861980 TokenProvider:<nil> GSSAPI:{AuthType:0 KeyTabPath: KerberosConfigPath: ServiceName: Username: Password: Realm: DisablePAFXFAST:false}} KeepAlive:0s LocalAddr:<nil> Proxy:{Enable:false Dialer:<nil>}} Metadata:{Retry:{Max:3 Backoff:250ms BackoffFunc:<nil>} RefreshFrequency:10m0s Full:true Timeout:0s AllowAutoTopicCreation:true} Producer:{MaxMessageBytes:1000000 RequiredAcks:1 Timeout:10s Compression:none CompressionLevel:-1000 Partitioner:0x17cf660 Idempotent:false Return:{Successes:true Errors:true} Flush:{Bytes:0 Messages:0 Frequency:0s MaxMessages:0} Retry:{Max:3 Backoff:100ms BackoffFunc:<nil>} Interceptors:[]} Consumer:{Group:{Session:{Timeout:10s} Heartbeat:{Interval:3s} Rebalance:{Strategy:0x2f67290 Timeout:1m0s Retry:{Max:4 Backoff:2s}} Member:{UserData:[]}} Retry:{Backoff:2s BackoffFunc:<nil>} Fetch:{Min:1 Default:1048576 Max:0} MaxWaitTime:250ms MaxProcessingTime:100ms Return:{Errors:true} Offsets:{CommitInterval:0s AutoCommit:{Enable:true Interval:1s} Initial:-2 Retention:0s Retry:{Max:3}} IsolationLevel:0 Interceptors:[]} ClientID:sarama RackID: ChannelBufferSize:256 ApiVersionsRequest:true Version:1.0.0 MetricRegistry:0xc002ca4080}", "commit": "394f005-dirty", "knative.dev/controller": "knative.dev.eventing-kafka.pkg.source.reconciler.source.Reconciler", "knative.dev/kind": "sources.knative.dev.KafkaSource", "knative.dev/traceid": "4d6b80c4-2116-4acb-b5bc-e1d074c2a380", "knative.dev/key": "coal-dev/uat-kafka-source" } { "level": "error", "ts": "2023-02-03T06:21:19.654Z", "logger": "kafka-controller", "caller": "source/kafkasource.go:184", "msg": "unable to create a kafka client", "commit": "394f005-dirty", "knative.dev/controller": "knative.dev.eventing-kafka.pkg.source.reconciler.source.Reconciler", "knative.dev/kind": "sources.knative.dev.KafkaSource", "knative.dev/traceid": "4d6b80c4-2116-4acb-b5bc-e1d074c2a380", "knative.dev/key": "coal-dev/uat-kafka-source", "error": "kafka: client has run out of available brokers to talk to (Is your cluster reachable?)", "stacktrace": "knative.dev/eventing-kafka/pkg/source/reconciler/source.(*Reconciler).ReconcileKind\n\tknative.dev/eventing-kafka/pkg/source/reconciler/source/kafkasource.go:184\nknative.dev/eventing-kafka/pkg/client/injection/reconciler/sources/v1beta1/kafkasource.(*reconcilerImpl).Reconcile\n\tknative.dev/eventing-kafka/pkg/client/injection/reconciler/sources/v1beta1/kafkasource/reconciler.go:239\nknative.dev/pkg/controller.(*Impl).processNextWorkItem\n\tknative.dev/pkg@v0.0.0-20220818004048-4a03844c0b15/controller/controller.go:542\nknative.dev/pkg/controller.(*Impl).RunContext.func3\n\tknative.dev/pkg@v0.0.0-20220818004048-4a03844c0b15/controller/controller.go:491" } { "level": "error", "ts": "2023-02-03T06:21:19.655Z", "logger": "kafka-controller", "caller": "kafkasource/reconciler.go:302", "msg": "Returned an error", "commit": "394f005-dirty", "knative.dev/controller": "knative.dev.eventing-kafka.pkg.source.reconciler.source.Reconciler", "knative.dev/kind": "sources.knative.dev.KafkaSource", "knative.dev/traceid": "4d6b80c4-2116-4acb-b5bc-e1d074c2a380", "knative.dev/key": "coal-dev/uat-kafka-source", "targetMethod": "ReconcileKind", "error": "kafka: client has run out of available brokers to talk to (Is your cluster reachable?)", "stacktrace": "knative.dev/eventing-kafka/pkg/client/injection/reconciler/sources/v1beta1/kafkasource.(*reconcilerImpl).Reconcile\n\tknative.dev/eventing-kafka/pkg/client/injection/reconciler/sources/v1beta1/kafkasource/reconciler.go:302\nknative.dev/pkg/controller.(*Impl).processNextWorkItem\n\tknative.dev/pkg@v0.0.0-20220818004048-4a03844c0b15/controller/controller.go:542\nknative.dev/pkg/controller.(*Impl).RunContext.func3\n\tknative.dev/pkg@v0.0.0-20220818004048-4a03844c0b15/controller/controller.go:491" } { "level": "error", "ts": "2023-02-03T06:21:19.655Z", "logger": "kafka-controller", "caller": "controller/controller.go:566", "msg": "Reconcile error", "commit": "394f005-dirty", "knative.dev/controller": "knative.dev.eventing-kafka.pkg.source.reconciler.source.Reconciler", "knative.dev/kind": "sources.knative.dev.KafkaSource", "knative.dev/traceid": "4d6b80c4-2116-4acb-b5bc-e1d074c2a380", "knative.dev/key": "coal-dev/uat-kafka-source", "duration": 0.813508291, "error": "kafka: client has run out of available brokers to talk to (Is your cluster reachable?)", "stacktrace": "knative.dev/pkg/controller.(*Impl).handleErr\n\tknative.dev/pkg@v0.0.0-20220818004048-4a03844c0b15/controller/controller.go:566\nknative.dev/pkg/controller.(*Impl).processNextWorkItem\n\tknative.dev/pkg@v0.0.0-20220818004048-4a03844c0b15/controller/controller.go:543\nknative.dev/pkg/controller.(*Impl).RunContext.func3\n\tknative.dev/pkg@v0.0.0-20220818004048-4a03844c0b15/controller/controller.go:491" } { "level": "info", "ts": "2023-02-03T06:21:19.655Z", "logger": "kafka-controller.event-broadcaster", "caller": "record/event.go:285", "msg": "Event(v1.ObjectReference{Kind:\"KafkaSource\", Namespace:\"coal-dev\", Name:\"uat-kafka-source\", UID:\"1b6dd5c4-539a-424a-811c-fd16a5d2468d\", APIVersion:\"sources.knative.dev/v1beta1\", ResourceVersion:\"56774522\", FieldPath:\"\"}): type: 'Warning' reason: 'InternalError' kafka: client has run out of available brokers to talk to (Is your cluster reachable?)", "commit": "394f005-dirty" }
解决方案
1. 启用TLS并验证CA证书
AWS MSK的SASL连接默认要求TLS加密,当前配置中tls.enable设为false,这会导致连接失败。修改配置:
net: tls: enable: true caCert: secretKeyRef: name: msk-secret-tf key: ca.crt
确保msk-secret-tf中的ca.crt是AWS MSK的正确CA证书。
2. 确认SASL机制匹配
从日志中看到Sarama配置使用的是SCRAM-SHA-512,需确认:
- MSK集群安全配置中已启用对应SCRAM机制;
msk-secret中的saslType值为SCRAM-SHA-512(注意大小写,Sarama对机制名称大小写敏感)。
3. 验证Secret内容正确性
执行以下命令检查Secret中的值是否与MSK配置一致:
kubectl get secret msk-secret -n my-ns -o jsonpath='{.data.user}' | base64 -d kubectl get secret msk-secret -n my-ns -o jsonpath='{.data.password}' | base64 -d kubectl get secret msk-secret -n my-ns -o jsonpath='{.data.saslType}' | base64 -d
4. 检查网络与安全组配置
- 确认K8s集群与MSK集群VPC连通(同VPC或已配置VPC peering);
- 检查MSK安全组是否开放9096端口给K8s节点;
- 在K8s集群内测试连通性:
kubectl run -it --rm --image=busybox:1.36 net-test -- nc -zv my-cluster-kafka-bootstrap.kafka 9096
5. 调整Sarama重试配置
增加元数据重试次数和超时时间,提升连接容错性:
spec: config: apiVersion: kafka.eventing.knative.dev/v1beta1 kind: KafkaSourceConfig metadata: name: kafka-source-config spec: sarama: metadata: retry: max: 5 backoff: 500ms net: dialTimeout: 60s
修正后的KafkaSource示例配置
kind: KafkaSource metadata: name: kafka-source spec: consumerGroup: kntive-groups bootstrapServers: - my-cluster-kafka-bootstrap.kafka:9096 #MSK broker topics: - knative-input-topic sink: ref: apiVersion: serving.knative.dev/v1 kind: Service name: my-app-service uri: /myAppUrl net: sasl: enable: true user: secretKeyRef: name: msk-secret key: user password: secretKeyRef: name: msk-secret key: password type: secretKeyRef: name: msk-secret key: saslType tls: enable: true caCert: secretKeyRef: name: msk-secret-tf key: ca.crt
内容的提问来源于stack exchange,提问作者Rajashekhar Meesala
相关产品推荐
相关产品推荐

