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

Kafka Connect无法识别MongoDB自定义WriteModelStrategy类问题求助

问题:MongoDB Kafka Sink Connector自定义WriteModelStrategy类找不到

尝试实现MongoDB Kafka Sink Connector的自定义WriteModelStrategy,完成代码编写、Docker镜像构建及plugin.path配置后,始终触发类找不到错误,但同Jar包内的自定义转换功能正常,怀疑类加载隔离导致该问题。


自定义WriteModelStrategy代码

package com.fu.connect.sink;

import org.bson.*;

import com.mongodb.client.model.UpdateOneModel;
import com.mongodb.client.model.UpdateOptions;
import com.mongodb.client.model.WriteModel;
import com.mongodb.kafka.connect.sink.converter.SinkDocument;
import com.mongodb.kafka.connect.sink.writemodel.strategy.WriteModelStrategy;

import org.apache.kafka.connect.errors.DataException;

public class CustomWriteModelStrategy implements WriteModelStrategy {

    private static final UpdateOptions UPDATE_OPTIONS = new UpdateOptions().upsert(true);

    // incoming json should have one message key e.g. { "message": "Hello World"}
    @Override
    public WriteModel<BsonDocument> createWriteModel(SinkDocument document) {
        
        // Retrieve the value part of the SinkDocument
        BsonDocument vd = document.getValueDoc().orElseThrow(
                () -> new DataException("Error: cannot build the WriteModel since the value document was missing unexpectedly"));

        // extract message from incoming document
        BsonString message = new BsonString("");
        if (vd.containsKey("message")) {
            message = vd.get("message").asString();
        }

        // Define the filter part of the update statement
        BsonDocument filters = new BsonDocument("counter", new BsonDocument("$lt", new BsonInt32(10)));

        // Define the update part of the update statement
        BsonDocument updateStatement = new BsonDocument();
        updateStatement.append("$inc", new BsonDocument("counter", new BsonInt32(1)));
        updateStatement.append("$push", new BsonDocument("messages", new BsonDocument("message", message)));

        // Return the full update
        return new UpdateOneModel<BsonDocument>(
                filters,
                updateStatement,
                UPDATE_OPTIONS
        );
    }
}

Dockerfile配置

FROM maven:3.6.0-jdk-11-slim AS build
COPY resources/custom_plugins /app/resources/custom_plugins
COPY resources/whitelist.csv /app/config/

WORKDIR /app/resources/custom_plugins
RUN mvn -e clean package

FROM confluentinc/cp-kafka-connect:7.2.2

ARG version
ENV VERSION=$version

USER appuser
RUN mkdir -p app/bin && \
    mkdir -p app/config
COPY --chown=appuser resources/truststore.jks app/config/
COPY --chown=appuser resources/whitelist.csv /app/config/

USER root

RUN confluent-hub install --no-prompt confluentinc/kafka-connect-avro-converter:5.5.3 && \
    confluent-hub install --no-prompt mongodb/kafka-connect-mongodb:1.8.0 && \
    mkdir /usr/share/confluent-hub-components/plugins && \
    mkdir /usr/share/confluent-hub-components/mongo_plugins && \
    cp /usr/share/confluent-hub-components/mongodb-kafka-connect-mongodb/lib/*.jar /usr/share/confluent-hub-components/mongo_plugins && \
    cp /usr/share/confluent-hub-components/confluentinc-kafka-connect-avro-converter/lib/*.jar /usr/share/confluent-hub-components/plugins && \
    cp /usr/share/filestream-connectors/*.jar /usr/share/confluent-hub-components/plugins

USER appuser

ENV ARTIFACT_ID=CustomPlugins-1.0-SNAPSHOT.jar
COPY --from=build /app/resources/custom_plugins/target/$ARTIFACT_ID /usr/share/confluent-hub-components/mongo_plugins/$ARTIFACT_ID

COPY --chown=appuser scripts/*.sh app/bin/
COPY --chown=appuser config/* app/config/

当前plugin.path配置

plugin.path=/usr/share/confluent-hub-components/plugins,/usr/share/confluent-hub-components/mongo_plugins/CustomPlugins-1.0-SNAPSHOT.jar,/usr/share/confluent-hub-components/mongo_plugins/mongo-kafka-connect-1.8.0-confluent.jar,

Sink配置关键项

writemodel.strategy=com.fu.connect.sink.CustomWriteModelStrategy 

报错信息

java.util.concurrent.ExecutionException: org.apache.kafka.connect.runtime.rest.errors.BadRequestException: Connector configuration is invalid and contains the following 1 error(s):
Invalid value com.fu.connect.sink.CustomWriteModelStrategy for configuration writemodel.strategy: Class not found: com.fu.connect.sink.CustomWriteModelStrategy

解决方案

问题核心是Kafka Connect的插件类加载隔离机制:MongoDB Connector的类加载器只能访问自身插件目录下的JAR文件,自定义WriteModelStrategy需要被该类加载器读取到。按以下步骤调整:

  1. 调整自定义JAR存放路径
    将自定义插件JAR直接放入MongoDB Connector的官方插件目录,而非单独的mongo_plugins目录。修改Dockerfile中COPY自定义JAR的步骤:

    ENV ARTIFACT_ID=CustomPlugins-1.0-SNAPSHOT.jar
    COPY --from=build /app/resources/custom_plugins/target/$ARTIFACT_ID /usr/share/confluent-hub-components/mongodb-kafka-connect-mongodb/lib/$ARTIFACT_ID
    
  2. 简化plugin.path配置
    plugin.path只需指定插件目录,Kafka Connect会自动扫描目录下所有JAR,无需单独指定JAR文件:

    plugin.path=/usr/share/confluent-hub-components/plugins,/usr/share/confluent-hub-components/mongodb-kafka-connect-mongodb
    
  3. 验证JAR包内容
    进入容器执行以下命令,确认自定义类已正确打包到JAR中:

    jar tf /usr/share/confluent-hub-components/mongodb-kafka-connect-mongodb/lib/CustomPlugins-1.0-SNAPSHOT.jar | grep com/fu/connect/sink/CustomWriteModelStrategy.class
    

    若输出对应类路径,说明打包正确。

  4. 确认依赖范围
    在pom.xml中将MongoDB Connector及Kafka Connect相关依赖设为provided,避免打包重复依赖导致冲突:

    <dependency>
        <groupId>org.mongodb.kafka</groupId>
        <artifactId>mongo-kafka-connect</artifactId>
        <version>1.8.0</version>
        <scope>provided</scope>
    </dependency>
    <dependency>
        <groupId>org.apache.kafka</groupId>
        <artifactId>connect-api</artifactId>
        <version>7.2.2</version>
        <scope>provided</scope>
    </dependency>
    
  5. 重启Kafka Connect服务
    所有配置修改后,重启Kafka Connect容器,确保类加载器重新加载插件。


内容的提问来源于stack exchange,提问作者Esben Folger Thomas

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.13 22:45:39