Strimzi Kafka Connect集群配置MongoDB Atlas Sink连接器问题求助
核心结论
官方MongoDB Kafka Sink连接器完全支持连接MongoDB Atlas,无需更换为Debezium连接器——Debezium主要用于MongoDB作为数据源的CDC(变更数据捕获)场景,而你的需求是从Kafka同步数据到MongoDB,官方Sink连接器是更匹配的选择。
现有配置的关键问题
你的配置无法运行,最可能的原因是:
- KafkaConnector与KafkaConnect集群的关联错误:
KafkaConnector的strimzi.io/cluster标签指向了Kafka集群名称my-cluster,但它应该匹配KafkaConnect资源的名称my-mongo-connect,否则连接器无法找到对应的Connect集群。 - 镜像构建不完整:如果
STRIMZI KAFKA CONNECT IMAGE WITH MONGODB PLUGIN未正确包含MongoDB连接器的所有依赖jar包,会导致连接器无法加载。 - 转换器配置冲突:
KafkaConnect全局配置了schemas.enable: true,但连接器局部设为false,若生产者发送的消息格式与配置不匹配,会导致数据解析失败。
分步配置指引
1. 构建包含MongoDB连接器的Strimzi Kafka Connect镜像
基于Strimzi官方镜像,将MongoDB连接器插件(含所有依赖)打包到指定目录:
# 选择与你的Kafka版本匹配的Strimzi镜像 FROM quay.io/strimzi/kafka:0.32.0-kafka-3.2.1 USER root # 创建插件目录 RUN mkdir -p /opt/kafka/plugins/mongodb-sink # 复制MongoDB连接器及依赖jar包(从MongoDB官网下载完整插件包) COPY mongodb-kafka-connect-mongodb-1.10.0.jar /opt/kafka/plugins/mongodb-sink/ COPY mongo-java-driver-4.11.1.jar /opt/kafka/plugins/mongodb-sink/ COPY bson-4.11.1.jar /opt/kafka/plugins/mongodb-sink/ USER 1001
构建并推送镜像到你的私有仓库:
docker build -t your-registry/strimzi-kafka-connect-mongodb:0.32.0-kafka-3.2.1 . docker push your-registry/strimzi-kafka-connect-mongodb:0.32.0-kafka-3.2.1
2. 修正KafkaConnect配置
确保镜像地址正确,调整全局配置与消息格式匹配:
apiVersion: kafka.strimzi.io/v1beta2 kind: KafkaConnect metadata: name: my-mongo-connect annotations: strimzi.io/use-connector-resources: "true" spec: image: your-registry/strimzi-kafka-connect-mongodb:0.32.0-kafka-3.2.1 version: 3.2.1 replicas: 1 bootstrapServers: my-cluster-kafka-bootstrap:9092 logging: type: inline loggers: connect.root.logger.level: "INFO" config: group.id: mongo-connect-group offset.storage.topic: mongo-connect-cluster-offsets config.storage.topic: mongo-connect-cluster-configs status.storage.topic: mongo-connect-cluster-status key.converter: org.apache.kafka.connect.json.JsonConverter value.converter: org.apache.kafka.connect.json.JsonConverter # 若生产者发送无schema的JSON,全局设为false更统一 key.converter.schemas.enable: false value.converter.schemas.enable: false # 根据你的Kafka集群副本数调整(如3),-1仅适用于单broker集群 config.storage.replication.factor: 3 offset.storage.replication.factor: 3 status.storage.replication.factor: 3
3. 修正KafkaConnector配置
关键修正:strimzi.io/cluster标签匹配KafkaConnect名称,确认Atlas连接字符串权限:
apiVersion: kafka.strimzi.io/v1beta2 kind: KafkaConnector metadata: name: mongodb-sink-connector labels: strimzi.io/cluster: my-mongo-connect # 必须匹配KafkaConnect的名称 spec: class: com.mongodb.kafka.connect.MongoSinkConnector tasksMax: 2 config: topics: my-topic # 替换为你的MongoDB Atlas连接字符串,确保用户有目标库的读写权限 connection.uri: "mongodb+srv://<USER>:<PASSWORD>@<CLUSTER-NAME>.mongodb.net/?retryWrites=true&w=majority" database: my_database collection: my_collection post.processor.chain: com.mongodb.kafka.connect.sink.processor.DocumentIdAdder,com.mongodb.kafka.connect.sink.processor.KafkaMetaAdder # 若全局已配置转换器,可省略以下重复配置 # key.converter: org.apache.kafka.connect.json.JsonConverter # key.converter.schemas.enable: false # value.converter: org.apache.kafka.connect.json.JsonConverter # value.converter.schemas.enable: false
4. 部署与验证
# 部署KafkaConnect集群 kubectl apply -f kafka-connect.yaml # 等待集群就绪 kubectl wait kafkaconnect/my-mongo-connect --for=condition=Ready --timeout=300s # 部署Sink连接器 kubectl apply -f mongodb-sink-connector.yaml # 查看连接器状态 kubectl get kafkaconnectors mongodb-sink-connector -o yaml # 查看Connect集群日志排查问题 kubectl logs -l strimzi.io/cluster=my-mongo-connect -c connect
额外排查要点
- 网络连通性:确保Kubernetes集群能访问MongoDB Atlas的公网地址(或配置VPC peering)。
- Atlas权限:确认连接字符串中的数据库用户拥有
my_database和my_collection的读写权限。 - 消息格式:若生产者发送的是带Schema的Avro/Protobuf,需对应调整转换器配置。
内容的提问来源于stack exchange,提问作者chirag
相关产品推荐
相关产品推荐

