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..e504246a4e 100644
--- a/fluss-flink/fluss-flink-tiering/pom.xml
+++ b/fluss-flink/fluss-flink-tiering/pom.xml
@@ -57,6 +57,77 @@
${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 +156,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();
+ }
}
}