Skip to content

Commit ec6991e

Browse files
authored
[minor][lake] Set default name for flink tiering service job (apache#1105)
1 parent 960b9d7 commit ec6991e

File tree

2 files changed

+9
-2
lines changed

2 files changed

+9
-2
lines changed

fluss-flink/fluss-flink-common/src/main/java/com/alibaba/fluss/flink/tiering/LakeTieringJobBuilder.java

Lines changed: 8 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -28,6 +28,7 @@
2828
import com.alibaba.fluss.lake.writer.LakeTieringFactory;
2929

3030
import org.apache.flink.api.common.eventtime.WatermarkStrategy;
31+
import org.apache.flink.configuration.PipelineOptions;
3132
import org.apache.flink.core.execution.JobClient;
3233
import org.apache.flink.streaming.api.datastream.DataStreamSource;
3334
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
@@ -42,6 +43,8 @@
4243
/** The builder to build Flink lake tiering job. */
4344
public class LakeTieringJobBuilder {
4445

46+
private static final String DEFAULT_TIERING_SERVICE_JOB_NAME = "Fluss Lake Tiering Service";
47+
4548
private final StreamExecutionEnvironment env;
4649
private final Configuration flussConfig;
4750
private final Configuration dataLakeConfig;
@@ -106,7 +109,11 @@ public JobClient build() throws Exception {
106109
.setParallelism(1)
107110
.setMaxParallelism(1)
108111
.sinkTo(new DiscardingSink());
112+
String jobName =
113+
env.getConfiguration()
114+
.getOptional(PipelineOptions.NAME)
115+
.orElse(DEFAULT_TIERING_SERVICE_JOB_NAME);
109116

110-
return env.executeAsync();
117+
return env.executeAsync(jobName);
111118
}
112119
}

fluss-flink/fluss-flink-tiering/src/main/java/com/alibaba/fluss/flink/tiering/FlussLakeTieringEntrypoint.java

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -30,7 +30,7 @@
3030
import static com.alibaba.fluss.utils.PropertiesUtils.extractAndRemovePrefix;
3131
import static org.apache.flink.runtime.executiongraph.failover.FailoverStrategyFactoryLoader.FULL_RESTART_STRATEGY_NAME;
3232

33-
/** The entrypoint for Flink to tiering fluss data to Paimon. */
33+
/** The entrypoint for Flink to tier fluss data to lake format like paimon. */
3434
public class FlussLakeTieringEntrypoint {
3535

3636
private static final String FLUSS_CONF_PREFIX = "fluss.";

0 commit comments

Comments
 (0)