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

Kafka Connect MongoDB源连接器配置报错及权限问题解决

Kafka Connect MongoDB源连接器权限错误解决:缺少changeStream权限

问题详情

使用以下JSON配置启动Kafka Connect MongoDB源连接器时,出现权限验证错误:

{
  "name": "lastlook-mongodb-source-connector",
  "config": {
    "connector.class": "com.mongodb.kafka.connect.MongoSourceConnector",
    "tasks.max": "4",
    "connection.uri": "mongodb://user:pass@machine1:27017,machine2:27017/?authSource=admin",
    "database": "DB",
    "collection": "Collection",
    "topic": "Topic",
    "startup_mode": "copy_existing",
    "pipeline":"[]",
    "change.stream.full.document":"updateLookup",
    "producer.override.sasl.mechanism": "SCRAM-SHA-256",
    "producer.override.security.protocol": "SASL_PLAINTEXT",
    "producer.override.sasl.jaas.config": "org.apache.kafka.common.security.scram.ScramLoginModule required username=\"user\" password=\"pass\";"
  }
}

错误提示:

{"error_code":400,"message":"Connector configuration is invalid and contains the following 1 error(s):\nInvalid user permissions. Missing the following action permissions: changeStream\nYou can also find the above list of errors at the endpoint /connector-plugins/{connectorType}/config/validate"}

解决步骤

错误原因是连接MongoDB的用户缺少changeStream权限,需要为该用户添加对应权限,操作如下:

  1. 登录MongoDB集群的权限节点,进入MongoDB Shell:
mongo "mongodb://admin_user:admin_pass@machine1:27017,machine2:27017/?authSource=admin"
  1. 切换到目标数据库(配置中的DB):
use DB
  1. 为用户授予changeStream权限,有两种方式可选:

    • 方式一:使用内置角色
      直接授予MongoDB内置的readChangeStream角色,该角色包含读取变更流的权限:
      db.grantRolesToUser(
        "user",
        [ { role: "readChangeStream", db: "DB" } ]
      )
      
    • 方式二:自定义精细权限
      如果需要更精准的权限控制,创建只包含目标集合所需权限的自定义角色:
      db.createRole(
        {
          role: "kafka_mongo_source_role",
          privileges: [
            { 
              resource: { db: "DB", collection: "Collection" }, 
              actions: [ "changeStream", "find" ] 
            }
          ],
          roles: []
        }
      )
      db.grantRolesToUser("user", [ { role: "kafka_mongo_source_role", db: "DB" } ])
      
  2. 验证权限是否生效:

db.getUser("user")

检查返回结果中的roles字段,确认已包含新增的权限角色。

说明

由于连接器配置中启用了startup_mode: copy_existing和变更流功能(change.stream.full.document),必须确保MongoDB用户拥有目标集合的changeStream权限,才能正常读取集合数据及后续变更。

内容的提问来源于stack exchange,提问作者Omegaspard

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.08 16:28:39