This is an automated email from the ASF dual-hosted git repository.

corgy-w 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 8ac4a022c4 [Feature][Connector-V2] Add Hudi timer flush (#11747)
8ac4a022c4 is described below

commit 8ac4a022c4e7c6820be4b576d4bc98084778dc73
Author: zhiwei.niu <[email protected]>
AuthorDate: Sun Aug 16 13:13:21 2026 +0800

    [Feature][Connector-V2] Add Hudi timer flush (#11747)
---
 docs/en/connectors/sink/Hudi.md                    |  20 +++-
 docs/zh/connectors/sink/Hudi.md                    |  19 +++-
 .../seatunnel/hudi/sink/writer/HudiSinkWriter.java |   9 ++
 .../hudi/sink/writer/HudiSinkWriterTest.java       | 122 +++++++++++++++++++++
 .../e2e/connector/hudi/HudiSinkCDCIT.java          | 114 +++++++++++++++++++
 .../hudi/mysql_cdc_to_hudi_timer_flush.conf        |  51 +++++++++
 6 files changed, 333 insertions(+), 2 deletions(-)

diff --git a/docs/en/connectors/sink/Hudi.md b/docs/en/connectors/sink/Hudi.md
index 85d89f05bc..617533c4be 100644
--- a/docs/en/connectors/sink/Hudi.md
+++ b/docs/en/connectors/sink/Hudi.md
@@ -110,7 +110,8 @@ Note: When this configuration corresponds to a single 
table, you can flatten the
 
 ### batch_interval_ms [Int]
 
-`batch_interval_ms` The maximum interval, in milliseconds, between two flushes 
to Hudi.
+`batch_interval_ms` is retained for compatibility. To schedule time-based 
flushes on Zeta, configure
+`sink.flush.interval` in the job `env` block.
 
 ### batch_size [Int]
 
@@ -159,6 +160,23 @@ Choose how to handle existing data before the 
synchronization task starts.
 
 Sink plugin common parameters, please refer to [Sink Common 
Options](../common-options/sink-common-options.md) for details.
 
+## Timer Flush
+
+Timer flush is an engine-level feature supported only by Zeta. Configure 
`sink.flush.interval` in the job `env` block
+to write pending Hudi records even when `batch_size` has not been reached. 
Spark and Flink do not inject `FlushSignal`
+records and therefore do not trigger this scheduled flush.
+
+```hocon
+env {
+  sink.flush.interval = 5000
+}
+```
+
+Hudi timer flush reuses the connector's synchronized batch flush and the Hudi 
client's auto-commit behavior. The Hudi
+sink does not provide a 2PC exactly-once writer, so timer flush provides 
at-least-once delivery. Retries can create
+additional commits. With `INSERT`, generated record keys can also produce 
duplicate rows after recovery; `UPSERT` with
+stable `record_key_fields` limits duplicate logical records.
+
 ## Examples
 
 ### Single Table Upsert
diff --git a/docs/zh/connectors/sink/Hudi.md b/docs/zh/connectors/sink/Hudi.md
index 1ab14baf08..7cd1224b73 100644
--- a/docs/zh/connectors/sink/Hudi.md
+++ b/docs/zh/connectors/sink/Hudi.md
@@ -108,7 +108,8 @@ SeaTunnel Hudi sink 会写入 Hudi 数据文件和 `.hoodie` 元数据,但不
 
 ### batch_interval_ms [Int]
 
-`batch_interval_ms` 两次刷新到 Hudi 的最大时间间隔,单位为毫秒。
+`batch_interval_ms` 为兼容性保留。在 Zeta 上需要定时刷新时,请在作业 `env` 中配置
+`sink.flush.interval`。
 
 ### batch_size [Int]
 
@@ -154,6 +155,22 @@ SeaTunnel Hudi sink 会写入 Hudi 数据文件和 `.hoodie` 元数据,但不
 
 Sink插件通用参数,请参考 [Sink Common Options](../common-options/sink-common-options.md) 
了解详细信息。
 
+## 定时刷新
+
+定时刷新是仅由 Zeta 支持的引擎级能力。在作业的 `env` 中配置 `sink.flush.interval` 后,即使尚未达到
+`batch_size`,Hudi Sink 也会写出待处理的记录。Spark 和 Flink 不会注入 `FlushSignal`,因此不会触发这种
+定时刷新。
+
+```hocon
+env {
+  sink.flush.interval = 5000
+}
+```
+
+Hudi 定时刷新复用连接器现有的同步批量刷新和 Hudi 客户端 auto-commit 行为。Hudi Sink 没有 2PC 精确一次
+写入器,因此定时刷新提供的是至少一次语义,重试可能产生额外的 commit。使用 `INSERT` 时,自动生成的
+record key 还可能在恢复后产生重复行;使用具有稳定 `record_key_fields` 的 `UPSERT` 可以减少逻辑记录重复。
+
 ## 示例
 
 ### 单表 UPSERT
diff --git 
a/seatunnel-connectors-v2/connector-hudi/src/main/java/org/apache/seatunnel/connectors/seatunnel/hudi/sink/writer/HudiSinkWriter.java
 
b/seatunnel-connectors-v2/connector-hudi/src/main/java/org/apache/seatunnel/connectors/seatunnel/hudi/sink/writer/HudiSinkWriter.java
index 130a79adab..3ab11cc360 100644
--- 
a/seatunnel-connectors-v2/connector-hudi/src/main/java/org/apache/seatunnel/connectors/seatunnel/hudi/sink/writer/HudiSinkWriter.java
+++ 
b/seatunnel-connectors-v2/connector-hudi/src/main/java/org/apache/seatunnel/connectors/seatunnel/hudi/sink/writer/HudiSinkWriter.java
@@ -67,6 +67,7 @@ public class HudiSinkWriter
                         sinkConfig, tableConfig.getTableName(), 
seaTunnelRowType);
         this.hudiRecordWriter =
                 new HudiRecordWriter(tableConfig, writeClientProvider, 
seaTunnelRowType);
+        context.registerFlushAction(this::timerFlush);
     }
 
     @Override
