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