From 755be62a76792b2b3be3e74c97a0f8a2cb42748a Mon Sep 17 00:00:00 2001 From: HuangXiao Date: Mon, 13 Jul 2026 20:43:37 +0800 Subject: [PATCH 1/2] [tiering] Use exponential delay restart strategy for Flink tiering job Configure the stateless tiering job to retry indefinitely when checkpointing is disabled, preventing transient failures from terminating the job permanently. Move the failover integration test to fluss-flink-tiering so that it covers the production job configuration. --- .../flink/tiering/FlinkTieringTestBase.java | 1 - fluss-flink/fluss-flink-tiering/pom.xml | 68 ++++++++++++++++++- .../fluss/flink/tiering/FlussLakeTiering.java | 7 ++ .../flink/tiering/TieringFailoverITCase.java | 30 +++++++- 4 files changed, 101 insertions(+), 5 deletions(-) rename fluss-flink/{fluss-flink-common => fluss-flink-tiering}/src/test/java/org/apache/fluss/flink/tiering/TieringFailoverITCase.java (79%) diff --git a/fluss-flink/fluss-flink-common/src/test/java/org/apache/fluss/flink/tiering/FlinkTieringTestBase.java b/fluss-flink/fluss-flink-common/src/test/java/org/apache/fluss/flink/tiering/FlinkTieringTestBase.java index 7e3ee48a78..aab355718e 100644 --- a/fluss-flink/fluss-flink-common/src/test/java/org/apache/fluss/flink/tiering/FlinkTieringTestBase.java +++ b/fluss-flink/fluss-flink-common/src/test/java/org/apache/fluss/flink/tiering/FlinkTieringTestBase.java @@ -103,7 +103,6 @@ void beforeEach() { execEnv = StreamExecutionEnvironment.getExecutionEnvironment(); execEnv.setRuntimeMode(RuntimeExecutionMode.STREAMING); execEnv.setParallelism(2); - execEnv.enableCheckpointing(500); } protected long createTable(TablePath tablePath, Schema schema) throws Exception { diff --git a/fluss-flink/fluss-flink-tiering/pom.xml b/fluss-flink/fluss-flink-tiering/pom.xml index 45917b8861..6dc98921dd 100644 --- a/fluss-flink/fluss-flink-tiering/pom.xml +++ b/fluss-flink/fluss-flink-tiering/pom.xml @@ -57,6 +57,72 @@ ${flink.minor.version} provided + + + + org.apache.fluss + fluss-flink-common + ${project.version} + test + test-jar + + + + org.apache.fluss + fluss-server + ${project.version} + test + + + + org.apache.fluss + fluss-server + ${project.version} + test + test-jar + + + + org.apache.fluss + fluss-common + ${project.version} + test + test-jar + + + + org.apache.fluss + fluss-test-utils + + + + + org.apache.curator + curator-test + ${curator.version} + test + + + + org.apache.flink + flink-clients + ${flink.minor.version} + test + + + + org.apache.flink + flink-connector-base + ${flink.minor.version} + test + + + + org.apache.flink + flink-table-common + ${flink.minor.version} + test + @@ -85,4 +151,4 @@ - \ No newline at end of file + diff --git a/fluss-flink/fluss-flink-tiering/src/main/java/org/apache/fluss/flink/tiering/FlussLakeTiering.java b/fluss-flink/fluss-flink-tiering/src/main/java/org/apache/fluss/flink/tiering/FlussLakeTiering.java index 340be397c6..aec35cf53f 100644 --- a/fluss-flink/fluss-flink-tiering/src/main/java/org/apache/fluss/flink/tiering/FlussLakeTiering.java +++ b/fluss-flink/fluss-flink-tiering/src/main/java/org/apache/fluss/flink/tiering/FlussLakeTiering.java @@ -23,12 +23,14 @@ import org.apache.fluss.flink.adapter.MultipleParameterToolAdapter; import org.apache.flink.configuration.JobManagerOptions; +import org.apache.flink.configuration.RestartStrategyOptions; import org.apache.flink.core.execution.JobClient; import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment; import java.util.Map; import java.util.ServiceLoader; +import static org.apache.flink.configuration.RestartStrategyOptions.RestartStrategyType.EXPONENTIAL_DELAY; import static org.apache.flink.runtime.executiongraph.failover.FailoverStrategyFactoryLoader.FULL_RESTART_STRATEGY_NAME; import static org.apache.fluss.flink.tiering.source.TieringSourceOptions.DATA_LAKE_CONFIG_PREFIX; import static org.apache.fluss.utils.PropertiesUtils.extractAndRemovePrefix; @@ -96,6 +98,11 @@ public FlussLakeTiering(String[] args) { org.apache.flink.configuration.Configuration flinkConfig = new org.apache.flink.configuration.Configuration(); flinkConfig.set(JobManagerOptions.EXECUTION_FAILOVER_STRATEGY, FULL_RESTART_STRATEGY_NAME); + // Configure restart strategy: the tiering job is completely stateless (offsets are + // persisted in Fluss Server's LakeSnapshot, not in Flink checkpoint state), so it is + // safe to restart indefinitely on failure. Without this, Flink defaults to "no-restart" + // when checkpoint is not enabled, causing the job to fail permanently on any exception. + flinkConfig.set(RestartStrategyOptions.RESTART_STRATEGY, EXPONENTIAL_DELAY.getMainValue()); execEnv = StreamExecutionEnvironment.getExecutionEnvironment(flinkConfig); } diff --git a/fluss-flink/fluss-flink-common/src/test/java/org/apache/fluss/flink/tiering/TieringFailoverITCase.java b/fluss-flink/fluss-flink-tiering/src/test/java/org/apache/fluss/flink/tiering/TieringFailoverITCase.java similarity index 79% rename from fluss-flink/fluss-flink-common/src/test/java/org/apache/fluss/flink/tiering/TieringFailoverITCase.java rename to fluss-flink/fluss-flink-tiering/src/test/java/org/apache/fluss/flink/tiering/TieringFailoverITCase.java index 07b2ddbd08..0f4d1a610e 100644 --- a/fluss-flink/fluss-flink-common/src/test/java/org/apache/fluss/flink/tiering/TieringFailoverITCase.java +++ b/fluss-flink/fluss-flink-tiering/src/test/java/org/apache/fluss/flink/tiering/TieringFailoverITCase.java @@ -18,16 +18,20 @@ package org.apache.fluss.flink.tiering; +import org.apache.fluss.config.ConfigOptions; import org.apache.fluss.lake.values.TestingValuesLake; +import org.apache.fluss.metadata.DataLakeFormat; import org.apache.fluss.metadata.Schema; import org.apache.fluss.metadata.TableBucket; import org.apache.fluss.metadata.TablePath; import org.apache.fluss.row.InternalRow; import org.apache.fluss.types.DataTypes; +import org.apache.flink.api.common.RuntimeExecutionMode; import org.apache.flink.core.execution.JobClient; import org.junit.jupiter.api.AfterAll; import org.junit.jupiter.api.BeforeAll; +import org.junit.jupiter.api.BeforeEach; import org.junit.jupiter.api.Test; import java.util.ArrayList; @@ -59,6 +63,24 @@ protected static void afterAll() throws Exception { FlinkTieringTestBase.afterAll(); } + @BeforeEach + @Override + void beforeEach() { + String bootstrapServers = String.join(",", clientConf.get(ConfigOptions.BOOTSTRAP_SERVERS)); + String[] args = { + "--fluss.bootstrap.servers", + bootstrapServers, + "--datalake.format", + DataLakeFormat.PAIMON.toString() + }; + FlussLakeTiering tiering = new FlussLakeTiering(args); + execEnv = tiering.execEnv; + execEnv.setRuntimeMode(RuntimeExecutionMode.STREAMING); + execEnv.setParallelism(2); + + assertThat(execEnv.getCheckpointConfig().isCheckpointingEnabled()).isFalse(); + } + @Test void testTiering() throws Exception { // create a pk table, write some records and wait until snapshot finished @@ -71,8 +93,8 @@ void testTiering() throws Exception { writeRows(t1, rows, false); FLUSS_CLUSTER_EXTENSION.triggerAndWaitSnapshot(t1); - // fail the first write to the pk table - TestingValuesLake.failWhen(t1.toString()).failWriteOnce(); + // fail the first two writes to the pk table + TestingValuesLake.failWhen(t1.toString()).failWriteNext(2); // then start tiering job JobClient jobClient = buildTieringJob(execEnv); @@ -96,7 +118,9 @@ void testTiering() throws Exception { checkDataInValuesTable(t1, expectedRows); } finally { - jobClient.cancel().get(); + if (!jobClient.getJobExecutionResult().isDone()) { + jobClient.cancel().get(); + } } } From dd6bcfdb39eee88bea8435960aaeae4809422c63 Mon Sep 17 00:00:00 2001 From: HuangXiao Date: Mon, 20 Jul 2026 19:53:48 +0800 Subject: [PATCH 2/2] Document tiering test dependency Explain why flink-table-common must be present on the fluss-flink-tiering test runtime classpath. --- fluss-flink/fluss-flink-tiering/pom.xml | 5 +++++ 1 file changed, 5 insertions(+) diff --git a/fluss-flink/fluss-flink-tiering/pom.xml b/fluss-flink/fluss-flink-tiering/pom.xml index 6dc98921dd..e504246a4e 100644 --- a/fluss-flink/fluss-flink-tiering/pom.xml +++ b/fluss-flink/fluss-flink-tiering/pom.xml @@ -117,6 +117,11 @@ test + org.apache.flink flink-table-common