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

Kafka MirrorMaker2无法连接启用TLS的Kafka节点求助

问题:Kafka MirrorMaker2无法连接目标集群节点

环境背景

  • 两个K8s环境运行Kafka集群:
    • 源集群:无需TLS或认证
    • 目标集群:配置SCRAM-SHA-512认证与TLS加密

MirrorMaker2配置

spec:
  clusters:
  - alias: source-dummy-six
    bootstrapServers: origin-kafka:9092
  - alias: target-dummy-six
    authentication:
      passwordSecret:
        password: target-password
        secretName: password-secret-dummy-six
      type: scram-sha-512
      username: kafka-username
    bootstrapServers: 10.11.12.13:32182
    config:
      config.storage.replication.factor: -1
      config.storage.topic: dummy-six-configs
      group.id: dummy-six-group-id
      offset.storage.replication.factor: -1
      offset.storage.topic: dummy-six-offset
      reconnect.backoff.max.ms: 2000
      reconnect.backoff.ms: 2000
      request.timeout.ms: 60000
      retry.backoff.max.ms: 2000
      retry.backoff.ms: 2000
      socket.connection.setup.timeout.ms: 30000
      ssl.endpoint.identification.algorithm: ""
      status.storage.replication.factor: -1
      status.storage.topic: dummy-six-status
    tls:
      trustedCertificates:
      - certificate: ca.cert
        secretName: target-tls-secret-dummy-six
  connectCluster: target-dummy-six
  logging:
    loggers:
      connect.root.logger.level: INFO
    type: inline
  metricsConfig:
    type: jmxPrometheusExporter
    valueFrom:
      configMapKeyRef:
        key: mirrormaker-metrics-config
        name: mirror-maker-2-metrics
  mirrors:
  - checkpointConnector:
      config:
        checkpoints.topic.replication.factor: 1
        offset-syncs.topic.location: target
        refresh.groups.interval.seconds: 20
        replication.policy.class: com.company.CustomRepPolicy
        replication.policy.dest.metric.topic.name: test_metric_con
        sync.group.offsets.enabled: false
    groupsPattern: .*
    heartbeatConnector:
      config:
        heartbeats.topic.replication.factor: 1
    sourceCluster: source-dummy-six
    sourceConnector:
      config:
        offset-syncs.topic.location: target
        offset-syncs.topic.replication.factor: 1
        refresh.topics.interval.seconds: 20
        replication.factor: 1
        replication.policy.class: com.company.CustomRepPolicy
        replication.policy.dest.metric.topic.name: test_metric_con
        replication.policy.separator: .
        sync.group.offsets.enabled: false
        sync.topic.acls.enabled: "true"
        topic.creation.default.message.format.version: 2.8-IV0
        topic.creation.default.partitions: -1
        topic.creation.default.replication.factor: -1
      tasksMax: 4
    targetCluster: target-dummy-six
    topicsPattern: my_target_topic

错误日志

