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"
+  }
+}


Reply via email to