如何从SASL_SSL Kafka集群读取并写入PLAINTEXT集群?MM2连接问题排查
问题描述
源Kafka集群采用SASL_SSL安全协议,目标集群无授权,仅使用PLAINTEXT协议。以Kafka Connect模式运行的MirrorMaker2输出如下连接日志:
INFO [AdminClient clientId=src->tgt|my_mm2_connector|replication-source-admin] Cancelled in-flight API_VERSIONS request with correlation id 0 due to node -1 being disconnected (elapsed time since creation: 334ms, elapsed time since send: 334ms, request timeout: 3600000ms) (org.apache.kafka.clients.NetworkClient:344)
日志中发现两个AdminClientConfigs:
- client.id = src->tgt|my_mm2_connector|offset-syncs-target-admin,连接至
kafka_plaintext:9092引导服务器 - client.id = src->tgt|my_mm2_connector|replication-source-admin,连接至
kafa_sasl_ssl:9092引导服务器
这两个配置均使用security.protocol=PLAINTEXT。若修改admin.security.protocol为SASL_SSL,会同时影响这两个AdminClientConfigs。
请问能否在安全与非安全Kafka集群间实现主题复制?
以下是MM2配置:
{ "name": "my_mm2_connector", "config": { "connector.class": "org.apache.kafka.connect.mirror.MirrorSourceConnector", "source.cluster.alias": "src", "target.cluster.alias": "tgt", "source.cluster.bootstrap.servers": "kafa_sasl_ssl:9092", "target.cluster.bootstrap.servers": "kafka_plaintext:9092", "topics": "test_topic", "tasks.max": 1, "replication.factor": 1, "offset-syncs.topic.replication.factor": 1, "offset-syncs.topic.location": "target", "enabled": true, "source.security.protocol": "SASL_SSL", "source.ssl.keystore.type": "JKS", "source.ssl.truststore.type": "JKS", "source.ssl.truststore.location": "/opt/ssl/kafka.truststore.jks", "source.ssl.keystore.password": "changeit", "source.ssl.keystore.location": "/opt/ssl/kafka.keystore.jks", "source.ssl.truststore.password": "changeit", "source.sasl.mechanism": "PLAIN", "source.sasl.jaas.config": "org.apache.kafka.common.security.plain.PlainLoginModule required username=\"kafka\" password=\"kafka-password\"", "target.security.protocol": "PLAINTEXT", "admin.security.protocol": "PLAINTEXT" } }
解决方案
可以实现安全与非安全集群间的主题复制,问题出在全局admin.security.protocol配置无法区分集群,导致连接源集群的AdminClient使用了错误的安全协议。
MirrorMaker2支持按集群别名单独配置AdminClient的安全参数,具体操作如下:
- 删除全局的
admin.security.protocol配置项 - 分别为源集群(src)和目标集群(tgt)添加专属的AdminClient安全配置:
- 源集群AdminClient配置SASL_SSL及对应的认证参数
- 目标集群AdminClient配置PLAINTEXT
修改后的MM2配置示例:
{ "name": "my_mm2_connector", "config": { "connector.class": "org.apache.kafka.connect.mirror.MirrorSourceConnector", "source.cluster.alias": "src", "target.cluster.alias": "tgt", "source.cluster.bootstrap.servers": "kafa_sasl_ssl:9092", "target.cluster.bootstrap.servers": "kafka_plaintext:9092", "topics": "test_topic", "tasks.max": 1, "replication.factor": 1, "offset-syncs.topic.replication.factor": 1, "offset-syncs.topic.location": "target", "enabled": true, "source.security.protocol": "SASL_SSL", "source.ssl.keystore.type": "JKS", "source.ssl.truststore.type": "JKS", "source.ssl.truststore.location": "/opt/ssl/kafka.truststore.jks", "source.ssl.keystore.password": "changeit", "source.ssl.keystore.location": "/opt/ssl/kafka.keystore.jks", "source.ssl.truststore.password": "changeit", "source.sasl.mechanism": "PLAIN", "source.sasl.jaas.config": "org.apache.kafka.common.security.plain.PlainLoginModule required username=\"kafka\" password=\"kafka-password\"", "target.security.protocol": "PLAINTEXT", // 源集群AdminClient安全配置 "src.admin.security.protocol": "SASL_SSL", "src.admin.ssl.keystore.type": "JKS", "src.admin.ssl.truststore.type": "JKS", "src.admin.ssl.truststore.location": "/opt/ssl/kafka.truststore.jks", "src.admin.ssl.keystore.password": "changeit", "src.admin.ssl.keystore.location": "/opt/ssl/kafka.keystore.jks", "src.admin.ssl.truststore.password": "changeit", "src.admin.sasl.mechanism": "PLAIN", "src.admin.sasl.jaas.config": "org.apache.kafka.common.security.plain.PlainLoginModule required username=\"kafka\" password=\"kafka-password\"", // 目标集群AdminClient安全配置 "tgt.admin.security.protocol": "PLAINTEXT" } }
核心原理:
MirrorMaker2允许通过{集群别名}.admin.{配置项}的格式为不同集群的AdminClient单独设置参数,这样连接源集群的AdminClient会使用SASL_SSL认证,连接目标集群的AdminClient使用PLAINTEXT,两者互不干扰,解决了节点断开的认证问题。
内容的提问来源于stack exchange,提问作者Cumba
相关产品推荐
相关产品推荐

