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>
- 若使用NodePort暴露,配置为
- 确保
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-sixSecret,且其中的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
相关产品推荐
相关产品推荐

