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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.22 13:17:35