This is an automated email from the ASF dual-hosted git repository.
dybyte pushed a commit to branch dev
in repository https://gitbox.apache.org/repos/asf/seatunnel.git
The following commit(s) were added to refs/heads/dev by this push:
new e55c843af2 [Feature][Connector-V2] Add ClickHouse timer flush (#11519)
e55c843af2 is described below
commit e55c843af27ec324f0af31840632b322bf091186
Author: zhiwei.niu <[email protected]>
AuthorDate: Thu Jul 23 01:43:46 2026 +0800
[Feature][Connector-V2] Add ClickHouse timer flush (#11519)
---
docs/en/connectors/sink/Clickhouse.md | 31 ++-
docs/zh/connectors/sink/Clickhouse.md | 31 ++-
.../sink/client/ClickhouseSinkWriter.java | 11 +-
.../sink/client/ClickhouseSinkWriterTest.java | 126 ++++++++++
.../connector-clickhouse-e2e/pom.xml | 20 ++
.../clickhouse/ClickhouseTimerFlushIT.java | 267 +++++++++++++++++++++
.../src/test/resources/ddl/shop.sql | 42 ++++
.../src/test/resources/docker/server-gtids/my.cnf | 31 +++
.../src/test/resources/docker/setup.sql | 21 ++
.../mysqlcdc_to_clickhouse_timer_flush.conf | 47 ++++
10 files changed, 621 insertions(+), 6 deletions(-)
diff --git a/docs/en/connectors/sink/Clickhouse.md
b/docs/en/connectors/sink/Clickhouse.md
index bfb5758182..440a16d0b8 100644
--- a/docs/en/connectors/sink/Clickhouse.md
+++ b/docs/en/connectors/sink/Clickhouse.md
@@ -18,7 +18,7 @@ import ChangeLog from '../changelog/connector-clickhouse.md';
> The Clickhouse sink can reduce duplicate effects through idempotent writing
> when the target table engine supports deduplication, such as
> `AggregatingMergeTree` or `ReplacingMergeTree`. It is not marked as
> exactly-once because the guarantee depends on the target table design.
- [x] [support multiple table
sink](../../introduction/concepts/connector-v2-features.md)
-- [ ] [timer flush](../../introduction/concepts/connector-v2-features.md)
+- [x] [timer flush](../../introduction/concepts/connector-v2-features.md)
## Description
@@ -131,6 +131,35 @@ The following placeholders can be used:
- `rowtype_unique_key`: Retrieves the unique key from the upstream schema
(this may be a list).
- `comment`: Retrieves the table comment from the upstream schema.
+### Zeta Timer Flush
+
+This engine-level capability is available only in Zeta. Configure
`sink.flush.interval` in `env` to periodically write buffered rows through
ClickHouse JDBC even when `bulk_size` has not been reached. Spark and Flink do
not trigger this scheduled flush.
+
+:::tip
+
+ClickHouse timer flush does not provide 2PC exactly-once semantics. The
ClickHouse Sink remains at-least-once, and retries or task restarts may insert
rows again. When duplicate handling is required, design the target table with a
suitable deduplication engine and deterministic keys.
+
+:::
+
+```hocon
+env {
+ job.mode = "STREAMING"
+ checkpoint.interval = 300000
+ sink.flush.interval = 5000
+}
+
+sink {
+ Clickhouse {
+ host = "localhost:8123"
+ database = "default"
+ table = "seatunnel_table"
+ username = "default"
+ password = ""
+ bulk_size = 10000
+ }
+}
+```
+
## Example Configurations and Cases
### How to Create a Clickhouse Data Synchronization Jobs
diff --git a/docs/zh/connectors/sink/Clickhouse.md
b/docs/zh/connectors/sink/Clickhouse.md
index 11c34aa1c9..fb97277384 100644
--- a/docs/zh/connectors/sink/Clickhouse.md
+++ b/docs/zh/connectors/sink/Clickhouse.md
@@ -17,7 +17,7 @@ import ChangeLog from '../changelog/connector-clickhouse.md';
> 当目标表引擎支持去重时,例如 `AggregatingMergeTree` 或 `ReplacingMergeTree`,Clickhouse Sink
> 可以通过幂等写入减少重复数据影响。这里未标记为精准一次,因为实际保证取决于目标表设计。
- [x] [支持多表写入](../../introduction/concepts/connector-v2-features.md)
-- [ ] [定时刷新](../../introduction/concepts/connector-v2-features.md)
+- [x] [定时刷新](../../introduction/concepts/connector-v2-features.md)
## 描述
@@ -132,6 +132,35 @@ CREATE TABLE IF NOT EXISTS `${database}`.`${table}` (
- rowtype_unique_key:用于获取上游模式中的唯一键(可能是列表)。
- comment:用于获取上游模式中的表注释。
+### Zeta 定时刷新
+
+该引擎级能力仅由 Zeta 支持。可以在 `env` 块中配置 `sink.flush.interval`,使尚未达到 `bulk_size`
的缓冲数据也能定时通过 ClickHouse JDBC 写出。Spark 和 Flink 不会触发该定时刷新。
+
+:::tip
+
+ClickHouse 定时刷新不提供基于 2PC 的精准一次语义,ClickHouse Sink
仍为至少一次语义,失败重试或任务重启后可能重复插入数据。如果业务需要处理重复数据,应为目标表选择合适的去重引擎并使用确定性的键。
+
+:::
+
+```hocon
+env {
+ job.mode = "STREAMING"
+ checkpoint.interval = 300000
+ sink.flush.interval = 5000
+}
+
+sink {
+ Clickhouse {
+ host = "localhost:8123"
+ database = "default"
+ table = "seatunnel_table"
+ username = "default"
+ password = ""
+ bulk_size = 10000
+ }
+}
+```
+
## 示例配置与案例
### 如何创建一个clickhouse 同步任务
diff --git
a/seatunnel-connectors-v2/connector-clickhouse/src/main/java/org/apache/seatunnel/connectors/seatunnel/clickhouse/sink/client/ClickhouseSinkWriter.java
b/seatunnel-connectors-v2/connector-clickhouse/src/main/java/org/apache/seatunnel/connectors/seatunnel/clickhouse/sink/client/ClickhouseSinkWriter.java
index ad83782832..fbbd78be65 100644
---
a/seatunnel-connectors-v2/connector-clickhouse/src/main/java/org/apache/seatunnel/connectors/seatunnel/clickhouse/sink/client/ClickhouseSinkWriter.java
+++
b/seatunnel-connectors-v2/connector-clickhouse/src/main/java/org/apache/seatunnel/connectors/seatunnel/clickhouse/sink/client/ClickhouseSinkWriter.java
@@ -52,7 +52,6 @@ public class ClickhouseSinkWriter
implements SinkWriter<SeaTunnelRow, CKCommitInfo, ClickhouseSinkState>,
SupportMultiTableSinkWriter<Void> {
- private final Context context;
private final ReaderOption option;
private final ShardRouter shardRouter;
private final transient ClickhouseProxy proxy;
@@ -60,11 +59,10 @@ public class ClickhouseSinkWriter
ClickhouseSinkWriter(ReaderOption option, Context context) {
this.option = option;
- this.context = context;
-
this.proxy = new
ClickhouseProxy(option.getShardMetadata().getDefaultShard().getNode());
this.shardRouter = new ShardRouter(proxy, option.getShardMetadata());
this.statementMap = initStatementMap();
+ context.registerFlushAction(this::flushPendingStatements);
}
@Override
@@ -93,6 +91,12 @@ public class ClickhouseSinkWriter
@Override
public Optional<CKCommitInfo> prepareCommit() throws IOException {
+ flushPendingStatements();
+ return Optional.empty();
+ }
+
+ /** Flushes pending JDBC batches when the Zeta engine delivers a timer
flush signal. */
+ private void flushPendingStatements() {
for (ClickhouseBatchStatement batchStatement : statementMap.values()) {
JdbcBatchStatementExecutor statement =
batchStatement.getJdbcBatchStatementExecutor();
IntHolder intHolder = batchStatement.getIntHolder();
@@ -101,7 +105,6 @@ public class ClickhouseSinkWriter
intHolder.setValue(0);
}
}
- return Optional.empty();
}
@Override
diff --git
a/seatunnel-connectors-v2/connector-clickhouse/src/test/java/org/apache/seatunnel/connectors/seatunnel/clickhouse/sink/client/ClickhouseSinkWriterTest.java
b/seatunnel-connectors-v2/connector-clickhouse/src/test/java/org/apache/seatunnel/connectors/seatunnel/clickhouse/sink/client/ClickhouseSinkWriterTest.java
new file mode 100644
index 0000000000..2a74c14c92
--- /dev/null
+++
b/seatunnel-connectors-v2/connector-clickhouse/src/test/java/org/apache/seatunnel/connectors/seatunnel/clickhouse/sink/client/ClickhouseSinkWriterTest.java
@@ -0,0 +1,126 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.seatunnel.connectors.seatunnel.clickhouse.sink.client;
+
+import org.apache.seatunnel.api.sink.SinkWriter;
+import org.apache.seatunnel.api.table.type.BasicType;
+import org.apache.seatunnel.api.table.type.SeaTunnelDataType;
+import org.apache.seatunnel.api.table.type.SeaTunnelRow;
+import org.apache.seatunnel.api.table.type.SeaTunnelRowType;
+import org.apache.seatunnel.common.utils.function.RunnableWithException;
+import
org.apache.seatunnel.connectors.seatunnel.clickhouse.config.ReaderOption;
+import
org.apache.seatunnel.connectors.seatunnel.clickhouse.exception.ClickhouseConnectorException;
+import org.apache.seatunnel.connectors.seatunnel.clickhouse.shard.Shard;
+import
org.apache.seatunnel.connectors.seatunnel.clickhouse.shard.ShardMetadata;
+import
org.apache.seatunnel.connectors.seatunnel.clickhouse.util.ClickhouseProxy;
+
+import org.junit.jupiter.api.Assertions;
+import org.junit.jupiter.api.Test;
+import org.mockito.ArgumentCaptor;
+import org.mockito.MockedConstruction;
+import org.mockito.Mockito;
+
+import com.clickhouse.jdbc.internal.ClickHouseConnectionImpl;
+
+import java.sql.PreparedStatement;
+import java.sql.SQLException;
+import java.util.Collections;
+import java.util.Properties;
+
+import static org.mockito.Mockito.doThrow;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.times;
+import static org.mockito.Mockito.verify;
+import static org.mockito.Mockito.when;
+
+public class ClickhouseSinkWriterTest {
+
+ @Test
+ void shouldRegisterTimerFlushAction() throws Exception {
+ SinkWriter.Context context = mock(SinkWriter.Context.class);
+ PreparedStatement statement = mock(PreparedStatement.class);
+ ArgumentCaptor<RunnableWithException> actionCaptor =
+ ArgumentCaptor.forClass(RunnableWithException.class);
+
+ createWriterWithPendingRow(context, statement);
+
+ verify(context, times(1)).registerFlushAction(actionCaptor.capture());
+ actionCaptor.getValue().run();
+ verify(statement, times(1)).executeBatch();
+ }
+
+ @Test
+ void shouldPropagateTimerFlushFailure() throws Exception {
+ SinkWriter.Context context = mock(SinkWriter.Context.class);
+ PreparedStatement statement = mock(PreparedStatement.class);
+ ArgumentCaptor<RunnableWithException> actionCaptor =
+ ArgumentCaptor.forClass(RunnableWithException.class);
+ SQLException expected = new SQLException("timer flush failed");
+ doThrow(expected).when(statement).executeBatch();
+
+ createWriterWithPendingRow(context, statement);
+
+ verify(context).registerFlushAction(actionCaptor.capture());
+ ClickhouseConnectorException actual =
+ Assertions.assertThrows(
+ ClickhouseConnectorException.class,
actionCaptor.getValue()::run);
+ Assertions.assertSame(expected, actual.getCause());
+ }
+
+ private void createWriterWithPendingRow(SinkWriter.Context context,
PreparedStatement statement)
+ throws Exception {
+ Shard shard = mock(Shard.class);
+
when(shard.getJdbcUrl()).thenReturn("jdbc:clickhouse://localhost:8123/default");
+ ShardMetadata shardMetadata =
+ new ShardMetadata(
+ null,
+ null,
+ null,
+ "default",
+ "timer_flush",
+ "MergeTree",
+ false,
+ shard,
+ "default",
+ "");
+ SeaTunnelRowType rowType =
+ new SeaTunnelRowType(
+ new String[] {"id"}, new SeaTunnelDataType[]
{BasicType.INT_TYPE});
+ ReaderOption option =
+ ReaderOption.builder()
+ .shardMetadata(shardMetadata)
+ .properties(new Properties())
+ .seaTunnelRowType(rowType)
+ .tableEngine("MergeTree")
+ .tableSchema(Collections.singletonMap("id", "Int32"))
+ .bulkSize(10000)
+ .build();
+
+ try (MockedConstruction<ClickhouseProxy> ignored =
+ Mockito.mockConstruction(ClickhouseProxy.class);
+ MockedConstruction<ClickHouseConnectionImpl> ignoredConnection
=
+ Mockito.mockConstruction(
+ ClickHouseConnectionImpl.class,
+ (mock, constructionContext) ->
+
when(mock.prepareStatement(Mockito.anyString()))
+ .thenReturn(statement))) {
+ ClickhouseSinkWriter writer = new ClickhouseSinkWriter(option,
context);
+ writer.write(new SeaTunnelRow(new Object[] {1}));
+ }
+ }
+}
diff --git
a/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-clickhouse-e2e/pom.xml
b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-clickhouse-e2e/pom.xml
index 77f1592736..f69996ffe8 100644
--- a/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-clickhouse-e2e/pom.xml
+++ b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-clickhouse-e2e/pom.xml
@@ -27,6 +27,7 @@
<properties>
<clickhouse.jdbc.version>0.3.2-patch11</clickhouse.jdbc.version>
+ <mysql.connector.version>8.0.32</mysql.connector.version>
</properties>
<dependencies>
@@ -57,5 +58,24 @@
<version>${project.version}</version>
<scope>test</scope>
</dependency>
+ <dependency>
+ <groupId>com.mysql</groupId>
+ <artifactId>mysql-connector-j</artifactId>
+ <version>${mysql.connector.version}</version>
+ <scope>test</scope>
+ </dependency>
+ <dependency>
+ <groupId>org.apache.seatunnel</groupId>
+ <artifactId>connector-cdc-mysql</artifactId>
+ <version>${project.version}</version>
+ <type>test-jar</type>
+ <scope>test</scope>
+ </dependency>
+ <dependency>
+ <groupId>org.testcontainers</groupId>
+ <artifactId>mysql</artifactId>
+ <version>${testcontainer.version}</version>
+ <scope>test</scope>
+ </dependency>
</dependencies>
</project>
diff --git
a/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-clickhouse-e2e/src/test/java/org/apache/seatunnel/connectors/seatunnel/clickhouse/ClickhouseTimerFlushIT.java
b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-clickhouse-e2e/src/test/java/org/apache/seatunnel/connectors/seatunnel/clickhouse/ClickhouseTimerFlushIT.java
new file mode 100644
index 0000000000..734b4a2a2b
--- /dev/null
+++
b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-clickhouse-e2e/src/test/java/org/apache/seatunnel/connectors/seatunnel/clickhouse/ClickhouseTimerFlushIT.java
@@ -0,0 +1,267 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.seatunnel.connectors.seatunnel.clickhouse;
+
+import
org.apache.seatunnel.connectors.seatunnel.cdc.mysql.testutils.MySqlContainer;
+import
org.apache.seatunnel.connectors.seatunnel.cdc.mysql.testutils.MySqlVersion;
+import
org.apache.seatunnel.connectors.seatunnel.cdc.mysql.testutils.UniqueDatabase;
+import org.apache.seatunnel.e2e.common.TestResource;
+import org.apache.seatunnel.e2e.common.TestSuiteBase;
+import org.apache.seatunnel.e2e.common.container.ContainerExtendedFactory;
+import org.apache.seatunnel.e2e.common.container.EngineType;
+import org.apache.seatunnel.e2e.common.container.TestContainer;
+import org.apache.seatunnel.e2e.common.junit.DisabledOnContainer;
+import org.apache.seatunnel.e2e.common.junit.TestContainerExtension;
+import org.apache.seatunnel.e2e.common.util.JobIdGenerator;
+
+import org.junit.jupiter.api.AfterAll;
+import org.junit.jupiter.api.Assertions;
+import org.junit.jupiter.api.BeforeAll;
+import org.junit.jupiter.api.TestTemplate;
+import org.testcontainers.containers.ClickHouseContainer;
+import org.testcontainers.containers.Container;
+import org.testcontainers.containers.GenericContainer;
+import org.testcontainers.containers.output.Slf4jLogConsumer;
+import org.testcontainers.lifecycle.Startables;
+import org.testcontainers.utility.DockerLoggerFactory;
+import org.testcontainers.utility.MountableFile;
+
+import com.mysql.cj.jdbc.Driver;
+import lombok.extern.slf4j.Slf4j;
+
+import java.io.IOException;
+import java.nio.file.Files;
+import java.nio.file.Path;
+import java.nio.file.Paths;
+import java.sql.Connection;
+import java.sql.DriverManager;
+import java.sql.ResultSet;
+import java.sql.Statement;
+import java.util.Properties;
+import java.util.concurrent.CompletableFuture;
+import java.util.concurrent.TimeUnit;
+import java.util.stream.Stream;
+
+import static org.awaitility.Awaitility.await;
+
+@Slf4j
+@DisabledOnContainer(
+ value = {},
+ type = {EngineType.SPARK, EngineType.FLINK},
+ disabledReason =
+ "engine-level timer flush (sink.flush.interval) is only
supported on Zeta engine")
+public class ClickhouseTimerFlushIT extends TestSuiteBase implements
TestResource {
+
+ private static final String CLICKHOUSE_IMAGE =
"clickhouse/clickhouse-server:23.3.13.6";
+ private static final String CLICKHOUSE_HOST = "clickhouse-timer-flush-e2e";
+ private static final String CLICKHOUSE_DATABASE = "default";
+ private static final String CLICKHOUSE_TABLE = "clickhouse_timer_flush";
+ private static final String CLICKHOUSE_DRIVER =
"com.clickhouse.jdbc.ClickHouseDriver";
+ private static final String MYSQL_HOST =
"mysql_clickhouse_timer_flush_e2e";
+ private static final String MYSQL_USER_NAME = "mysqluser";
+ private static final String MYSQL_USER_PASSWORD = "mysqlpw";
+ private static final String MYSQL_DATABASE = "shop";
+ private static final String MYSQL_CDC_PLUGIN_LIB =
"/tmp/seatunnel/plugins/MySQL-CDC/lib";
+ private static final MySqlContainer MYSQL_CONTAINER =
createMySqlContainer(MySqlVersion.V8_0);
+
+ private final UniqueDatabase shopDatabase = new
UniqueDatabase(MYSQL_CONTAINER, MYSQL_DATABASE);
+ private ClickHouseContainer clickhouseContainer;
+ private Connection clickhouseConnection;
+
+ @TestContainerExtension
+ private final ContainerExtendedFactory extendedFactory =
this::copyMySQLDriverToContainer;
+
+ @BeforeAll
+ @Override
+ public void startUp() throws Exception {
+ clickhouseContainer =
+ new ClickHouseContainer(CLICKHOUSE_IMAGE)
+ .withNetwork(NETWORK)
+ .withNetworkAliases(CLICKHOUSE_HOST)
+ .withCreateContainerCmdModifier(
+ command ->
command.withHostName(CLICKHOUSE_HOST))
+ .withLogConsumer(
+ new Slf4jLogConsumer(
+
DockerLoggerFactory.getLogger(CLICKHOUSE_IMAGE)));
+ Startables.deepStart(Stream.of(clickhouseContainer,
MYSQL_CONTAINER)).join();
+
+ await().ignoreExceptions()
+ .atMost(360, TimeUnit.SECONDS)
+ .untilAsserted(this::initializeClickhouseConnection);
+ shopDatabase.createAndInitialize();
+ initializeTimerFlushTable();
+ }
+
+ @TestTemplate
+ public void testClickhouseTimerFlush(TestContainer testContainer) throws
Exception {
+ String jobId = String.valueOf(JobIdGenerator.newJobId());
+ CompletableFuture<Container.ExecResult> jobFuture =
+ CompletableFuture.supplyAsync(
+ () -> {
+ try {
+ return testContainer.executeJob(
+
"/mysqlcdc_to_clickhouse_timer_flush.conf", jobId);
+ } catch (Exception e) {
+ throw new RuntimeException(e);
+ }
+ });
+
+ try {
+ await().atMost(2, TimeUnit.MINUTES)
+ .pollInterval(2, TimeUnit.SECONDS)
+ .untilAsserted(
+ () -> {
+ assertJobStillRunning(
+ jobFuture,
+ "The streaming job terminated before
reaching RUNNING");
+ Assertions.assertEquals(
+ "RUNNING",
testContainer.getJobStatus(jobId));
+ });
+
+ await().atMost(120, TimeUnit.SECONDS)
+ .ignoreExceptions()
+ .pollInterval(2, TimeUnit.SECONDS)
+ .untilAsserted(
+ () -> {
+ assertJobStillRunning(
+ jobFuture,
+ "The streaming job terminated before
timer flush published the snapshot");
+ Assertions.assertEquals(9, tableCount());
+ });
+
+ try (Connection connection = shopDatabase.getJdbcConnection();
+ Statement statement = connection.createStatement()) {
+ statement.executeUpdate(
+ "INSERT INTO products (id, name, description, weight) "
+ + "VALUES (110, 'timer-flush', 'timer-flush
probe', 1.0)");
+ }
+
+ await().atMost(120, TimeUnit.SECONDS)
+ .ignoreExceptions()
+ .pollInterval(2, TimeUnit.SECONDS)
+ .untilAsserted(
+ () -> {
+ assertJobStillRunning(
+ jobFuture,
+ "The streaming job terminated before
timer flush published the binlog event");
+ Assertions.assertEquals(10, tableCount());
+ });
+ } finally {
+ if (!jobFuture.isDone()) {
+ Container.ExecResult cancelResult =
testContainer.cancelJob(jobId);
+ Assertions.assertEquals(0, cancelResult.getExitCode(),
cancelResult.getStderr());
+ }
+ }
+
+ Container.ExecResult jobResult = jobFuture.get(120, TimeUnit.SECONDS);
+ Assertions.assertEquals(0, jobResult.getExitCode(),
jobResult.getStderr());
+ }
+
+ private void assertJobStillRunning(
+ CompletableFuture<Container.ExecResult> jobFuture, String message)
throws Exception {
+ if (jobFuture.isDone()) {
+ Container.ExecResult jobResult = jobFuture.get();
+ Assertions.fail(message + ":\n" + jobResult.getStderr());
+ }
+ }
+
+ private void initializeClickhouseConnection() throws Exception {
+ Properties properties = new Properties();
+ properties.put("user", clickhouseContainer.getUsername());
+ properties.put("password", clickhouseContainer.getPassword());
+ Class.forName(CLICKHOUSE_DRIVER);
+ clickhouseConnection =
+ DriverManager.getConnection(clickhouseContainer.getJdbcUrl(),
properties);
+ Assertions.assertNotNull(clickhouseConnection);
+ }
+
+ private void initializeTimerFlushTable() throws Exception {
+ String table = CLICKHOUSE_DATABASE + "." + CLICKHOUSE_TABLE;
+ try (Statement statement = clickhouseConnection.createStatement()) {
+ statement.execute(
+ "CREATE TABLE IF NOT EXISTS "
+ + table
+ + " (id Int32, name String, description
Nullable(String), "
+ + "weight Nullable(Float32)) ENGINE=MergeTree
ORDER BY id");
+ statement.execute("TRUNCATE TABLE " + table);
+ }
+ }
+
+ private int tableCount() throws Exception {
+ try (Statement statement = clickhouseConnection.createStatement();
+ ResultSet resultSet =
+ statement.executeQuery(
+ "SELECT COUNT(*) FROM "
+ + CLICKHOUSE_DATABASE
+ + "."
+ + CLICKHOUSE_TABLE)) {
+ Assertions.assertTrue(resultSet.next());
+ return resultSet.getInt(1);
+ }
+ }
+
+ private void copyMySQLDriverToContainer(GenericContainer<?> container)
+ throws IOException, InterruptedException {
+ Path driverJarPath = mysqlDriverJarPath();
+ Assertions.assertTrue(
+ Files.isRegularFile(driverJarPath),
+ "MySQL JDBC driver should be resolved from the test classpath
before E2E runs: "
+ + driverJarPath);
+ Container.ExecResult extraCommands =
+ container.execInContainer("bash", "-c", "mkdir -p " +
MYSQL_CDC_PLUGIN_LIB);
+ Assertions.assertEquals(0, extraCommands.getExitCode(),
extraCommands.getStderr());
+ container.copyFileToContainer(
+ MountableFile.forHostPath(driverJarPath),
+ MYSQL_CDC_PLUGIN_LIB + "/" + driverJarPath.getFileName());
+ }
+
+ private Path mysqlDriverJarPath() {
+ try {
+ return Paths.get(
+
Driver.class.getProtectionDomain().getCodeSource().getLocation().toURI());
+ } catch (Exception e) {
+ throw new RuntimeException(
+ "Failed to resolve MySQL JDBC driver jar from the test
classpath", e);
+ }
+ }
+
+ private static MySqlContainer createMySqlContainer(MySqlVersion version) {
+ return new MySqlContainer(version)
+ .withConfigurationOverride("docker/server-gtids/my.cnf")
+ .withSetupSQL("docker/setup.sql")
+ .withNetwork(NETWORK)
+ .withNetworkAliases(MYSQL_HOST)
+ .withDatabaseName(MYSQL_DATABASE)
+ .withUsername(MYSQL_USER_NAME)
+ .withPassword(MYSQL_USER_PASSWORD)
+ .withLogConsumer(
+ new
Slf4jLogConsumer(DockerLoggerFactory.getLogger("mysql-docker-image")));
+ }
+
+ @AfterAll
+ @Override
+ public void tearDown() throws Exception {
+ if (clickhouseConnection != null) {
+ clickhouseConnection.close();
+ }
+ if (clickhouseContainer != null) {
+ clickhouseContainer.close();
+ }
+ MYSQL_CONTAINER.close();
+ }
+}
diff --git
a/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-clickhouse-e2e/src/test/resources/ddl/shop.sql
b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-clickhouse-e2e/src/test/resources/ddl/shop.sql
new file mode 100644
index 0000000000..8451311108
--- /dev/null
+++
b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-clickhouse-e2e/src/test/resources/ddl/shop.sql
@@ -0,0 +1,42 @@
+--
+-- Licensed to the Apache Software Foundation (ASF) under one or more
+-- contributor license agreements. See the NOTICE file distributed with
+-- this work for additional information regarding copyright ownership.
+-- The ASF licenses this file to You under the Apache License, Version 2.0
+-- (the "License"); you may not use this file except in compliance with
+-- the License. You may obtain a copy of the License at
+--
+-- http://www.apache.org/licenses/LICENSE-2.0
+--
+-- Unless required by applicable law or agreed to in writing, software
+-- distributed under the License is distributed on an "AS IS" BASIS,
+-- WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+-- See the License for the specific language governing permissions and
+-- limitations under the License.
+--
+
+--
----------------------------------------------------------------------------------------------------------------
+-- DATABASE: shop
+--
----------------------------------------------------------------------------------------------------------------
+CREATE DATABASE IF NOT EXISTS `shop`;
+use shop;
+
+drop table if exists products;
+-- Create and populate our products using a single insert with many rows
+CREATE TABLE products (
+ id INTEGER NOT NULL AUTO_INCREMENT PRIMARY KEY,
+ name VARCHAR(255) NOT NULL DEFAULT 'SeaTunnel',
+ description VARCHAR(512),
+ weight FLOAT
+);
+
+INSERT INTO products
+VALUES (101,"scooter","Small 2-wheel scooter",3.14),
+ (102,"car battery","12V car battery",8.1),
+ (103,"12-pack drill bits","12-pack of drill bits with sizes ranging
from #40 to #3",0.8),
+ (104,"hammer","12oz carpenter's hammer",0.75),
+ (105,"hammer","14oz carpenter's hammer",0.875),
+ (106,"hammer","16oz carpenter's hammer",1.0),
+ (107,"rocks","box of assorted rocks",5.3),
+ (108,"jacket","water resistent black wind breaker",0.1),
+ (109,"spare tire","24 inch spare tire",22.2);
diff --git
a/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-clickhouse-e2e/src/test/resources/docker/server-gtids/my.cnf
b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-clickhouse-e2e/src/test/resources/docker/server-gtids/my.cnf
new file mode 100644
index 0000000000..bfe5641b9c
--- /dev/null
+++
b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-clickhouse-e2e/src/test/resources/docker/server-gtids/my.cnf
@@ -0,0 +1,31 @@
+#
+# Licensed to the Apache Software Foundation (ASF) under one or more
+# contributor license agreements. See the NOTICE file distributed with
+# this work for additional information regarding copyright ownership.
+# The ASF licenses this file to You under the Apache License, Version 2.0
+# (the "License"); you may not use this file except in compliance with
+# the License. You may obtain a copy of the License at
+#
+# http://www.apache.org/licenses/LICENSE-2.0
+#
+# Unless required by applicable law or agreed to in writing, software
+# distributed under the License is distributed on an "AS IS" BASIS,
+# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+# See the License for the specific language governing permissions and
+# limitations under the License.
+#
+
+[mysqld]
+skip-host-cache
+skip-name-resolve
+secure-file-priv=/var/lib/mysql
+user=mysql
+symbolic-links=0
+
+server-id = 223344
+log_bin = mysql-bin
+expire_logs_days = 1
+binlog_format = row
+
+gtid_mode = on
+enforce_gtid_consistency = on
diff --git
a/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-clickhouse-e2e/src/test/resources/docker/setup.sql
b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-clickhouse-e2e/src/test/resources/docker/setup.sql
new file mode 100644
index 0000000000..95b8e560b9
--- /dev/null
+++
b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-clickhouse-e2e/src/test/resources/docker/setup.sql
@@ -0,0 +1,21 @@
+--
+-- Licensed to the Apache Software Foundation (ASF) under one or more
+-- contributor license agreements. See the NOTICE file distributed with
+-- this work for additional information regarding copyright ownership.
+-- The ASF licenses this file to You under the Apache License, Version 2.0
+-- (the "License"); you may not use this file except in compliance with
+-- the License. You may obtain a copy of the License at
+--
+-- http://www.apache.org/licenses/LICENSE-2.0
+--
+-- Unless required by applicable law or agreed to in writing, software
+-- distributed under the License is distributed on an "AS IS" BASIS,
+-- WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+-- See the License for the specific language governing permissions and
+-- limitations under the License.
+--
+
+GRANT ALL PRIVILEGES ON *.* TO 'mysqluser'@'%';
+
+CREATE USER 'st_user_source' IDENTIFIED BY 'mysqlpw';
+GRANT SELECT, RELOAD, SHOW DATABASES, REPLICATION SLAVE, REPLICATION CLIENT,
DROP, LOCK TABLES ON *.* TO 'st_user_source'@'%';
diff --git
a/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-clickhouse-e2e/src/test/resources/mysqlcdc_to_clickhouse_timer_flush.conf
b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-clickhouse-e2e/src/test/resources/mysqlcdc_to_clickhouse_timer_flush.conf
new file mode 100644
index 0000000000..6d4cb03ee1
--- /dev/null
+++
b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-clickhouse-e2e/src/test/resources/mysqlcdc_to_clickhouse_timer_flush.conf
@@ -0,0 +1,47 @@
+#
+# Licensed to the Apache Software Foundation (ASF) under one or more
+# contributor license agreements. See the NOTICE file distributed with
+# this work for additional information regarding copyright ownership.
+# The ASF licenses this file to You under the Apache License, Version 2.0
+# (the "License"); you may not use this file except in compliance with
+# the License. You may obtain a copy of the License at
+#
+# http://www.apache.org/licenses/LICENSE-2.0
+#
+# Unless required by applicable law or agreed to in writing, software
+# distributed under the License is distributed on an "AS IS" BASIS,
+# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+# See the License for the specific language governing permissions and
+# limitations under the License.
+#
+
+env {
+ parallelism = 1
+ job.mode = "STREAMING"
+ checkpoint.interval = 300000
+ sink.flush.interval = 500
+}
+
+source {
+ MySQL-CDC {
+ parallelism = 1
+ server-id = 5690
+ username = "st_user_source"
+ password = "mysqlpw"
+ table-names = ["shop.products"]
+ url = "jdbc:mysql://mysql_clickhouse_timer_flush_e2e:3306/shop"
+ }
+}
+
+sink {
+ Clickhouse {
+ host = "clickhouse-timer-flush-e2e:8123"
+ database = "default"
+ table = "clickhouse_timer_flush"
+ username = "default"
+ password = ""
+ bulk_size = 100000
+ schema_save_mode = "IGNORE"
+ data_save_mode = "APPEND_DATA"
+ }
+}