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

Strimzi Kafka Connect集群配置MongoDB Atlas Sink连接器问题求助

核心结论

官方MongoDB Kafka Sink连接器完全支持连接MongoDB Atlas,无需更换为Debezium连接器——Debezium主要用于MongoDB作为数据源的CDC(变更数据捕获)场景,而你的需求是从Kafka同步数据到MongoDB,官方Sink连接器是更匹配的选择。

现有配置的关键问题

你的配置无法运行,最可能的原因是:

  1. KafkaConnector与KafkaConnect集群的关联错误:KafkaConnector的strimzi.io/cluster标签指向了Kafka集群名称my-cluster,但它应该匹配KafkaConnect资源的名称my-mongo-connect,否则连接器无法找到对应的Connect集群。
  2. 镜像构建不完整:如果STRIMZI KAFKA CONNECT IMAGE WITH MONGODB PLUGIN未正确包含MongoDB连接器的所有依赖jar包,会导致连接器无法加载。
  3. 转换器配置冲突: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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.18 13:35:23