@@ -112,6 +113,14 @@ public class HudiSinkWriter
                 new HudiRecordWriter(tableConfig, writeClientProvider, 
seaTunnelRowType);
     }
 
+    /**
+     * Flushes buffered records when the sink receives a timer-generated 
FlushSignal. The signal is
+     * processed on the sink task thread in order with data records and 
checkpoint barriers.
+     */
+    private void timerFlush() {
+        hudiRecordWriter.flush();
+    }
+
     private void tryOpen() {
         if (!isOpen) {
             isOpen = true;
diff --git 
a/seatunnel-connectors-v2/connector-hudi/src/test/java/org/apache/seatunnel/connectors/seatunnel/hudi/sink/writer/HudiSinkWriterTest.java
 
b/seatunnel-connectors-v2/connector-hudi/src/test/java/org/apache/seatunnel/connectors/seatunnel/hudi/sink/writer/HudiSinkWriterTest.java
new file mode 100644
index 0000000000..21d3729d02
--- /dev/null
+++ 
b/seatunnel-connectors-v2/connector-hudi/src/test/java/org/apache/seatunnel/connectors/seatunnel/hudi/sink/writer/HudiSinkWriterTest.java
@@ -0,0 +1,122 @@
+/*
+ * 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.hudi.sink.writer;
+
+import org.apache.seatunnel.api.sink.MultiTableResourceManager;
+import org.apache.seatunnel.api.sink.SinkWriter;
+import org.apache.seatunnel.api.table.type.SeaTunnelRowType;
+import org.apache.seatunnel.common.utils.function.RunnableWithException;
+import org.apache.seatunnel.connectors.seatunnel.hudi.config.HudiSinkConfig;
+import org.apache.seatunnel.connectors.seatunnel.hudi.config.HudiTableConfig;
+import 
org.apache.seatunnel.connectors.seatunnel.hudi.exception.HudiConnectorException;
+import org.apache.seatunnel.connectors.seatunnel.hudi.exception.HudiErrorCode;
+import org.apache.seatunnel.connectors.seatunnel.hudi.sink.HudiClientManager;
+
+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 java.util.Optional;
+
+import static org.mockito.Mockito.doThrow;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.never;
+import static org.mockito.Mockito.times;
+import static org.mockito.Mockito.verify;
+import static org.mockito.Mockito.when;
+
+class HudiSinkWriterTest {
+
+    @Test
+    void shouldRegisterAndExecuteTimerFlush() throws Exception {
+        SinkWriter.Context context = mock(SinkWriter.Context.class);
+        ArgumentCaptor<RunnableWithException> actionCaptor =
+                ArgumentCaptor.forClass(RunnableWithException.class);
+
+        try (MockedConstruction<HudiRecordWriter> recordWriters = 
createWriter(context)) {
+            verify(context, 
times(1)).registerFlushAction(actionCaptor.capture());
+
+            actionCaptor.getValue().run();
+
+            verify(recordWriters.constructed().get(0), times(1)).flush();
+        }
+    }
+
+    @Test
+    void shouldPropagateTimerFlushFailure() throws Exception {
+        SinkWriter.Context context = mock(SinkWriter.Context.class);
+        ArgumentCaptor<RunnableWithException> actionCaptor =
+                ArgumentCaptor.forClass(RunnableWithException.class);
+
+        try (MockedConstruction<HudiRecordWriter> recordWriters = 
createWriter(context)) {
+            HudiConnectorException expected =
+                    new HudiConnectorException(
+                            HudiErrorCode.FLUSH_DATA_FAILED, "timer flush 
failed");
+            doThrow(expected).when(recordWriters.constructed().get(0)).flush();
+            verify(context).registerFlushAction(actionCaptor.capture());
+
+            HudiConnectorException actual =
+                    Assertions.assertThrows(
+                            HudiConnectorException.class, 
actionCaptor.getValue()::run);
+
+            Assertions.assertSame(expected, actual);
+        }
+    }
+
+    @Test
+    void shouldFlushCurrentRecordWriterAfterResourceManagerReplacement() 
throws Exception {
+        SinkWriter.Context context = mock(SinkWriter.Context.class);
+        ArgumentCaptor<RunnableWithException> actionCaptor =
+                ArgumentCaptor.forClass(RunnableWithException.class);
+        MultiTableResourceManager<HudiClientManager> resourceManager =
+                mock(MultiTableResourceManager.class);
+        when(resourceManager.getSharedResource())
+                .thenReturn(Optional.of(mock(HudiClientManager.class)));
+
+        try (MockedConstruction<HudiRecordWriter> recordWriters =
+                Mockito.mockConstruction(HudiRecordWriter.class)) {
+            HudiSinkWriter writer =
+                    new HudiSinkWriter(
+                            context,
+                            mock(SeaTunnelRowType.class),
+                            mock(HudiSinkConfig.class),
+                            mock(HudiTableConfig.class));
+            verify(context).registerFlushAction(actionCaptor.capture());
+
+            writer.setMultiTableResourceManager(resourceManager, 0);
+            actionCaptor.getValue().run();
+
+            Assertions.assertEquals(2, recordWriters.constructed().size());
+            verify(recordWriters.constructed().get(0), never()).flush();
+            verify(recordWriters.constructed().get(1)).flush();
+        }
+    }
+
+    private MockedConstruction<HudiRecordWriter> 
createWriter(SinkWriter.Context context) {
+        MockedConstruction<HudiRecordWriter> recordWriters =
+                Mockito.mockConstruction(HudiRecordWriter.class);
+        new HudiSinkWriter(
+                context,
+                mock(SeaTunnelRowType.class),
+                mock(HudiSinkConfig.class),
+                mock(HudiTableConfig.class));
+        return recordWriters;
+    }
+}
diff --git 
a/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-hudi-e2e/src/test/java/org/apache/seatunnel/e2e/connector/hudi/HudiSinkCDCIT.java
 
b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-hudi-e2e/src/test/java/org/apache/seatunnel/e2e/connector/hudi/HudiSinkCDCIT.java
index ab7c91a4b4..e5abaf33a1 100644
--- 
a/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-hudi-e2e/src/test/java/org/apache/seatunnel/e2e/connector/hudi/HudiSinkCDCIT.java
+++ 
b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-hudi-e2e/src/test/java/org/apache/seatunnel/e2e/connector/hudi/HudiSinkCDCIT.java
@@ -28,6 +28,7 @@ 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.apache.hadoop.conf.Configuration;
 import org.apache.hadoop.fs.FileUtil;
@@ -87,6 +88,8 @@ public class HudiSinkCDCIT extends TestSuiteBase implements 
TestResource {
 
     private static final String DATABASE = "st";
     private static final String TABLE_NAME = "st_test";
+    private static final String TIMER_FLUSH_DATABASE = "timer_flush_db";
+    private static final String TIMER_FLUSH_TABLE = "timer_flush_table";
     private static final String TABLE_PATH = HOST_VOLUME_MOUNT_PATH + "/hudi/";
     private static final String NAMESPACE = "hudi";
     private static final String NAMESPACE_TAR = "hudi.tar.gz";
@@ -173,6 +176,117 @@ public class HudiSinkCDCIT extends TestSuiteBase 
implements TestResource {
         upsertAndCheckData(container);
     }
 
+    @TestTemplate
+    public void testHudiTimerFlush(TestContainer container) throws Exception {
+        clearTable(MYSQL_DATABASE, SOURCE_TABLE);
+        FileUtil.fullyDelete(
+                new File(
+                        TABLE_PATH
+                                + File.separator
+                                + TIMER_FLUSH_DATABASE
+                                + File.separator
+                                + TIMER_FLUSH_TABLE));
+        String jobId = String.valueOf(JobIdGenerator.newJobId());
+        CompletableFuture<Container.ExecResult> jobFuture =
+                CompletableFuture.supplyAsync(
+                        () -> {
+                            try {
+                                return container.executeJob(
+                                        
"/hudi/mysql_cdc_to_hudi_timer_flush.conf", jobId);
+                            } catch (Exception e) {
+                                throw new RuntimeException(e);
+                            }
+                        });
+
+        try {
+            given().ignoreExceptions()
+                    .await()
+                    .atMost(2, TimeUnit.MINUTES)
+                    .pollInterval(2, TimeUnit.SECONDS)
+                    .untilAsserted(
+                            () -> {
+                                assertJobStillRunning(
+                                        jobFuture,
+                                        "The streaming job terminated before 
reaching RUNNING");
+                                Assertions.assertEquals("RUNNING", 
container.getJobStatus(jobId));
+                            });
+
+            executeSql(
+                    "INSERT INTO "
+                            + MYSQL_DATABASE
+                            + "."
+                            + SOURCE_TABLE
+                            + " (id, f_bigint, f_json) VALUES (1001, 1, 
JSON_OBJECT('phase', 1))");
+            awaitTimerFlush(jobFuture, 1);
+
+            executeSql(
+                    "INSERT INTO "
+                            + MYSQL_DATABASE
+                            + "."
+                            + SOURCE_TABLE
+                            + " (id, f_bigint, f_json) VALUES (1002, 2, 
JSON_OBJECT('phase', 2))");
+            awaitTimerFlush(jobFuture, 2);
+        } finally {
+            if (!jobFuture.isDone()) {
+                Container.ExecResult cancelResult = container.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 awaitTimerFlush(
+            CompletableFuture<Container.ExecResult> jobFuture, long 
expectedRows) {
+        given().ignoreExceptions()
+                .await()
+                .atMost(120, TimeUnit.SECONDS)
+                .pollInterval(2, TimeUnit.SECONDS)
+                .untilAsserted(
+                        () -> {
+                            assertJobStillRunning(
+                                    jobFuture,
+                                    "The streaming job terminated before timer 
flush published "
+                                            + expectedRows
+                                            + " rows");
+                            Assertions.assertEquals(
+                                    expectedRows,
+                                    countNewestCommitRows(),
+                                    "Timer flush should publish buffered rows 
while the job is running");
+                        });
+    }
+
+    private long countNewestCommitRows() throws IOException {
+        File tablePath =
+                new File(
+                        TABLE_PATH
+                                + File.separator
+                                + TIMER_FLUSH_DATABASE
+                                + File.separator
+                                + TIMER_FLUSH_TABLE);
+        Configuration configuration = new Configuration();
+        configuration.set("fs.defaultFS", LocalFileSystem.DEFAULT_FS);
+        long rowCount = 0;
+        try (ParquetReader<Group> reader =
+                ParquetReader.builder(new GroupReadSupport(), 
getNewestCommitFilePath(tablePath))
+                        .withConf(configuration)
+                        .build()) {
+            while (reader.read() != null) {
+                rowCount++;
+            }
+        }
+        return rowCount;
+    }
+
+    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 insertAndCheckData(TestContainer container) throws 
InterruptedException {
         // Init table data
         initSourceTableData(MYSQL_DATABASE, SOURCE_TABLE);
diff --git 
a/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-hudi-e2e/src/test/resources/hudi/mysql_cdc_to_hudi_timer_flush.conf
 
b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-hudi-e2e/src/test/resources/hudi/mysql_cdc_to_hudi_timer_flush.conf
new file mode 100644
index 0000000000..69fbf1b801
--- /dev/null
+++ 
b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-hudi-e2e/src/test/resources/hudi/mysql_cdc_to_hudi_timer_flush.conf
@@ -0,0 +1,51 @@
+#
+# 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 {
+    catalog {
+      factory = Mysql
+    }
+    database-names = ["mysql_cdc"]
+    table-names = ["mysql_cdc.mysql_cdc_e2e_source_table"]
+    format = DEFAULT
+    username = "st_user"
+    password = "seatunnel"
+    url = "jdbc:mysql://mysql_cdc_e2e:3306/mysql_cdc"
+  }
+}
+
+sink {
+  Hudi {
+    op_type = "UPSERT"
+    table_dfs_path = "/tmp/seatunnel_mnt/hudi"
+    database = "timer_flush_db"
+    table_name = "timer_flush_table"
+    table_type = "COPY_ON_WRITE"
+    record_key_fields = "id"
+    cdc_enabled = true
+    batch_size = 100
+    batch_interval_ms = 300000
+  }
+}

Reply via email to