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

使用TestContainers测试依赖外部库的Kafka Connect自定义SMT遇阻

解决Kafka Connect自定义SMT(含外部依赖)在TestContainers中类加载失败问题

问题背景

我用TestContainers对依赖Apache Commons的自定义Kafka Connect SMT做集成测试,流程如下:

  • 通过TestContainers启动Kafka、ZooKeeper、PostgreSQL、Kafka Connect容器
  • 配置PostgreSQL JDBC Source Connector并搭配自定义SMT
  • 向测试数据库表插入数据,期望Connector生成消息经SMT转换后发送到Kafka Topic,最后验证消息字段正确性

但执行./gradlew test时,未使用SMT的测试通过,使用SMT的测试db_insert_creates_kakfa_message_with_smt失败,返回400错误,提示Class smt.test.RandomField$Value could not be found。排查发现Kafka Connect已扫描绑定的自定义JAR目录,但未注册插件,推测生成的JAR缺失Apache Commons等依赖,但日志无报错。

问题根源

普通JAR打包仅包含SMT自身代码,未将外部依赖(如Apache Commons)一同打包。Kafka Connect容器的类加载器无法找到这些依赖类,导致SMT类无法完成初始化,最终触发“类找不到”的错误(表面上是找不到SMT的内部类,实际是依赖缺失导致类加载失败)。

解决方案

1. 构建包含所有依赖的Fat Jar(Gradle项目)

使用Gradle的Shadow插件生成包含所有依赖的Fat Jar,同时排除Kafka Connect自带的依赖避免版本冲突:

  • 在build.gradle中引入Shadow插件:
plugins {
    id 'java'
    id 'com.github.johnrengelman.shadow' version '7.1.2'
}
  • 配置Shadow Jar任务,排除Kafka Connect核心依赖:
shadowJar {
    mergeServiceFiles()
    dependencies {
        // 排除Kafka Connect自带的API和转换依赖,避免冲突
        exclude(dependency('org.apache.kafka:connect-api'))
        exclude(dependency('org.apache.kafka:connect-transforms'))
    }
    // 指定Fat Jar的名称后缀,方便区分普通JAR
    archiveClassifier.set('all')
}
  • 执行./gradlew shadowJar生成Fat Jar,产物路径通常为build/libs/xxx-all.jar

2. 调整TestContainers挂载的JAR路径

在测试代码中,将原来挂载普通JAR的路径替换为Shadow Jar的路径,确保Kafka Connect容器加载的是包含所有依赖的Fat Jar。例如:

// 替换原有的JAR路径为Shadow Jar路径
File smtJar = new File("build/libs/kafka-connect-testcontainers-custom-smt-all.jar");
kafkaConnectContainer.withPlugin(smtJar);

3. 开启Kafka Connect详细日志排查

修改Kafka Connect容器的日志级别为DEBUG,查看插件加载的详细过程:

kafkaConnectContainer.withEnv("LOG4J_LOGGER_LEVEL", "DEBUG");

启动测试后,查看容器日志,会看到插件加载时的依赖缺失报错,确认问题是否由未打包外部依赖导致。

4. 验证Fat Jar的内容

用以下命令检查Fat Jar是否包含所有必要的类:

jar tf build/libs/kafka-connect-testcontainers-custom-smt-all.jar | grep -E "smt/test/RandomField|org/apache/commons"

确保输出中包含smt/test/RandomField.class、smt/test/RandomField$Value.class以及Apache Commons的相关类。

注意事项

  • 必须排除Kafka Connect自带的依赖(如connect-api),否则Fat Jar中的依赖会与容器内的版本冲突,导致插件加载失败
  • 确保Kafka Connect容器的plugin.path环境变量包含自定义插件的挂载目录,默认情况下TestContainers会自动处理,也可手动指定:
kafkaConnectContainer.withEnv("CONNECT_PLUGIN_PATH", "/usr/share/java,/usr/share/confluent-hub-components,/tmp/custom-plugins");

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.10 15:18:55