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

MongoDB Kafka Sink Connector连接失败问题求助

Kafka到MongoDB数据写入失败排查(MongoDB Kafka Sink Connector)

环境配置详情

1. Connect全局配置文件(connect-standalone-demo.properties)

# Licensed to the Apache Software Foundation (ASF) under one or more
# contributor license agreements.  See the NOTICE file distributed with
# this work for additional information regarding copyright ownership.
# The ASF licenses this file to You under the Apache License, Version 2.0
# (the "License"); you may not use this file except in compliance with
# the License.  You may obtain a copy of the License at
#
#    http://www.apache.org/licenses/LICENSE-2.0
#
# Unless required by applicable law or agreed to in writing, software
# distributed under the License is distributed on an "AS IS" BASIS,
# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
# See the License for the specific language governing permissions and
# limitations under the License.

# These are defaults. This file just demonstrates how to override some settings.
bootstrap.servers=localhost:9092

# The converters specify the format of data in Kafka and how to translate it into Connect data. Every Connect user will
# need to configure these based on the format they want their data in when loaded from or stored into Kafka
key.converter=org.apache.kafka.connect.storage.StringConverter
value.converter=org.apache.kafka.connect.storage.StringConverter
#value.converter=org.apache.kafka.connect.json.JsonConverter
#key.converter=org.apache.kafka.connect.json.JsonConverter
# Converter-specific settings can be passed in by prefixing the Converter's setting with the converter we want to apply
# it to
key.converter.schemas.enable=true
value.converter.schemas.enable=true

offset.storage.file.filename=/tmp/connect.offsets
# Flush much faster than normal, which is useful for testing/debugging
offset.flush.interval.ms=10000

# Set to a list of filesystem paths separated by commas (,) to enable class loading isolation for plugins
# (connectors, converters, transformations). The list should consist of top level directories that include 
# any combination of: 
# a) directories immediately containing jars with plugins and their dependencies
# b) uber-jars with plugins and their dependencies
# c) directories immediately containing the package directory structure of classes of plugins and their dependencies
# Note: symlinks will be followed to discover dependencies or plugins.
# Examples: 
# plugin.path=/usr/local/share/java,/usr/local/share/kafka/plugins,/opt/connectors,
# plugin.path=/home/adminacl/Kafka/kafka_2.13-3.1.0/libs

2. Sink Connector配置文件(file-sink-standalone.properties)

当前配置文件格式错误,混入了REST API请求内容:

curl -X POST -H "Content-Type: application/json" -d ' {
      "connector.class":"com.mongodb.kafka.connect.MongoSinkConnector",
      "tasks.max":"1",
      "topics":"departments",
      "connection.uri":"mongodb://localhost:27017",
      "database":"hrmdb",
      "collection":"departments",
      "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
    
}

3. 启动命令

bin/connect-standalone.sh config/connect-standalone-demo.properties config/file-sink-standalone.properties 

错误信息

ERROR Failed to create job for config/file-sink-standalone.properties (org.apache.kafka.connect.cli.ConnectStandalone:107)
[2022-07-27 17:04:20,424] ERROR Stopping after connector error (org.apache.kafka.connect.cli.ConnectStandalone:117)

问题排查与修复

  1. 修正Sink配置文件格式
    Standalone模式下,Connector配置文件必须使用properties格式,而非REST API的JSON格式。将file-sink-standalone.properties修改为:
connector.class=com.mongodb.kafka.connect.MongoSinkConnector
tasks.max=1
topics=departments
connection.uri=mongodb://localhost:27017
database=hrmdb
collection=departments
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
  1. 启用Plugin路径配置
    全局配置文件中plugin.path被注释,需取消注释并指向存放MongoDB Connector jar包的目录。若jar包在Kafka的libs目录,配置如下:
plugin.path=/home/adminacl/Kafka/kafka_2.13-3.1.0/libs

建议生产环境单独创建plugins目录存放连接器jar包,避免与Kafka核心依赖冲突。

  1. 统一Converter配置(可选)
    全局配置使用StringConverter且启用Schema,而Sink配置覆盖为JsonConverter并禁用Schema。需确保Kafka主题中的消息格式与Converter匹配:
  • 若主题消息为无Schema的JSON,当前Sink配置有效;
  • 若为带Schema的JSON,需将schemas.enable设为true。
  1. 查看详细日志定位深层问题
    当前错误信息过于简略,需查看Kafka Connect的详细日志(默认路径logs/connect.log),日志中会包含具体异常栈信息(如类找不到、参数非法等),帮助精准定位问题。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.25 05:24:11