chl-wxp commented on code in PR #11212:
URL: https://github.com/apache/seatunnel/pull/11212#discussion_r3497628131
##########
docs/zh/developer/test-coding-guide.md:
##########
@@ -0,0 +1,518 @@
+---
+title: 测试编码指南
+---
+
+# 测试编码指南
+
+本指南介绍如何为 Apache SeaTunnel 编写一个高质量、稳定的端到端(E2E)测试或单元测试。
+它是对[编码指南](coding-guide.md)的补充:编码指南覆盖通用的 PR 质量要求,而本指南专注于测试的稳定性与资源安全。
+
+一个优秀的 SeaTunnel 测试应当是确定性的(每次运行结果一致)、无泄漏的(释放所有打开的资源)、
+低成本的(启动尽可能少的容器)。下面的规范将这些原则落实到 SeaTunnel 贡献者编写的两类测试上。
+
+SeaTunnel 将测试分为两层:
+
+- **单元测试**(`*Test`,由 Surefire 插件运行)在隔离环境中校验 connector 或 core 逻辑,不依赖外部系统。
+ 参见[单元测试规范](#单元测试规范)。
+- **端到端(E2E)测试**(`*IT`,由 Failsafe 插件运行)借助 Testcontainers,针对真实服务校验 source、
+ transform、sink 的行为。参见 [E2E 测试规范](#e2e-测试规范)。
+
+当一个 Pull Request 同时改动了逻辑与集成行为时,请在两层都补充测试,并在提交前分别运行。
+
+## 单元测试规范
+
+### 1. 测行为和契约,不测实现细节
+
+单元测试应验证输入输出行为和配置契约,不要绑定私有方法内部实现或临时代码结构。
+
+真实案例:`JdbcSourceFactoryTest`(模块:`connector-jdbc`)
+
+```java
+private Map<String, Object> baseConfig() {
+ Map<String, 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() {
+ Map<String, Object> cfg = baseConfig();
+ cfg.put("table_path", "test.users");
+ Assertions.assertDoesNotThrow(() -> validate(cfg));
+}
+```
+
+这个测试保护的是对外配置契约,即使工厂内部实现重构,测试仍然稳定。
+
+### 2. 单元测试必须确定且本地可跑
+
+单元测试应快速且确定:
+
+- 不使用 `Thread.sleep`
+- 不使用无固定种子的随机断言
+- 不依赖墙钟时间
+
+如果用例引入跨进程或外部依赖,请明确归类到对应测试层级,并保证断言可重复、可诊断。
+
+### 3. 用 mock、stub、fake 隔离依赖
+
+单元测试要把目标行为与外部 IO 隔离。对协作依赖优先使用内存实现或轻量测试替身。
+
+每个测试方法尽量只覆盖一个断言范围,失败时能直接定位单一行为。
+
+引擎侧真实 mock 案例:`JobInfoServiceNullSafetyTest`(模块:`seatunnel-engine-server`)
+
+```java
+@BeforeEach
+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 依赖,无需启动引擎服务;
+- 通过 `when(...).thenReturn(...)` 精准构造单一业务分支;
+- 断言只关注可观察输出(结果 JSON 契约),而非内部调用。
+
+### 4. 错误路径同时校验异常类型与关键报错信息
+
+负向用例不仅要断言异常类,还要断言关键报错内容,以保护用户可见错误质量。
+
+真实案例:`MongodbIncrementalSourceFactoryTest`(模块:`connector-cdc-mongodb`)
+
+```java
+Assertions.assertThrows(
+ MongodbConnectorException.class,
+ () ->
+ MongodbSourceConfigProvider.newBuilder()
+ .startupOptions(
+ new StartupConfig(StartupMode.EARLIEST, null,
null, null)));
+```
+
+### 5. 命名清晰,结构采用 Arrange-Act-Assert
+
+- 类名使用 `*Test`,方法名明确表达行为或错误场景。
+- 每个测试按 Arrange-Act-Assert 组织。
+- 避免在测试方法里写大量准备逻辑,可提取成小型复用 helper。
+
+## 单元测试检查清单
+
+- [ ] 类名符合 `*Test`,方法名可直接表达行为或错误场景
+- [ ] 单元测试保持确定性执行,避免不必要的运行时依赖
+- [ ] 断言目标是行为和契约,而非实现细节
+- [ ] 负向用例同时断言异常类型和关键报错信息
+- [ ] 测试数据准备最小且可复用
+
+## E2E 测试规范
+
+一个 E2E 测试继承 `TestSuiteBase`,在 Testcontainers 管理的引擎容器内运行 connector,并校验真实作业的结果。
+下面六条规则用于保证此类测试的确定性、无泄漏与低成本。
+
+### 1. 使用动态端口,禁止硬编码
+
+Testcontainers 在启动时会分配一个随机的宿主机端口。硬编码端口会在 CI 中引发冲突——端口可能已被占用,
+或者多个测试套件并行运行。
+
+应在运行时从容器解析宿主机端口:
+
+```java
+// 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
+```
+
+对所有外部服务都适用此规则:
+
+```java
+// 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.sleep`
+
+`Thread.sleep` 是非确定性的:时间太短测试会不稳定,太长则拖慢 CI。当期望状态始终未到达时,它也无法给出有用的
+错误信息。应使用 [Awaitility](https://github.com/awaitility/awaitility) 轮询真实条件。
+
+```java
+// 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()));
+```
+
+应根据场景选择超时时间,而非随手填一个整数:
+
+| 场景 | `atMost` | `pollInterval` | 原因
|
+|---------------------------|------------|----------------|--------------------------|
+| 容器 / 客户端就绪 | 2 分钟 | 1 秒 | 镜像拉取 + 服务初始化较慢 |
+| 作业进入 `RUNNING` 状态 | 1 分钟 | 2 秒 | 调度开销
|
+| 批作业结果校验 | 60 秒 | 2 秒 | 小数据集完成 |
+| Kafka / MQ 消息消费 | 30–60 秒 | 1 秒 | 消费组再平衡
|
+| CDC / Schema 变更传播 | 60–120 秒 | 2–5 秒 | binlog 延迟 + 快照
|
+
+当客户端在服务就绪前会抛异常时,加上 `.ignoreExceptions()`:
+
+```java
+Awaitility.given()
+ .ignoreExceptions()
+ .pollInterval(500, TimeUnit.MILLISECONDS)
+ .atMost(180, TimeUnit.SECONDS)
+ .untilAsserted(this::initProducer);
+```
+
+一个可复用的「等待作业 RUNNING」辅助方法:
+
+```java
+private void awaitJobRunning(TestContainer container, String jobId) {
+ Awaitility.await()
+ .pollInterval(2, TimeUnit.SECONDS)
+ .atMost(1, TimeUnit.MINUTES)
+ .untilAsserted(
+ () -> Assertions.assertEquals("RUNNING",
container.getJobStatus(jobId)));
+}
+```
+
+### 3. 及时释放资源
+
+泄漏的连接会耗尽容器资源和宿主机可用端口,并可能导致 CI 挂起。你打开的每一个资源都必须关闭。大多数 E2E
+测试会实现 `TestResource` 接口并重写其 `tearDown()` 方法,该方法在类中所有测试结束后执行一次。应在其中按
+创建的逆序关闭资源,并对每个资源做 null 检查(测试可能在资源创建前就已失败),同时显式停止手动启动的容器。
+
+```java
+@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 而非手动 close:
+
+```java
+private 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` 状态,执行
+测试步骤,校验结果,最后取消作业。
+
+```java
+String jobId = "streaming-cdc-job";
+
+CompletableFuture<Void> 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());
+
+job.cancel(true); // cancel the job
Review Comment:
```suggestion
container.cancelJob(jobId); // cancel the job
```
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]