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权限,需要为该用户添加对应权限,操作如下:
- 登录MongoDB集群的权限节点,进入MongoDB Shell:
mongo "mongodb://admin_user:admin_pass@machine1:27017,machine2:27017/?authSource=admin"
- 切换到目标数据库(配置中的
DB):
use DB
为用户授予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" } ])
- 方式一:使用内置角色
验证权限是否生效:
db.getUser("user")
检查返回结果中的roles字段,确认已包含新增的权限角色。
说明
由于连接器配置中启用了startup_mode: copy_existing和变更流功能(change.stream.full.document),必须确保MongoDB用户拥有目标集合的changeStream权限,才能正常读取集合数据及后续变更。
内容的提问来源于stack exchange,提问作者Omegaspard
相关产品推荐
相关产品推荐

