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(); + } } }