Skip to content
Closed
43 changes: 43 additions & 0 deletions .github/workflows/connector-it.yml
Original file line number Diff line number Diff line change
@@ -0,0 +1,43 @@
# Connector 集成测试(开源库触发入口)
#
# 触发时机:push 任意分支(含 PR 合并后的推送)+ PR 提交/更新 + 手工触发。
# paths 过滤保证仅连接器相关路径变更时启动,文档/README 等改动不触发。
# 实际测试逻辑在 tapdata-it 的 connector-it-reusable.yml 中统一维护。
# 运行环境:office-build 自托管 runner(runs-on 在 reusable workflow 中配置)。
name: Connector Integration Tests

on:
push:
branches: ['**']
paths:
- 'connectors/**'
- 'connectors-common/**'
- 'connectors-javascript/**'
- 'pom.xml'
- '.github/**'
pull_request:
types: [opened, synchronize, reopened]
paths:
- 'connectors/**'
- 'connectors-common/**'
- 'connectors-javascript/**'
- 'pom.xml'
- '.github/**'
workflow_dispatch:
inputs:
connectors:
description: '手工指定模块列表(逗号分隔,如 connectors/mysql-connector, connectors/mongodb-connector)'
required: false
type: string

concurrency:
group: connector-it-${{ github.workflow }}-${{ github.ref }}
cancel-in-progress: true

jobs:
connector-it:
uses: tapdata/tapdata-it/.github/workflows/connector-it-reusable.yml@main
with:
java_version: '17'
connectors_override: ${{ inputs.connectors || '' }}
secrets: inherit
2 changes: 1 addition & 1 deletion connectors-common/mysql-core/pom.xml
Original file line number Diff line number Diff line change
Expand Up @@ -126,7 +126,7 @@
<dependency>
<groupId>io.tapdata</groupId>
<artifactId>tapdata-common</artifactId>
<version>2.0-SNAPSHOT</version>
<version>2.9-SNAPSHOT</version>
<scope>provided</scope>
</dependency>
<dependency>
Expand Down
67 changes: 67 additions & 0 deletions connectors/mongodb-connector/pom.xml
Original file line number Diff line number Diff line change
Expand Up @@ -15,6 +15,7 @@
<java.version>8</java.version>
<tapdata.pdk.api.version>2.0.7-SNAPSHOT</tapdata.pdk.api.version>
<bouncycastle.version>1.84</bouncycastle.version>
<skipITs>true</skipITs>
</properties>
<dependencies>
<dependency>
Expand Down Expand Up @@ -116,6 +117,13 @@
</dependency>


<!-- 通用 Connector 集成测试框架(test scope,提供 ConnectorIT 基类与全部通用用例) -->
<dependency>
<groupId>io.tapdata</groupId>
<artifactId>tapdata-connector-it</artifactId>
<version>1.0-SNAPSHOT</version>
<scope>test</scope>
</dependency>
<!-- Mockito -->
<dependency>
<groupId>org.mockito</groupId>
Expand Down Expand Up @@ -289,6 +297,65 @@
</execution>
</executions>
</plugin>
<plugin>
<groupId>org.codehaus.mojo</groupId>
<artifactId>build-helper-maven-plugin</artifactId>
<version>3.2.0</version>
<executions>
<execution>
<id>add-integration-test-source</id>
<phase>generate-test-sources</phase>
<goals>
<goal>add-test-source</goal>
<goal>add-test-resource</goal>
</goals>
<configuration>
<sources>
<source>src/it/java</source>
</sources>
<resources>
<resource>
<directory>src/it/resources</directory>
</resource>
</resources>
</configuration>
</execution>
</executions>
</plugin>
<!-- 集成测试(src/it/java,*IT 命名,与单测 surefire 分离) -->
<plugin>
<groupId>org.apache.maven.plugins</groupId>
<artifactId>maven-failsafe-plugin</artifactId>
<version>3.5.6</version>
<configuration>
<skipITs>${skipITs}</skipITs>
<!-- MongodbWriter 依赖 AppType(引擎启动参数 app_type),fork JVM 需显式设置,
否则构造 writer 时抛 "app_type is blank" -->
<argLine>
-Xmx4g
-Dapp_type=DAAS
--add-opens=java.base/java.lang=ALL-UNNAMED
--add-opens=java.base/java.lang.reflect=ALL-UNNAMED
--add-opens=java.base/java.util=ALL-UNNAMED
--add-opens=java.base/java.security=ALL-UNNAMED
--add-opens=java.base/java.io=ALL-UNNAMED
--add-opens=java.base/java.time=ALL-UNNAMED
--add-opens=java.base/jdk.internal.loader=ALL-UNNAMED
</argLine>
<reportsDirectory>
${project.build.directory}/failsafe-reports
</reportsDirectory>
</configuration>
<executions>
<execution>
<id>integration-test</id>
<goals>
<goal>integration-test</goal>
<goal>verify</goal>
</goals>
</execution>
</executions>
</plugin>
</plugins>
</build>
</project>
Original file line number Diff line number Diff line change
@@ -0,0 +1,139 @@
package io.tapdata.mongodb;

