如何将Kafka Connect的Debezium Postgres连接器对接第三方Kafka集群?
跨集群同步Debezium生成的Topic到第三方Kafka集群(SASL/OAUTHBEARER认证)
要实现将Debezium Postgres连接器生成的Topic同步到第三方Kafka集群,无需对方部署Kafka Connect,最直接的方案是使用Kafka Connect内置的MirrorSourceConnector(镜像源连接器),以下是具体实现步骤:
1. 确认连接器可用性
先检查你的Kafka Connect集群是否已包含MirrorSourceConnector(该连接器属于Kafka官方自带插件,通常默认已包含):
curl -X GET http://<你的Connect集群REST地址>/connector-plugins
返回结果中需存在org.apache.kafka.connect.mirror.MirrorSourceConnector,如果没有,需确保Connect集群的插件目录包含kafka-mirror-connect相关JAR包。
2. 编写MirrorSourceConnector配置
创建JSON配置文件(示例如下),重点配置源集群(你的Kafka集群)和目标集群(第三方)的连接及认证信息:
{ "name": "mirror-debezium-to-thirdparty", "config": { "connector.class": "org.apache.kafka.connect.mirror.MirrorSourceConnector", "tasks.max": "3", // 源集群配置(你的Kafka集群) "source.cluster.alias": "our-kafka", "source.cluster.bootstrap.servers": "你的Kafka引导地址:9092", // 目标集群配置(第三方) "target.cluster.alias": "thirdparty-kafka", "target.cluster.bootstrap.servers": "example1.kafka.connect.com:9092", // 目标集群SASL/OAUTHBEARER认证配置 "target.cluster.security.protocol": "SASL_SSL", "target.cluster.sasl.mechanism": "OAUTHBEARER", "target.cluster.sasl.jaas.config": "org.apache.kafka.common.security.oauthbearer.OAuthBearerLoginModule required oauth.token.endpoint.uri=\"第三方OAUTH令牌端点\" oauth.client.id=\"分配给你的客户端ID\" oauth.client.secret=\"分配给你的客户端密钥\";", "target.cluster.sasl.login.callback.handler.class": "org.apache.kafka.common.security.oauthbearer.OAuthBearerLoginCallbackHandler", // 指定要同步的Debezium生成的Topic(匹配Debezium默认Topic格式) "topics": "dbserver1.*", // 保持Topic名称不变(如需添加前缀可自定义ReplicationPolicy类) "replication.policy.class": "org.apache.kafka.connect.mirror.IdentityReplicationPolicy", // 同步Topic配置(如保留原Topic的分区、副本数等) "sync.topic.configs.enabled": "true", "sync.topic.acls.enabled": "false" } }
关键配置说明
topics:匹配Debezium生成的Topic,Debezium默认Topic格式为<前缀>.<数据库模式>.<表名>,可使用通配符批量匹配(如dbserver1.*)。- 认证部分:
sasl.jaas.config需替换为第三方提供的OAUTH参数,包括令牌端点、客户端ID、客户端密钥等,若第三方有自定义认证逻辑,需替换对应的callback.handler.class。
3. 部署连接器
通过Connect的REST接口提交配置:
curl -X POST -H "Content-Type: application/json" --data @mirror-source-config.json http://<你的Connect集群REST地址>/connectors
4. 验证同步状态
检查连接器运行状态,确认无错误:
curl -X GET http://<你的Connect集群REST地址>/connectors/mirror-debezium-to-thirdparty/status
同时通知第三方团队检查其Kafka集群中是否出现同步后的Topic及数据。
注意事项
- 确保你的Kafka集群与第三方集群网络连通,9092端口(或第三方指定的端口)可正常访问。
- 若Debezium Topic包含大量历史数据,可配置集群化offset存储,避免同步中断。
- 如需过滤特定数据,可结合Transforms(如
Filter)对同步的消息进行预处理后再发送到第三方集群。
内容的提问来源于stack exchange,提问作者langepass88
相关产品推荐
相关产品推荐

