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需要被该类加载器读取到。按以下步骤调整:
调整自定义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简化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验证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若输出对应类路径,说明打包正确。
确认依赖范围
在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>重启Kafka Connect服务
所有配置修改后,重启Kafka Connect容器,确保类加载器重新加载插件。
内容的提问来源于stack exchange,提问作者Esben Folger Thomas
相关产品推荐
相关产品推荐

