You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何将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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.03 06:35:44