Node 2 disconnected. (org.apache.kafka.clients.NetworkClient) [kafka-admin-client-thread | adminclient-1]
2024-11-19 16:11:17,103 WARN [AdminClient clientId=adminclient-1] Connection to node 2 (kafka-target-cluster/10.23.52.37:32187) could not be established. Node may not be available. (org.apache.kafka.clients.NetworkClient) [kafka-admin-client-thread | adminclient-1]
2024-11-19 16:11:18,769 INFO [AdminClient clientId=adminclient-1] Node 0 disconnected. (org.apache.kafka.clients.NetworkClient) [kafka-admin-client-thread | adminclient-1]
2024-11-19 16:11:18,769 WARN [AdminClient clientId=adminclient-1] Connection to node 0 (kafka-target-cluster/10.23.52.37:32185) could not be established. Node may not be available. (org.apache.kafka.clients.NetworkClient) [kafka-admin-client-thread | adminclient-1]
2024-11-19 16:11:18,770 INFO App info kafka.admin.client for adminclient-1 unregistered (org.apache.kafka.common.utils.AppInfoParser) [kafka-admin-client-thread | adminclient-1]
2024-11-19 16:11:18,771 INFO [AdminClient clientId=adminclient-1] Metadata update failed (org.apache.kafka.clients.admin.internals.AdminMetadataManager) [kafka-admin-client-thread | adminclient-1]
org.apache.kafka.common.errors.TimeoutException: The AdminClient thread has exited. Call: fetchMetadata
2024-11-19 16:11:18,773 INFO [AdminClient clientId=adminclient-1] Timed out 1 remaining operation(s) during close. (org.apache.kafka.clients.admin.KafkaAdminClient) [kafka-admin-client-thread | adminclient-1]
2024-11-19 16:11:18,779 INFO Metrics scheduler closed (org.apache.kafka.common.metrics.Metrics) [kafka-admin-client-thread | adminclient-1]
2024-11-19 16:11:18,779 INFO Closing reporter org.apache.kafka.common.metrics.JmxReporter (org.apache.kafka.common.metrics.Metrics) [kafka-admin-client-thread | adminclient-1]
2024-11-19 16:11:18,779 INFO Metrics reporters closed (org.apache.kafka.common.metrics.Metrics) [kafka-admin-client-thread | adminclient-1]
2024-11-19 16:11:18,779 ERROR Stopping due to error (org.apache.kafka.connect.cli.AbstractConnectCli) [main]
org.apache.kafka.connect.errors.ConnectException: Failed to connect to and describe Kafka cluster. Check worker's broker connection and security properties.
    at org.apache.kafka.connect.runtime.WorkerConfig.lookupKafkaClusterId(WorkerConfig.java:305)
    at org.apache.kafka.connect.runtime.WorkerConfig.lookupKafkaClusterId(WorkerConfig.java:285)
    at org.apache.kafka.connect.runtime.WorkerConfig.kafkaClusterId(WorkerConfig.java:415)
    at org.apache.kafka.connect.cli.AbstractConnectCli.startConnect(AbstractConnectCli.java:124)
    at org.apache.kafka.connect.cli.AbstractConnectCli.run(AbstractConnectCli.java:94)
    at org.apache.kafka.connect.cli.ConnectDistributed.main(ConnectDistributed.java:116)
Caused by: java.util.concurrent.ExecutionException: org.apache.kafka.common.errors.TimeoutException: Timed out waiting for a node assignment. Call: listNodes
    at java.base/java.util.concurrent.CompletableFuture.reportGet(Unknown Source)
    at java.base/java.util.concurrent.CompletableFuture.get(Unknown Source)
    at org.apache.kafka.common.internals.KafkaFutureImpl.get(KafkaFutureImpl.java:165)
    at org.apache.kafka.connect.runtime.WorkerConfig.lookupKafkaClusterId(WorkerConfig.java:299)
    ... 5 more
Caused by: org.apache.kafka.common.errors.TimeoutException: Timed out waiting for a node assignment. Call: listNodes

已完成的排查

  • 确认bootstrap地址10.11.12.13:32182可通,但无法连接集群节点10.23.52.37:32185/32187,说明MM2能获取集群元数据但无法访问实际节点
  • 执行openssl s_client -connect 10.11.12.13:32182 -showcerts验证TLS证书有效
  • 执行./kafka-acls.sh --list --bootstrap-server 10.11.12.13:32182确认目标集群未配置ACL,默认权限全开(返回SecurityDisabledException)

解决步骤

1. 检查目标Kafka集群的advertised.listeners配置

这是K8s环境中此类问题的最常见原因:

  • Kafka broker会通过advertised.listeners向客户端返回自身的访问地址,若配置为集群内部IP(如Pod IP/ClusterIP),外部客户端(如跨K8s集群的MM2)无法访问
  • 需将advertised.listeners配置为外部可访问的地址:
    • 若使用NodePort暴露,配置为SSL://<K8s-node-ip>:<node-port>(对应TLS端口)
    • 若使用LoadBalancer,配置为SSL://<LB-ip>:<LB-port>
  • 确保listeners配置包含对应的监听端口,与advertised.listeners匹配

2. 网络连通性验证

  • 在MM2的Pod内执行nc -zv 10.23.52.37 32185和nc -zv 10.23.52.37 32187,确认节点端口是否可达
  • 检查K8s NetworkPolicy是否限制了MM2 Pod到目标集群节点的流量
  • 若跨K8s集群,检查集群间防火墙/安全组是否开放了节点端口的访问权限

3. 验证SCRAM认证配置

  • 确认MM2 Pod已正确挂载password-secret-dummy-six Secret,且其中的target-password值正确
  • 可在MM2 Pod内执行cat /path/to/secret/target-password查看密码内容(路径取决于Secret挂载方式)

4. 调整MM2超时配置

当前配置的超时参数可适当调大,避免因网络延迟导致连接失败:

  • 增大request.timeout.ms(如改为120000)
  • 增大socket.connection.setup.timeout.ms(如改为60000)

内容的提问来源于stack exchange,提问作者om shreenidhi

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.16 04:29:53