Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
42 changes: 42 additions & 0 deletions .github/workflows/connector-it.yml
Original file line number Diff line number Diff line change
@@ -0,0 +1,42 @@
# 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:
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*7;
Assertions.assertEquals(70, result);
}
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,4 @@
{
"uri": "",
"database": ""
}
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