import io.tapdata.entity.codec.TapCodecsRegistry;
import io.tapdata.entity.logger.TapLog;
import io.tapdata.entity.utils.DataMap;
import io.tapdata.it.ConnectorIT;
import io.tapdata.it.ConnectorTestContext;
import io.tapdata.it.support.TestStateMap;
import io.tapdata.pdk.apis.context.TapConnectorContext;
import io.tapdata.pdk.apis.functions.ConnectorFunctions;
import org.bson.BsonRegularExpression;
import org.bson.BsonTimestamp;
import org.bson.types.Binary;
import org.bson.types.Code;
import org.bson.types.Decimal128;
import org.bson.types.ObjectId;
import org.bson.types.Symbol;
import org.junit.jupiter.api.Assertions;
import org.junit.jupiter.api.DisplayName;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.TestInstance;

import java.util.HashSet;
import java.util.LinkedHashMap;
import java.util.List;
import java.util.Map;
import java.util.Set;
import java.util.stream.Collectors;
import java.util.stream.Stream;

/**
* MongoDB 连接器通用集成测试。
* <p>
* 继承 {@link ConnectorIT} 后自动运行框架内置的全部通用集成测试用例
* (连接元数据、集合 DDL、数据读写、事务、流式读取、命令等),
* 对 MongoDB 不支持的 ConnectorFunctions 能力自动跳过。
* 必实现能力由 {@link #requiredCapabilities()} 主动声明(原则 3:声明式能力),
* 框架据此校验:声明必实现但未实现、或已实现但无用例覆盖 → 测试失败。
* <p>
* 连接配置读取 src/it/resources/config/mongodb-connection.json
* (uri/database),支持系统属性/环境变量覆盖
* (如 -Dconnector.it.uri=xxx 或 CONNECTOR_IT_URI=xxx)。
*/
@TestInstance(TestInstance.Lifecycle.PER_CLASS)
public class MongoDBConnectorIT extends ConnectorIT {

@Override
protected Set<String> requiredCapabilities() {
// MongoDB 同时承担源与目标角色:声明对外承诺的全部能力(与 registerCapabilities 注册保持一致;
// memoryFetcher 为引擎诊断钩子,由基类 ignoredCapabilities() 排除,不参与校验)
return Stream.of("connectionTest", "discoverSchema", "getTableNames",
"createTableV2", "dropTable", "batchCount", "batchRead", "streamRead",
"timestampToStreamOffset", "queryByAdvanceFilter", "countByPartitionFilter",
"writeRecord", "createIndex", "queryIndexes", "errorHandle",
"executeCommand", "getTableInfo", "getReadPartitions", "queryFieldMinMaxValue",
"transactionBegin", "transactionCommit", "transactionRollback").collect(Collectors.toSet());
}

@Override
protected ConnectorTestContext createContext() throws Throwable {
// 1. 创建连接器
MongodbConnector connector = new MongodbConnector();

// 2. 连接配置(URI 模式,与产品连接表单默认一致)
DataMap config = readConnectionConfig("config/mongodb-connection.json");

// 3. 构建 NodeContext:specification + connectionConfig + nodeConfig + log
// specification 提供 dataTypesMap,discoverSchema 的 tapType 自动填充依赖它
TapConnectorContext nodeContext = new TapConnectorContext(
loadSpecification("spec.json"), config, DataMap.create(), new TapLog());
nodeContext.setStateMap(new TestStateMap());

// 4. 注册能力(registerCapabilities 产物,基类据此动态检测并驱动用例)
ConnectorFunctions functions = new ConnectorFunctions();
TapCodecsRegistry codecRegistry = TapCodecsRegistry.create();
connector.registerCapabilities(functions, codecRegistry);

// 5. 组装测试上下文(特性开关:MongoDB NoSQL 特性适配)
ConnectorTestContext ctx = ConnectorTestContext.builder()
.connector(connector)
.nodeContext(nodeContext)
.connectorFunctions(functions)
.codecRegistry(codecRegistry)
.config(config)
.log(new TapLog())
// MongoDB createTableV2 幂等建集合固定返回 tableExists=false
.createTableReportsTableExists(false)
// schema-free:空集合无法推断字段,需先写入采样数据再 discoverSchema
.schemaDiscoveryRequiresSampleData(true)
// discoverSchema 会带出隐式 _id 字段
.schemaAllowsExtraFields(true)
// 主键为库自动生成的 _id,业务字段不被标记为主键
.schemaPrimaryKeyStrict(false)
// executeCommand 仅支持 execute/executeQuery/count/aggregate,不支持 ping
.executeCommandSupportsPing(false)
// queryFieldMinMaxValue 基于分区索引字段计算 min/max
.fieldMinMaxRequiresPartitionIndex(true)
.build();
ctx.getLog().info("[IT] MongoDB connection: uri={}, database={}",
config.getString("uri"), config.getString("database"));
return ctx;
}

@Override
protected List<Map<String, Object>> beforeWrite(List<Map<String, Object>> rows) {
// unique 索引用例(idx_c_int)要求 c_int 唯一:随机值冲突时改为超出随机范围的序列值
Set<Long> seen = new HashSet<>();
long seq = 1_000_001L;
for (Map<String, Object> row : rows) {
long v = ((Number) row.get("c_int")).longValue();
if (!seen.add(v)) {
row.put("c_int", seq++);
}
}
return rows;
}

@Override
protected Map<String, Object> specialValueSamples() {
// MongoDB 特有 BSON 类型样本:须与 registerCapabilities 注册的 ToTapValueCodec 一一对应
// (U8 用例验证这些特殊值 wrap 后被连接器 codec 识别,而非 TapRawValue 兜底)
Map<String, Object> samples = new LinkedHashMap<>();
samples.put("obj_objectid", new ObjectId());
samples.put("obj_binary", new Binary((byte) 0x80, new byte[]{1, 2, 3}));
samples.put("obj_code", new Code("function() { return 1; }"));
samples.put("obj_decimal128", Decimal128.parse("12345.6789"));
samples.put("obj_symbol", new Symbol("sym"));
samples.put("obj_bson_timestamp", new BsonTimestamp(1_700_000_000, 1));
samples.put("obj_regex", new BsonRegularExpression("^tap.*", "i"));
return samples;
}

@DisplayName("Test github action")
@Test
public void testGithubAction() {
int result = 10*5;
Assertions.assertEquals(50, result);
}
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,4 @@
{
"uri": "mongodb://root:Gotapd8@113.98.206.142:14017/test?authSource=admin",
"database": "test"
}
71 changes: 69 additions & 2 deletions connectors/mysql-connector/pom.xml
Original file line number Diff line number Diff line change
Expand Up @@ -22,8 +22,9 @@
</modules>

<properties>
<java.version>8</java.version>
<java.version>17</java.version>
<tapdata.pdk.api.version>2.0.5-SNAPSHOT</tapdata.pdk.api.version>
<skipITs>true</skipITs>
</properties>
<dependencies>
<dependency>
Expand Down Expand Up @@ -54,11 +55,26 @@
<version>5.8.1</version>
<scope>test</scope>
</dependency>
<!-- 通用 Connector 集成测试框架(test scope,提供 ConnectorIT 基类与全部通用用例) -->
<dependency>
<groupId>io.tapdata</groupId>
<artifactId>tapdata-connector-it</artifactId>
<version>1.0-SNAPSHOT</version>
<scope>test</scope>
</dependency>
<dependency>
<groupId>io.tapdata</groupId>
<artifactId>mysql-core</artifactId>
<version>1.0-SNAPSHOT</version>
</dependency>
<!-- mysql-core 中 MysqlReaderV2(highPerformance 路径)依赖 tapdata-common 的 ConcurrentProcessor/TapExecutors,
该依赖在 mysql-core 内为 provided(生产由引擎提供,不随 connector 打包);IT 测试运行时需显式补齐 -->
<dependency>
<groupId>io.tapdata</groupId>
<artifactId>tapdata-common</artifactId>
<version>2.9-SNAPSHOT</version>
<scope>test</scope>
</dependency>
<dependency>
<groupId>io.tapdata</groupId>
<artifactId>pdk-error-code</artifactId>
Expand Down Expand Up @@ -209,6 +225,24 @@
</sources>
</configuration>
</execution>
<execution>
<id>add-integration-test-source</id>
<phase>generate-test-sources</phase>
<goals>
<goal>add-test-source</goal>
<goal>add-test-resource</goal>
</goals>
<configuration>
<sources>
<source>src/it/java</source>
</sources>
<resources>
<resource>
<directory>src/it/resources</directory>
</resource>
</resources>
</configuration>
</execution>
</executions>
</plugin>
<plugin>
Expand All @@ -232,7 +266,40 @@
</execution>
</executions>
</plugin>
<!-- <plugin>-->
<!-- 集成测试(src/it/java,*IT 命名,与单测 surefire 分离) -->
<plugin>
<groupId>org.apache.maven.plugins</groupId>
<artifactId>maven-failsafe-plugin</artifactId>
<version>3.5.6</version>
<configuration>
<skipITs>${skipITs}</skipITs>
<!-- 与引擎启动参数一致:cglib/反射框架在 JDK 17+ 需要 opens(否则 BeanMap 初始化失败);
55 个 IT 用例全量运行时需较大堆,默认 fork 堆不足会 OOM -->
<argLine>
-Xmx8g
--add-opens=java.base/java.lang=ALL-UNNAMED
--add-opens=java.base/java.lang.reflect=ALL-UNNAMED
--add-opens=java.base/java.util=ALL-UNNAMED
--add-opens=java.base/java.security=ALL-UNNAMED
--add-opens=java.base/java.io=ALL-UNNAMED
--add-opens=java.base/java.time=ALL-UNNAMED
--add-opens=java.base/jdk.internal.loader=ALL-UNNAMED
</argLine>
<reportsDirectory>
${project.build.directory}/failsafe-reports
</reportsDirectory>
</configuration>
<executions>
<execution>
<id>integration-test</id>
<goals>
<goal>integration-test</goal>
<goal>verify</goal>
</goals>
</execution>
</executions>
</plugin>
<!-- <plugin>-->
<!-- <groupId>org.codehaus.mojo</groupId>-->
<!-- <artifactId>exec-maven-plugin</artifactId>-->
<!-- <version>3.0.0</version>-->
Expand Down
Loading
Loading