
Apache SeaTunnel 测试套件实战指南编写稳定、无泄漏、确定性的 E2E 与单元测试【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel导读Apache SeaTunnel 是一个多模态、高性能、分布式的海量数据集成工具其连接器数量庞大、运行引擎多样Zeta、Flink、Spark因此对测试的质量要求极高一个测试若写得不够规范轻则静默地什么都没测重则在 CI 上随机抖动、泄漏线程导致构建失败。本文以仓库中 .skills/seatunnel-test-suite/SKILL.md 这份测试套件操作手册为骨架结合seatunnel-e2e模块的真实源码系统讲解 SeaTunnel E2ETestcontainers与单元测试的编写规范、审查流程和快速修复方法。读完本文你将掌握 E1–E7E2E与 U1–U6单元测试两套可引用、可执行的规则体系以及如何在本地用mvnw稳定复跑指定测试类。适用范围何时使用这套规范seatunnel-test-suite技能主要面向以下三类场景覆盖了 SeaTunnel 连接器测试生命周期中的新建—审查—修复新建编写一个新的连接器 E2E 测试*IT.java或单元测试*Test.java审查在提交 PR 前审查测试 diff检查是否存在抖动flaky、资源泄漏leak或布局违规修复稳定一个已经抖动或泄漏的测试或把旧测试升级到本规范约定的风格。核心原则是测试必须稳定、无泄漏、确定性stable, leak-free, deterministic。这三个词贯穿全文所有规则。E2E 框架契约为什么Test会让测试静默失效E2E 测试集成测试由 Maven Failsafe 插件运行必须严格遵循以下结构性契约。违反其中任何一条测试都会静默地什么都不做——而不是报错。要求细节继承TestSuiteBase提供引擎容器的生命周期、NETWORK与测试分发能力TestInstance(Lifecycle.PER_CLASS)从TestSuiteBase继承所有测试方法共享同一个类实例字段无需static测试方法用TestTemplate不能用Test。TestCaseInvocationContextProvider只分发TestTemplate方法TestContainer参数每个TestTemplate方法必须接收唯一的TestContainer container参数DisabledOnContainer可选仅当某个场景在某引擎上确实无法运行时才添加默认不加注解让测试在所有引擎上运行这些契约在源码中有明确印证。在 TestSuiteBase.java 中基类通过ExtendWith注册了四个 JUnit 扩展ContainerTestingExtension、TestLoggerExtension、TestCaseInvocationContextProvider和TimingExtension并声明了TestInstance(TestInstance.Lifecycle.PER_CLASS)与protected static final Network NETWORK TestContainer.NETWORK。分发机制的关键在 TestCaseInvocationContextProvider.javaOverride public boolean supportsTestTemplate(ExtensionContext context) { // Only support test cases with TestContainer as parameter Class?[] parameterTypes context.getRequiredTestMethod().getParameterTypes(); return parameterTypes.length 1 Arrays.stream(parameterTypes).anyMatch(TestContainer.class::isAssignableFrom); }也就是说只有参数数量为 1、且参数类型是TestContainer的方法才会被该 Provider 识别并分发到每个引擎容器上执行。若你误用了Test方法根本不会被分发到任何引擎容器测试通过只是假象。最小骨架一个可在所有引擎上运行的 E2E 测试下面是最小可用骨架。它对每个引擎Zeta、Flink、Spark都会执行一次这也是大多数批式 source/sink 测试的默认且正确的形态public class MyConnectorIT extends TestSuiteBase { TestTemplate public void testSourceToSink(TestContainer container) throws Exception { Container.ExecResult result container.executeJob(/my_connector_to_assert.conf); Assertions.assertEquals(0, result.getExitCode()); } }executeJob(/xxx.conf)的路径以前导斜杠引用文件实际位于测试类路径根目录src/test/resources/下详见下文布局规范。何时才应该添加DisabledOnContainer不要默认排除某个引擎。只有当某个场景在特定引擎上确实无法运行时才添加该注解且disabledReason要针对具体场景措辞而不是笼统写引擎 X 不支持。排除可以放在类上排除整个套件也可以放在单个TestTemplate方法上只排除该用例。仓库中真实存在的原因如下排除的引擎真实约束原因SparkCDC / 流式任务——Spark 不支持连续 / changelog 模式Spark Flinkcheckpoint 恢复类测试需要 Zeta 独有特性的场景Spark Flink仅剩 Zeta持续发现的长时间运行任务仅 SeaTunnel 自身的行为SparkSpark 会丢弃记录的 RowKind导致 changelog 断言失败注解定义见 DisabledOnContainer.java支持valueTestContainerId[]、typeEngineType[]与disabledReason三个属性且可标注在类或方法上。// 只有这个 CDC 用例无法在 Spark 上运行——排除放在方法上而不是类上 TestTemplate DisabledOnContainer( value {}, type {EngineType.SPARK}, disabledReason Spark does not support the CDC streaming job) public void testCdcStreaming(TestContainer container) throws Exception { // ... }规则总览E1–E7 与 U1–U6 两套命名空间规则分为两个命名空间必须与文件类型严格对应E1–E7约束 E2E*IT测试引用格式[E3]U1–U6约束单元*Test测试引用格式[U5]E6 仅适用于 Zeta 引擎。审查时永远选择与待审文件匹配的命名空间——绝不把 E 规则套到单元测试上反之亦然。#规则检查点坏味道修复方式E1动态端口Java 代码中出现字面量127.0.0.1:port/localhost:port使用container.getHost()container.getMappedPort(p).conf文件引用网络别名 内部端口见 E7E2基于条件的等待任何Thread.sleep(...)使用Awaitility.await().atMost(...).pollInterval(...).untilAsserted(...)超时取自下方场景表E3释放资源打开的 client/connection/container 没有 close实现TestResource在tearDown()中逆序 空安全关闭方法级资源用 try-with-resourcesE4异步任务提交对流式/CDC 任务内联调用executeJob永不返回用CompletableFuture.supplyAsync(...)以RUNNING为门控执行动作、验证后cancel(true)E5共享单个容器每个测试方法都新起一个容器在一个类内所有方法间共享一个容器TestInstance(PER_CLASS)E5b固定镜像版本使用:latest镜像标签固定到具体版本如mysql:8.0.32绝不用mysql:latest见子节E5c覆盖全部数据类型只测了String/int覆盖连接器的完整数据类型集合而不是只测简单的E6线程白名单Zeta报错There are still threads running in the container先关闭 clientE3若库线程确实无法回收才在isIssueWeAlreadyKnow(...)中白名单其前缀E7Docker 网络引擎容器无法访问外部服务容器.withNetwork(NETWORK).withNetworkAliases(my-alias)用Startables.deepStart(...)启动添加Slf4jLogConsumer(DockerLoggerFactory.getLogger(IMAGE))获取调试日志完整设置见下E7网络设置——引擎在容器里服务容器必须同网SeaTunnel 引擎运行在它自己的 Docker 容器内部。外部服务容器必须与引擎共享同一个 Docker 网络引擎才能访问到它们——这也是 E1 中.conf文件使用网络别名而非宿主机映射端口的原因。完整的设置模式如下// 在 IT 类的 BeforeAll 中NETWORK 从 TestSuiteBase 继承 container new GenericContainer(DockerImageName.parse(my-service:1.2.3)) .withNetwork(NETWORK) // 继承自 TestSuiteBase 的字段 .withNetworkAliases(my-service-host) // 引擎容器内用这个名字可达 .withExposedPorts(3306) .withLogConsumer(new Slf4jLogConsumer(DockerLoggerFactory.getLogger(my-service:1.2.3))); Startables.deepStart(Stream.of(container)).join(); // 并行启动在.conf文件里在引擎容器内部执行host my-service-host port 3306在 Java 测试设置代码里在宿主机上执行String jdbcUrl String.format(jdbc:mysql://%s:%d/test, container.getHost(), container.getMappedPort(3306));三层各用各的地址JVM 测试代码用宿主机映射地址job 配置用网络别名 内部端口。混用二者是 E2E 测试最常见的失败根源之一。E2超时参考表用 Awaitility 替换Thread.sleep时atMost与pollInterval不能随手填整数必须从下表按场景取值场景atMostpollInterval容器 / 客户端就绪2 min1 s任务进入RUNNING1 min2 s批式任务结果验证60 s2 sMQ / Kafka 消费30–60 s1 sCDC / schema 变更传播60–120 s2–5 s另外当客户端在服务起来之前会抛异常时务必在Awaitility.await()上链式调用.ignoreExceptions()这是 E2 的标准修复——否则第一次轮询的异常会直接让等待失败而不是继续重试。绝不要选一个没有场景支撑的光秃秃的整数作为超时。E5b镜像版本固定绝不使用:latest——上游镜像可能在没有代码 diff 可定位的情况下悄悄改变行为让 CI 随机挂掉// 坏——不可预测地破裂 private static final String IMAGE postgres:latest; // 好——可复现 private static final String IMAGE postgres:14.5;E6线程白名单一段话讲透最后一个任务结束后SeaTunnelContainer会对服务器 JVM 做快照若有非系统线程存活超过 120 秒则测试失败。第三方客户端的守护线程JDBC 驱动、HTTP 连接池有时确实无法回收。只针对这些线程在SeaTunnelContainer.isIssueWeAlreadyKnow(String)中添加一个具体的前缀并注释注明库名。一个你能够关闭的线程应该放在tearDown()里关闭而不是进白名单。系统线程hz.main、pool-N-thread-N等已由isSystemThread(...)处理不要重复添加。这条规则的源码实现在 SeaTunnelContainer.java 的doExecuteJob中任务执行完成后先取前后两次 JVM 线程快照剔除系统线程与执行前已存在的线程若仍有剩余则用Awaitility.await().atMost(120, TimeUnit.SECONDS).untilAsserted(...)等待线程释放并输出There are still threads running in the container: ...的详细线程堆栈。值得注意的实践细节见 SeaTunnelContainer.javaisSystemThread维护了一份很长的系统线程前缀清单hz.main、pool-N-thread-N正则、seatunnel-coordinator-service、GC task thread、OkHttp ConnectionPool、qtp等而isIssueWeAlreadyKnow中的白名单则精确到库级例如 Couchbase SDK 的SimplePauseDetectorThread/dnsjava NIO selector/cb-cleaner、ClickHouse 的ClickHouseClientWorker、InfluxDB 的Okio Watchdog、Oracle 驱动的BlockReleaser、RocketMQ 的AsyncAppender-Dispatcher-Thread、Paimon 的MANIFEST-READ-THREAD-POOL等。代码还展示了一个更精细的防护思路对于线程名与其它连接器共享的场景如 Reactor 的parallel-N、boundedElastic-evictor-、OpenCensus 的导出线程仓库通过enableXxxThreadExemption()/disableXxxThreadExemption()配合couchbaseE2eActive、azureQueueE2eActive、gcsE2eActive等生命周期标志把豁免限定在对应 E2E 测试存活期间避免静默掩盖其它连接器的真实泄漏。单元测试规则U1–U6*Test类走 Surefire不涉及引擎容器。用下面这张表来审查#规则坏味道修复方式U1测试行为而非内部实现断言私有字段或内部调用顺序断言可观察的输出/状态仅当副作用本身就是契约时才验证void调用U2无时序或 IO 依赖Thread.sleep、System.currentTimeMillis()、文件系统、网络移除时序mock/伪造 IO 边界必要时注入时钟U3只在边界处 mockmock 领域对象或可用new创建的值类型只 mock 类通过构造器或 setter 依赖的外部接口/client/map/服务U4一个方法只测一个行为单个测试内有多个分支的if/for拆成一个分支一个方法方法名即其证明的行为U5负向路径要断言消息assertThrows(FooException.class, ...)但不检查消息追加assertThat(ex.getMessage()).contains(expected fragment)U6命名规范testFoo、testBar、test123类名*Test方法名如shouldReturn404WhenUserMissing、throwsOnNullHostMock 设置的标准形态Arrange-Act-Assert每个测试方法都应该是 Arrange-Act-Assert 三段式Arrange——在BeforeEach中对每个外部边界执行mock(...)用when(...).thenReturn(...)把调用链接好然后构造被测对象SUT。每个测试只 stub 能选中目标分支的那些值。Act——只调用被测的那个方法。Assert——断言可观察的输出U1而不是调用链。让 mock 测试错误地通过的三个坑链路 stub当 SUT 通过调用链访问协作者a.getB().getC(key)时要 stub 整条链且每次查找都用生产代码里的精确 key 常量——key 不匹配会返回null并悄悄改变所走的分支上游守卫要走到下游分支必须先把所有上游守卫 stub 成非空——上游值为null会提前短路只 stub 下游值等于什么都没测负向路径断言异常消息片段而不只是类型U5——否则同类型但原因错误的异常也会通过。何时不要 mock简单的值对象、POJO、DTO——直接用真实实例被测类本身——永远不要 mock SUT当内存伪实现比如用HashMap代替接口更简单且隔离程度相同时。布局规范测试文件该放哪每个连接器对应一个 Maven 模块位于seatunnel-e2e/seatunnel-connector-v2-e2e/connector-name-e2e/。文件从测试类路径根目录src/test/resources/加载测试类→src/test/java/org/apache/seatunnel/e2e/connector/name/NameIT.javaCDC 连接器使用...connectors.seatunnel.cdc.name包。*IT 集成测试Failsafe*Test 单元测试Surefire任务配置→src/test/resources/name_source_to_sink.conf用前导斜杠引用executeJob(/name_source_to_sink.conf)。每个新文件都要有 Apache License 头#注释DDL→src/test/resources/ddl/table.sqlCDC 的UniqueDatabase会解析ddl/template.sql容器挂载→src/test/resources/docker/...通过容器 API 引用无前导斜杠模块注册→ 在父级pom.xml中注册该模块新连接器还需更新 CI 标签和plugin-mapping.properties。审查工作流10 步检查清单用 grep 扫明显坏味道结果干净只是提示不是保证——127\.0\.0\.1和localhost要连端口一起匹配以捕捉端口被单独模板化的宿主引用grep -nE Thread\.sleep|127\.0\.0\.1|localhost|:latest test-file[E1, E7]外部容器是否加入了NETWORK并带别名.conf是否使用别名 内部端口框架契约是否继承TestSuiteBase是否是TestTemplateTestContainer参数而不是TestDisabledOnContainer——只出现在确实无法运行的场景无笼统排除且当只有一个用例不兼容时是否收敛到该方法的注解上[E3]每个打开的资源是否在tearDown()中关闭或 try-with-resources容器是否显式停止[E4]流式/CDC 任务是否通过CompletableFuture.supplyAsync提交并门控在RUNNING[E2]超时是否有场景表支撑——没有光秃秃的整数[E5, E5b, E5c]每个类是否共享一个容器镜像版本是否固定无:latest是否覆盖完整类型集资源是否放在规范路径新模块是否注册到父 pom[E6]对于 Zeta E2E跑完整任务集若线程泄漏检查触发先关闭 client。单元测试逐条走[U1–U6]表负向路径确认[U5]断言了消息。多次重复运行目标测试类确认稳定性./mvnw -pl seatunnel-e2e/seatunnel-connector-v2-e2e/connector-name-e2e \ -DskipUT -DskipITfalse -Dit.testITClassName verify快速修复速查表审查后每一条发现都映射到下面的速查表直接对应输出格式中的规则 ID。规则坏味道替换为E2Thread.sleep(n)后断言Awaitility.await().atMost(...).untilAsserted(...)E1tcp://localhost:61616container.getHost() : container.getMappedPort(61616)E3缺少tearDown()实现TestResource空安全、逆序关闭 container.stop()E4流式任务内联executeJobCompletableFuture.supplyAsync(...) 门控RUNNINGE4consumer.receive()无超时consumer.receive(timeoutMs)或 Awaitility 轮询E4异步块内吞掉异常以CompletionException重新抛出E6Zeta E2E 因残留 client 线程失败在tearDown()中关闭若无法回收则在isIssueWeAlreadyKnow白名单—E2E 方法上用了Test改为TestTemplate 添加TestContainer container参数E5bpostgres:latestpostgres:14.5固定到具体版本E7引擎容器访问不到服务容器.withNetwork(NETWORK).withNetworkAliases(...)U5assertThrows(...)未检查消息追加assertTrue(ex.getMessage().contains(...))审查输出格式可引用、可定位审查测试时每条发现按如下格式输出RuleId对*IT测试取 E 规则E1–E7对*Test单元测试取 U 规则U1–U6[RuleId] file:line — issue例如[E3] FooIT.java:88 — Kafka consumer never closed。只有当修复不显然时才附上 before/after 代码片段。写新类时需满足所有适用规则并为每个新文件包括.conf和.sql添加 Apache License 头。与框架源码的对应关系TestSuiteBase.javaE2E 基类注册四大 JUnit 扩展、PER_CLASS生命周期、共享NETWORKTestCaseInvocationContextProvider.java只分发带TestContainer参数的TestTemplate方法并过滤被DisabledOnContainer排除的容器DisabledOnContainer.java类/方法级排除注解支持EngineType[]与disabledReasonTestResource.javaE3 规则中可关闭资源的接口契约SeaTunnelContainer.javaZeta 引擎容器内置 120 秒线程泄漏检查、系统线程过滤与库级白名单实现。掌握了这套规则再配合 .skills/seatunnel-test-suite/SKILL.md 中的规则表无论是编写新连接器的 E2E 测试、审查 PR 中的测试 diff还是修复 CI 上反复出现的 flaky 用例你都有了完整的、可引用、可落地的操作依据。【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考