
Apache SeaTunnel 测试编码指南编写高质量、稳定、可重复的单元测试与 E2E 测试【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel本指南面向 Apache SeaTunnel 的贡献者与二次开发者系统讲解如何编写确定性每次运行结果一致、无泄漏释放所有打开的资源、低成本启动尽可能少的容器的单元测试与端到端E2E测试。文中所有规则、代码示例与结构约定均来自 docs/zh/developer/test-coding-guide.md并结合仓库内真实测试类、测试基类与容器实现进行源码级佐证读者学完后可以直接动手为新的 connector、transform 或引擎能力补齐两层测试。本指南是编码指南的补充编码指南覆盖通用的 PR 质量要求本指南专注于测试的稳定性与资源安全。测试的两层体系SeaTunnel 将测试明确分为两层由两个 Maven 插件分别驱动层级命名后缀运行插件依赖外部系统职责单元测试*TestSurefire否在隔离环境中校验 connector 或 core 逻辑端到端E2E测试*ITFailsafe是Testcontainers针对真实服务校验 source、transform、sink 的行为当一个 Pull Request 同时改动了逻辑与集成行为时请在两层都补充测试并在提交前分别运行。从仓库源码看这一约定被严格贯彻例如 JdbcSourceFactoryTest.java 位于src/test/java下且以Test结尾由 Surefire 执行而各 connector 的connector-*-e2e模块如 connector-assert-e2e下的NameIT.java则由 Failsafe 以集成测试方式执行。单元测试规范1. 测行为和契约不测实现细节单元测试应验证输入输出行为和配置契约不要绑定私有方法内部实现或临时代码结构。这保证了即使工厂内部实现被重构测试仍然稳定。真实案例来自 JdbcSourceFactoryTest.java模块connector-jdbcprivate MapString, Object baseConfig() { MapString, Object cfg new HashMap(); cfg.put(url, jdbc:mysql://localhost:3306/test); cfg.put(driver, com.mysql.cj.jdbc.Driver); return cfg; } Test void testValidConfigWithTablePath() { MapString, Object cfg baseConfig(); cfg.put(table_path, test.users); Assertions.assertDoesNotThrow(() - validate(cfg)); }该测试保护的是对外配置契约。结合源码可以看到validate经由ConfigValidator.of(ReadonlyConfig.fromMap(config)).validate(rule)完成rule来自factory.optionRule()这正是 SeaTunnel 提交任务时的真实校验入口FactoryUtil.createAndPrepareSource路径因此即使工厂创建逻辑重构只要对外配置契约不变测试依旧稳定。2. 单元测试必须确定且本地可跑单元测试应快速且确定不使用Thread.sleep不使用无固定种子的随机断言不依赖墙钟时间如果用例引入跨进程或外部依赖请明确归类到对应测试层级并保证断言可重复、可诊断。3. 用 mock、stub、fake 隔离依赖单元测试要把目标行为与外部 IO 隔离。对协作依赖优先使用内存实现或轻量测试替身。每个测试方法尽量只覆盖一个断言范围失败时能直接定位单一行为。引擎侧真实 mock 案例来自 JobInfoServiceNullSafetyTest.java模块seatunnel-engine-serverBeforeEach void setUp() { nodeEngine mock(NodeEngineImpl.class); hazelcastInstance mock(HazelcastInstance.class); runningJobInfoMap mock(IMap.class); finishedJobStateMap mock(IMap.class); finishedJobMetricsMap mock(IMap.class); finishedJobVertexInfoMap mock(IMap.class); when(nodeEngine.getHazelcastInstance()).thenReturn(hazelcastInstance); when(hazelcastInstance.getMap(Constant.IMAP_RUNNING_JOB_INFO)).thenReturn(runningJobInfoMap); when(hazelcastInstance.getMap(Constant.IMAP_FINISHED_JOB_STATE)).thenReturn(finishedJobStateMap); when(hazelcastInstance.getMap(Constant.IMAP_FINISHED_JOB_METRICS)) .thenReturn(finishedJobMetricsMap); when(hazelcastInstance.getMap(Constant.IMAP_FINISHED_JOB_VERTEX_INFO)) .thenReturn(finishedJobVertexInfoMap); jobInfoService new JobInfoService(nodeEngine); } Test void shouldReturnJobIdOnlyWhenFinishedMetricsIsMissing() { when(runningJobInfoMap.get(jobId)).thenReturn(null); when(finishedJobStateMap.get(jobId)).thenReturn(jobState); when(finishedJobMetricsMap.get(jobId)).thenReturn(null); JsonObject result jobInfoService.getJobInfoJson(jobId); Assertions.assertEquals(jobId.toString(), result.getString(RestConstant.JOB_ID, null)); }该示例体现了三项实践在边界处 mock 所有外部 map 依赖Hazelcast 的IMap无需启动引擎服务通过when(...).thenReturn(...)精准构造单一业务分支例如“运行中信息缺失、已完成状态存在、指标缺失”这一条路径断言只关注可观察输出结果 JSON 契约中的JOB_ID字段而非内部调用过程。4. 错误路径同时校验异常类型与关键报错信息负向用例不仅要断言异常类还要断言关键报错内容以保护用户可见错误质量。真实案例来自 MongodbIncrementalSourceFactoryTest.java模块connector-cdc-mongodbAssertions.assertThrows( MongodbConnectorException.class, () - MongodbSourceConfigProvider.newBuilder() .startupOptions( new StartupConfig(StartupMode.EARLIEST, null, null, null)));5. 命名清晰结构采用 Arrange-Act-Assert类名使用*Test方法名明确表达行为或错误场景如上述shouldReturnJobIdOnlyWhenFinishedMetricsIsMissing、testValidConfigWithTablePath每个测试按 Arrange-Act-Assert 组织避免在测试方法里写大量准备逻辑可提取成小型复用 helper如JdbcSourceFactoryTest中的baseConfig()。单元测试检查清单类名符合*Test方法名可直接表达行为或错误场景单元测试保持确定性执行避免不必要的运行时依赖断言目标是行为和契约而非实现细节负向用例同时断言异常类型和关键报错信息测试数据准备最小且可复用E2E 测试规范一个 E2E 测试继承TestSuiteBase在 Testcontainers 管理的引擎容器内运行 connector并校验真实作业的结果。TestSuiteBase定义于 TestSuiteBase.java它通过ExtendWith组合了ContainerTestingExtension、TestLoggerExtension、TestCaseInvocationContextProvider、TimingExtension四个 JUnit 扩展并提供了共享的NETWORKTestContainer.NETWORK网络、Docker 客户端以及startContainerTcpProxy(...)代理工具。下面六条规则用于保证此类测试的确定性、无泄漏与低成本。1. 使用动态端口禁止硬编码Testcontainers 在启动时会分配一个随机的宿主机端口。硬编码端口会在 CI 中引发冲突——端口可能已被占用或者多个测试套件并行运行。应在运行时从容器解析宿主机端口// Bad — collides in CI String brokerUrl tcp://127.0.0.1:61616; // Good — let Testcontainers assign the host port String brokerUrl tcp:// container.getHost() : container.getMappedPort(61616); config.setPort(container.getFirstMappedPort()); // when only one port is exposed对所有外部服务都适用此规则// Database String jdbcUrl String.format( jdbc:mysql://%s:%d/%s, mysqlContainer.getHost(), mysqlContainer.getMappedPort(3306), DATABASE); // HTTP service String endpoint String.format( http://%s:%d, serviceContainer.getHost(), serviceContainer.getMappedPort(8080));:::caution 作业配置文件是个例外SeaTunnel 作业运行在 Docker 网络内部因此其配置.conf必须引用容器的网络别名withNetworkAliases(activemq-host)和内部端口如61616而不是映射后的宿主机端口。:::2. 基于条件等待禁止Thread.sleepThread.sleep是非确定性的时间太短测试会不稳定太长则拖慢 CI。当期望状态始终未到达时它也无法给出有用的错误信息。应使用 Awaitility 轮询真实条件。// Bad — flaky, wasteful, and silent on failure container.executeJob(/job.conf); Thread.sleep(30000); Assertions.assertIterableEquals(expected, query()); // Good — returns as soon as the condition holds, fails with the last assertion error Awaitility.await() .atMost(60, TimeUnit.SECONDS) .pollInterval(2, TimeUnit.SECONDS) .pollDelay(Duration.ZERO) .untilAsserted(() - Assertions.assertIterableEquals(expected, query()));应根据场景选择超时时间而非随手填一个整数场景atMostpollInterval原因容器 / 客户端就绪2 分钟1 秒镜像拉取 服务初始化较慢作业进入RUNNING状态1 分钟2 秒调度开销批作业结果校验60 秒2 秒小数据集完成Kafka / MQ 消息消费30–60 秒1 秒消费组再平衡CDC / Schema 变更传播60–120 秒2–5 秒binlog 延迟 快照当客户端在服务就绪前会抛异常时加上.ignoreExceptions()Awaitility.given() .ignoreExceptions() .pollInterval(500, TimeUnit.MILLISECONDS) .atMost(180, TimeUnit.SECONDS) .untilAsserted(this::initProducer);一个可复用的「等待作业 RUNNING」辅助方法private void awaitJobRunning(TestContainer container, String jobId) { Awaitility.await() .pollInterval(2, TimeUnit.SECONDS) .atMost(1, TimeUnit.MINUTES) .untilAsserted( () - Assertions.assertEquals(RUNNING, container.getJobStatus(jobId))); }从源码看getJobStatus(jobId)是SeaTunnelContainer提供的标准接口见 SeaTunnelContainer.java配合executeJob(confFile, jobId)、cancelJob(jobId)即可完整控制作业生命周期。3. 及时释放资源泄漏的连接会耗尽容器资源和宿主机可用端口并可能导致 CI 挂起。你打开的每一个资源都必须关闭。大多数 E2E 测试会实现TestResource接口并重写其tearDown()方法该方法在类中所有测试结束后执行一次。应在其中按创建的逆序关闭资源并对每个资源做 null 检查测试可能在资源创建前就已失败同时显式停止手动启动的容器。AfterAll Override public void tearDown() throws Exception { if (producer ! null) { producer.close(); } if (session ! null) { session.close(); } if (connection ! null) { connection.close(); } if (container ! null) { container.stop(); } }对于仅在单个方法内使用的资源优先使用 try-with-resources 而非手动 closeprivate void executeDml(String sql) { try (Connection conn getJdbcConnection(); Statement stmt conn.createStatement()) { stmt.execute(sql); } catch (SQLException e) { throw new RuntimeException(Execute DML failed: sql, e); } }资源清理检查清单JDBC 连接 / Statement 已关闭消息中间件的 connection、session、producer、consumer 已关闭自定义客户端HTTP、gRPC已关闭ExecutorService带超时关闭异步的CompletableFuture作业已取消临时文件 / 目录已删除手动创建的 Docker 网络已移除4. 异步提交长时间运行的作业流作业或 CDC 作业会一直运行直到被取消因此内联调用executeJob会使测试线程无限期阻塞也就没有机会注入数据或断言中间状态。这类作业应使用CompletableFuture.supplyAsync提交等待作业进入RUNNING状态执行测试步骤校验结果最后取消作业。String jobId streaming-cdc-job; CompletableFutureVoid job CompletableFuture.supplyAsync( () - { try { container.executeJob(/streaming_job.conf, jobId); } catch (Exception e) { log.error(Job execution failed, e); throw new CompletionException(e); // propagate, never swallow } return null; }); awaitJobRunning(container, jobId); // wait for RUNNING before proceeding insertCdcData(); // run test steps while the job is running Awaitility.await() .atMost(60, TimeUnit.SECONDS) .pollInterval(2, TimeUnit.SECONDS) .untilAsserted(() - verifySinkResults()); container.cancelJob(jobId); // cancel the job用下表判断作业类型模式带数据注入的流作业CompletableFuture.supplyAsync流式 CDC 作业CompletableFuture.supplyAsync动态分区发现CompletableFuture.supplyAsync批作业source → sink 会自然结束内联executeJob对于批作业断言退出码并校验结果Container.ExecResult result container.executeJob(/batch_job.conf); Assertions.assertEquals(0, result.getExitCode(), result.getStderr()); Awaitility.await().atMost(30, TimeUnit.SECONDS).untilAsserted(() - verifyResults());5. 保持测试确定且低成本继承TestSuiteBase并实现TestResource接口以获得startUp()/tearDown()生命周期钩子使用带TestContainer参数的TestTemplate使测试在每个已配置的引擎如 Zeta、Flink、Spark上运行启动尽可能少的容器。把 source 和 sink 的用例放在同一个类中让 Docker 镜像只初始化一次参见编码指南第 11 条覆盖完整的数据类型集合不要只测String和int用TestMethodOrder(MethodOrderer.OrderAnnotation.class)和Order(n)对有依赖的步骤排序加上Slf4j并记录关键阶段数据插入、校验的日志便于排查 CI 失败。6. 将无法回收的客户端线程加入白名单Zeta 引擎测试中最后一个作业结束后Zeta 的SeaTunnelContainer会对服务端 JVM 的存活线程做一次快照如果在 120 秒宽限期之后仍有非系统线程在运行就会让测试失败。这是为了防止 connector 泄漏线程。大多数客户端线程在你于tearDown()中关闭客户端后规则 3就会消亡——请务必先尝试这种方式。源码佐证SeaTunnelContainer.doExecuteJob在remainingJobs 0所有作业结束时会先取作业前的线程快照beforeThreads通过removeSystemThread过滤再用Awaitility以 120 秒为上限轮询等待残留线程清空否则断言失败见 SeaTunnelContainer.java 的线程检查逻辑。但有些第三方客户端库JDBC 驱动、HTTP 客户端、连接池会启动永远不会被回收、且测试无法关闭的守护线程。这类线程会让泄漏检查失败而这并非测试本身的问题。仅针对这种情况可以通过扩展SeaTunnelContainer中的isIssueWeAlreadyKnow(String)方法把该线程的名称前缀加入白名单——项目将这些例外集中维护于此// seatunnel-e2e/seatunnel-e2e-common/src/test/java/org/apache/seatunnel/e2e/common/ // container/seatunnel/SeaTunnelContainer.java /** The thread should be recycled but not, we should fix it in the future. */ protected boolean isIssueWeAlreadyKnow(String threadName) { return threadName.startsWith(ClickHouseClientWorker) // ClickHouse client || threadName.startsWith(Okio Watchdog) // InfluxDB / OkHttp || threadName.startsWith(iceberg-worker-pool) // Iceberg worker pool // Add the new connectors thread prefix here, with a comment naming the library. || threadName.startsWith(your-connector-thread-prefix); }从当前仓库源码看实际维护的白名单远不止文档示例中的三项已经覆盖了 Couchbase SDKSimplePauseDetectorThread、dnsjava NIO selector、cb-cleaner且parallel-N仅在 Couchbase E2E 生命周期内豁免、Azure Queue 的 Reactor Netty、GCS 的 OpenCensus、ClickHouseClickHouseClientWorker、InfluxDB/OkHttpOkio Watchdog、OkHttp TaskRunner、IoTDBSessionExecutor、Icebergiceberg-worker-pool等、Oracle JDBC 驱动BlockReleaser、RocketMQAsyncAppender-Dispatcher-Thread、MongoDBBufferPoolPruner、MaintenanceTimer、cluster-、JNA Cleaner、gRPC、PaimonAsyncOutputStream、MANIFEST-READ-THREAD-POOL等。这印证了文档中的三条要求。添加白名单条目时匹配具体的前缀绝不要用宽泛的子串——过松的匹配会掩盖其他测试的真实泄漏按现有风格加上注释说明所属 connector 以及启动该线程的全限定类名只对确实无法关闭的线程加白名单。能够关闭的线程应放在tearDown()中处理而不是白名单。:::note引擎与 JVM 系统线程hz.main、seatunnel-coordinator-service、pool-N-thread-N、GC 线程等已由isSystemThread(...)排除请勿在isIssueWeAlreadyKnow(...)中重复添加。:::从源码看isSystemThread的排除清单非常细致除了文档提到的三类还包括seatunnel-metrics-fetch-、pending-job-schedule-runner、SeaTunnel-CompletableFuture-Thread-、Connector-Scheduler-、heartbeat、LeaseRenewerHDFS、commons-pool-evictorRedis 连接池、process reaper、Log4j2-TF-等凡是引擎与 JVM 自身的线程都已统一在此过滤。整体示例下面的示例综合应用了上述规则测试方法共用一个固定版本镜像的容器、使用动态宿主机端口、在启动阶段基于条件等待、异步处理作业并在tearDown()中清理资源。Slf4j TestMethodOrder(MethodOrderer.OrderAnnotation.class) public class ExampleConnectorIT extends TestSuiteBase implements TestResource { private static final String CONTAINER_IMAGE example/example-server:1.2.3; private static final String CONTAINER_HOST example-host; private static final int CONTAINER_PORT 8080; private GenericContainer? container; private Connection connection; BeforeAll Override public void startUp() throws Exception { container new GenericContainer(DockerImageName.parse(CONTAINER_IMAGE)) .withNetwork(NETWORK) .withNetworkAliases(CONTAINER_HOST) .withExposedPorts(CONTAINER_PORT) .withLogConsumer( new Slf4jLogConsumer(DockerLoggerFactory.getLogger(CONTAINER_IMAGE))); Startables.deepStart(Stream.of(container)).join(); log.info(Example container started); // Wait until the client can connect, rather than sleeping for a fixed time. Awaitility.given() .ignoreExceptions() .await() .atMost(2, TimeUnit.MINUTES) .pollInterval(1, TimeUnit.SECONDS) .untilAsserted(this::initConnection); initializeTestData(); } AfterAll Override public void tearDown() throws Exception { if (connection ! null) { connection.close(); } if (container ! null) { container.stop(); } } TestTemplate Order(1) public void testBatchSourceToSink(TestContainer container) throws Exception { Container.ExecResult result container.executeJob(/example_batch.conf); Assertions.assertEquals(0, result.getExitCode(), result.getStderr()); Awaitility.await() .atMost(30, TimeUnit.SECONDS) .pollInterval(2, TimeUnit.SECONDS) .untilAsserted(this::verifyResults); } private void initConnection() { String endpoint String.format( http://%s:%d, container.getHost(), container.getMappedPort(CONTAINER_PORT)); connection createConnection(endpoint); } }模块与资源结构一个 connector 的 E2E 测试是位于seatunnel-e2e/seatunnel-connector-v2-e2e/下的独立 Maven 模块。引擎从测试的 classpath 根目录src/test/resources/加载.conf及其他资源文件因此这些路径是必须遵循的约定而非随意摆放。seatunnel-e2e/seatunnel-connector-v2-e2e/ └── connector-name-e2e/ # 每个 connector 一个模块 ├── pom.xml # artifactId: connector-name-e2e └── src/test/ ├── java/org/apache/seatunnel/e2e/connector/name/ │ └── NameIT.java # IT 测试类 └── resources/ ├── name_source_to_sink.conf # 作业配置classpath 根目录 ├── ddl/table.sql # 建表 / 初始化 DDL可选 ├── docker/server.cnf # 挂载到容器的配置可选 └── log4j2-test.properties # 测试日志可选仓库中 seatunnel-e2e/seatunnel-connector-v2-e2e 下的近百个connector-*-e2e模块均遵循该目录约定可对照connector-kafka-e2e、connector-jdbc-e2e等模块查看实际落地的文件排布。测试类位置src/test/java/org/apache/seatunnel/e2e/connector/name/。CDC connector 改用org.apache.seatunnel.connectors.seatunnel.cdc.name包——请与同级的 CDC 模块保持一致。命名测试类命名为NameIT。Failsafe 插件以集成测试方式运行*IT类而 Surefire 以单元测试方式运行*Test类。每个行为一个类例如NameWithSchemaChangeIT。配置.conf文件作业配置直接放在src/test/resources/下并在代码中用前导斜线引用container.executeJob(/name_source_to_sink.conf)。前导斜线表示 classpath 绝对路径而非文件系统路径。按用途与数据流向命名例如fake_to_sink.conf、name_multitable.conf或namecdc_to_name_with_schema_change.conf。避免使用test.conf这类泛化名称。在.conf内部服务必须使用容器的网络别名和内部端口如host rabbitmq-e2e、port 5672而不是映射后的宿主机端口参见 E2E 规则 1。.conf文件属于源文件必须以#注释形式带上 Apache License 头。DDL / SQL 及挂载资源建表与初始化脚本放在src/test/resources/ddl/下按其所建的数据库或表命名如ddl/inventory.sql。UniqueDatabase等 CDC 辅助类按 classpath 名称ddl/template.sql解析请保持该路径不变。挂载到容器的文件如my.cnf、setup.sql放在src/test/resources/docker/下通过容器 API 引用如.withSetupSQL(docker/setup.sql)。每个新增的.sql文件同样必须带上 Apache License 头。注册模块新模块在被纳入构建前不会被编译在seatunnel-e2e/seatunnel-connector-v2-e2e/pom.xml中加入moduleconnector-name-e2e/module。在模块的pom.xml中将artifactId设为connector-name-e2e给定形如SeaTunnel : E2E : Connector V2 : Name的name并按作业配置所需声明对被测 connector 以及connector-console、connector-assert、connector-fake的测试依赖。对于全新的 connector还需更新 CI label 和 plugin-mapping.properties参见贡献检查清单。提交前验证# Format ./mvnw spotless:apply # Run a specific E2E test class (run more than once to confirm stability) ./mvnw -pl seatunnel-e2e/seatunnel-connector-v2-e2e/connector-name-e2e \ -DskipUT -DskipITfalse -Dit.testTestClassName verify注意-pl指定的是模块路径需按实际模块名替换connector-name-e2e-DskipUT跳过单元测试阶段、-DskipITfalse显式开启集成测试、-Dit.test精确指定要运行的 IT 类。由于 E2E 测试依赖 Docker 环境与 Testcontainers请确保运行机器上 Docker 可用。建议对同一测试类连续运行多次以确认其稳定性确定性是合格 E2E 测试的基本要求。检查清单UT 遵循*Test命名、确定性执行和依赖隔离UT 负向用例同时断言异常类型和关键报错信息无硬编码宿主机端口——全部使用getMappedPort()/getFirstMappedPort()无Thread.sleep()——所有等待都使用 Awaitility并按场景设定超时每个打开的资源都在AfterAll/AfterEach或 try-with-resources 中关闭流 / CDC 作业通过CompletableFuture.supplyAsync提交并以RUNNING为前置条件source 与 sink 共用一个容器覆盖完整数据类型资源遵循约定结构.conf位于 classpath 根目录DDL 在ddl/下挂载文件在docker/下新模块已在父pom.xml中注册无 Zeta 线程泄漏失败无法回收的客户端线程已在isIssueWeAlreadyKnow中加入白名单所有新文件含.conf和.sql都带有 Apache License 头以上检查清单连同文中的六条 E2E 规则、五条单元测试规范与资源目录约定共同构成了 SeaTunnel 测试稳定性的完整基线。实际提交 PR 时可以对照 contributing 相关文档 与编码指南 进一步确认整体流程。【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考