This is an automated email from the ASF dual-hosted git repository.
bamaer pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/hop.git
The following commit(s) were added to refs/heads/main by this push:
new b9b062a083 Issue #3131 : Expose the submitted Flink job id on the
pipeline engine (#8679)
b9b062a083 is described below
commit b9b062a0839beb2f85af6adb0853827b6d336459
Author: Matt Casters <[email protected]>
AuthorDate: Thu Oct 1 11:33:50 2026 +0200
Issue #3131 : Expose the submitted Flink job id on the pipeline engine
(#8679)
* Issue #3131 : Expose the submitted Flink job id on the pipeline engine
* Issue #3131 : Only fix the Flink job id for local and host:port masters
---------
Co-authored-by: Bart Maertens <[email protected]>
---
.../beam-flink-pipeline-engine.adoc | 13 +++
.../engines/flink/BeamFlinkPipelineEngine.java | 100 ++++++++++++++++
.../beam/engines/flink/FlinkJobConfiguration.java | 126 +++++++++++++++++++++
.../beam/gui/PipelineExecutionViewerUpdateXP.java | 18 +++
.../flink/messages/messages_en_US.properties | 1 +
.../beam/gui/messages/messages_en_US.properties | 1 +
.../engines/flink/BeamFlinkPipelineEngineTest.java | 100 ++++++++++++++++
.../engines/flink/FlinkJobConfigurationTest.java | 118 +++++++++++++++++++
8 files changed, 477 insertions(+)
diff --git
a/docs/hop-user-manual/modules/ROOT/pages/pipeline/pipeline-run-configurations/beam-flink-pipeline-engine.adoc
b/docs/hop-user-manual/modules/ROOT/pages/pipeline/pipeline-run-configurations/beam-flink-pipeline-engine.adoc
index 92b060b191..df4ea9377b 100644
---
a/docs/hop-user-manual/modules/ROOT/pages/pipeline/pipeline-run-configurations/beam-flink-pipeline-engine.adoc
+++
b/docs/hop-user-manual/modules/ROOT/pages/pipeline/pipeline-run-configurations/beam-flink-pipeline-engine.adoc
@@ -88,6 +88,19 @@ Set this to BATCH_FORCED if pipelines get blocked, see
https://issues.apache.org
|Fat jar file location|Fat jar location.|
|===
+== Job id
+
+Beam does not return the Flink job id on the pipeline result. Hop chooses the
id before submitting the job. After `prepareExecution()` it is available from
the engine:
+
+[source,java]
+----
+if (pipeline instanceof BeamFlinkPipelineEngine flinkEngine) {
+ String jobId = flinkEngine.getFlinkJobId();
+}
+----
+
+The same value is written to the pipeline log. When the run configuration has
an execution information location, it is stored as `flink.job.id` and shown on
the execution viewer info tab. A master of `[collection]` does not submit a
cluster job, so no id is assigned.
+
== Running with Flink Run
You can also execute using the 'bin/flink run' command.
diff --git
a/plugins/engines/beam/src/main/java/org/apache/hop/beam/engines/flink/BeamFlinkPipelineEngine.java
b/plugins/engines/beam/src/main/java/org/apache/hop/beam/engines/flink/BeamFlinkPipelineEngine.java
index 0a98162a9a..334e468d30 100644
---
a/plugins/engines/beam/src/main/java/org/apache/hop/beam/engines/flink/BeamFlinkPipelineEngine.java
+++
b/plugins/engines/beam/src/main/java/org/apache/hop/beam/engines/flink/BeamFlinkPipelineEngine.java
@@ -17,8 +17,17 @@
package org.apache.hop.beam.engines.flink;
+import java.io.IOException;
+import java.nio.file.Path;
+import java.util.HashMap;
+import lombok.Getter;
+import org.apache.beam.runners.flink.FlinkPipelineOptions;
+import org.apache.commons.lang3.StringUtils;
+import org.apache.flink.api.common.JobID;
import org.apache.hop.beam.engines.BeamPipelineEngine;
import org.apache.hop.core.exception.HopException;
+import org.apache.hop.execution.ExecutionState;
+import org.apache.hop.i18n.BaseMessages;
import org.apache.hop.pipeline.PipelineMeta;
import org.apache.hop.pipeline.config.IPipelineEngineRunConfiguration;
import org.apache.hop.pipeline.engine.IPipelineEngine;
@@ -30,6 +39,24 @@ import org.apache.hop.pipeline.engine.PipelineEnginePlugin;
description = "This is a Flink pipeline engine provided by the Apache Beam
community")
public class BeamFlinkPipelineEngine extends BeamPipelineEngine
implements IPipelineEngine<PipelineMeta> {
+
+ private static final Class<?> PKG = BeamFlinkPipelineEngine.class;
+
+ /** Execution-information detail for the Flink job id submitted for this
pipeline. */
+ public static final String DETAIL_FLINK_JOB_ID = "flink.job.id";
+
+ private static final String COLLECTION_MASTER = "[collection]";
+
+ private static final String AUTO_MASTER = "[auto]";
+
+ /**
+ * Job id Flink will submit. Null for the collection master, which does not
start a job, and for
+ * the auto master, where Flink ignores the configuration directory that
carries the id.
+ */
+ @Getter private String flinkJobId;
+
+ private Path flinkJobConfDir;
+
@Override
public IPipelineEngineRunConfiguration
createDefaultPipelineEngineRunConfiguration() {
BeamFlinkPipelineRunConfiguration runConfiguration = new
BeamFlinkPipelineRunConfiguration();
@@ -46,4 +73,77 @@ public class BeamFlinkPipelineEngine extends
BeamPipelineEngine
+ engineRunConfiguration.getClass().getName());
}
}
+
+ @Override
+ public void prepareExecution() throws HopException {
+ super.prepareExecution();
+ assignFlinkJobId();
+ }
+
+ @Override
+ public void startThreads() throws HopException {
+ try {
+ super.startThreads();
+ } finally {
+ deleteFlinkJobConfiguration();
+ }
+ }
+
+ @Override
+ protected ExecutionState capturePipelineExecutionState() {
+ ExecutionState executionState = super.capturePipelineExecutionState();
+ if (StringUtils.isNotEmpty(flinkJobId)) {
+ if (executionState.getDetails() == null) {
+ executionState.setDetails(new HashMap<>());
+ }
+ executionState.getDetails().put(DETAIL_FLINK_JOB_ID, flinkJobId);
+ }
+ return executionState;
+ }
+
+ /**
+ * Fix the id Flink submits. {@code FlinkRunnerResult} keeps no job id,
including after a failed
+ * attached run.
+ */
+ private void assignFlinkJobId() throws HopException {
+ deleteFlinkJobConfiguration();
+ flinkJobId = null;
+ if (getBeamPipeline() == null || getBeamPipeline().getOptions() == null) {
+ return;
+ }
+ FlinkPipelineOptions options =
getBeamPipeline().getOptions().as(FlinkPipelineOptions.class);
+ if (!acceptsFixedJobId(options.getFlinkMaster())) {
+ return;
+ }
+ String jobId = new JobID().toHexString();
+ flinkJobConfDir = FlinkJobConfiguration.create(options.getFlinkConfDir(),
jobId);
+ options.setFlinkConfDir(flinkJobConfDir.toAbsolutePath().toString());
+ flinkJobId = jobId;
+ logChannel.logBasic(BaseMessages.getString(PKG,
"BeamEnginesFlink.JobId.Log", jobId));
+ }
+
+ /**
+ * Only a local or host:port master builds its Flink environment from the
configuration directory
+ * Hop hands over. The auto master (also Beam's default for an empty master)
takes the environment
+ * of {@code flink run} or the Kubernetes operator instead, so a fixed job
id would never reach
+ * Flink and the reported id would be wrong.
+ */
+ static boolean acceptsFixedJobId(String flinkMaster) {
+ return StringUtils.isNotBlank(flinkMaster)
+ && !COLLECTION_MASTER.equals(flinkMaster.trim())
+ && !AUTO_MASTER.equals(flinkMaster.trim());
+ }
+
+ private void deleteFlinkJobConfiguration() {
+ Path dir = flinkJobConfDir;
+ flinkJobConfDir = null;
+ if (dir == null) {
+ return;
+ }
+ try {
+ FlinkJobConfiguration.delete(dir);
+ } catch (IOException e) {
+ logChannel.logDebug("Could not remove temporary Flink configuration
directory " + dir, e);
+ }
+ }
}
diff --git
a/plugins/engines/beam/src/main/java/org/apache/hop/beam/engines/flink/FlinkJobConfiguration.java
b/plugins/engines/beam/src/main/java/org/apache/hop/beam/engines/flink/FlinkJobConfiguration.java
new file mode 100644
index 0000000000..92da00700e
--- /dev/null
+++
b/plugins/engines/beam/src/main/java/org/apache/hop/beam/engines/flink/FlinkJobConfiguration.java
@@ -0,0 +1,126 @@
+/*
+ * 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.hop.beam.engines.flink;
+
+import java.io.IOException;
+import java.nio.charset.StandardCharsets;
+import java.nio.file.Files;
+import java.nio.file.Path;
+import java.util.Comparator;
+import java.util.stream.Stream;
+import org.apache.commons.lang3.StringUtils;
+import org.apache.flink.api.common.JobID;
+import org.apache.flink.configuration.GlobalConfiguration;
+import org.apache.flink.configuration.PipelineOptionsInternal;
+import org.apache.hop.core.exception.HopException;
+
+/**
+ * Beam's {@code FlinkRunnerResult} does not expose the Flink job id. Flink
applies {@code
+ * $internal.pipeline.job-id} from the configuration directory when it builds
the job graph, so the
+ * id can be fixed before submission.
+ *
+ * <p>Flink only reads that directory from the local filesystem.
+ */
+final class FlinkJobConfiguration {
+
+ private static final String JOB_ID_KEY =
PipelineOptionsInternal.PIPELINE_FIXED_JOB_ID.key();
+
+ private FlinkJobConfiguration() {}
+
+ static Path create(String configuredConfDir, String jobId) throws
HopException {
+ String source = configuredConfDir;
+ if (StringUtils.isEmpty(source)) {
+ source = System.getenv("FLINK_CONF_DIR");
+ }
+ Path sourceDir = StringUtils.isEmpty(source) ? null : Path.of(source);
+ return createFromSource(sourceDir, jobId);
+ }
+
+ static Path createFromSource(Path sourceDir, String jobId) throws
HopException {
+ try {
+ JobID.fromHexString(jobId);
+ Path dir = Files.createTempDirectory("hop-flink-");
+ if (sourceDir == null) {
+ Files.writeString(
+ dir.resolve(GlobalConfiguration.FLINK_CONF_FILENAME),
+ standardEntry(jobId),
+ StandardCharsets.UTF_8);
+ } else if (!Files.isDirectory(sourceDir)) {
+ throw new HopException("Flink configuration directory does not exist:
" + sourceDir);
+ } else {
+ Path legacy =
sourceDir.resolve(GlobalConfiguration.LEGACY_FLINK_CONF_FILENAME);
+ Path standard =
sourceDir.resolve(GlobalConfiguration.FLINK_CONF_FILENAME);
+ Path sourceFile;
+ String entry;
+ if (Files.isRegularFile(legacy)) {
+ sourceFile = legacy;
+ entry = legacyEntry(jobId);
+ } else if (Files.isRegularFile(standard)) {
+ sourceFile = standard;
+ entry = standardEntry(jobId);
+ } else {
+ throw new HopException(
+ "Flink configuration directory '"
+ + sourceDir
+ + "' does not contain "
+ + GlobalConfiguration.FLINK_CONF_FILENAME
+ + " or "
+ + GlobalConfiguration.LEGACY_FLINK_CONF_FILENAME);
+ }
+ Path target = dir.resolve(sourceFile.getFileName());
+ Files.copy(sourceFile, target);
+ appendJobId(target, entry);
+ }
+ return dir;
+ } catch (HopException e) {
+ throw e;
+ } catch (Exception e) {
+ throw new HopException(
+ "Unable to create a Flink configuration directory for job id " +
jobId, e);
+ }
+ }
+
+ static void delete(Path dir) throws IOException {
+ if (dir == null || !Files.exists(dir)) {
+ return;
+ }
+ try (Stream<Path> walk = Files.walk(dir)) {
+ for (Path path : walk.sorted(Comparator.reverseOrder()).toList()) {
+ Files.deleteIfExists(path);
+ }
+ }
+ }
+
+ private static void appendJobId(Path file, String entry) throws IOException {
+ StringBuilder kept = new StringBuilder();
+ for (String line : Files.readString(file).split("\\R", -1)) {
+ if (!line.contains(JOB_ID_KEY)) {
+ kept.append(line).append('\n');
+ }
+ }
+ Files.writeString(file, kept + entry, StandardCharsets.UTF_8);
+ }
+
+ private static String standardEntry(String jobId) {
+ return "\"" + JOB_ID_KEY + "\": \"" + jobId + "\"\n";
+ }
+
+ private static String legacyEntry(String jobId) {
+ return JOB_ID_KEY + ": " + jobId + "\n";
+ }
+}
diff --git
a/plugins/engines/beam/src/main/java/org/apache/hop/beam/gui/PipelineExecutionViewerUpdateXP.java
b/plugins/engines/beam/src/main/java/org/apache/hop/beam/gui/PipelineExecutionViewerUpdateXP.java
index a5f47d3e60..d72c8a8436 100644
---
a/plugins/engines/beam/src/main/java/org/apache/hop/beam/gui/PipelineExecutionViewerUpdateXP.java
+++
b/plugins/engines/beam/src/main/java/org/apache/hop/beam/gui/PipelineExecutionViewerUpdateXP.java
@@ -20,12 +20,14 @@ package org.apache.hop.beam.gui;
import org.apache.commons.lang3.StringUtils;
import org.apache.hop.beam.engines.dataflow.BeamDataFlowPipelineEngine;
+import org.apache.hop.beam.engines.flink.BeamFlinkPipelineEngine;
import org.apache.hop.core.exception.HopException;
import org.apache.hop.core.extension.ExtensionPoint;
import org.apache.hop.core.extension.IExtensionPoint;
import org.apache.hop.core.logging.ILogChannel;
import org.apache.hop.core.variables.IVariables;
import org.apache.hop.execution.ExecutionState;
+import org.apache.hop.i18n.BaseMessages;
import org.apache.hop.ui.hopgui.perspective.execution.PipelineExecutionViewer;
@ExtensionPoint(
@@ -34,6 +36,9 @@ import
org.apache.hop.ui.hopgui.perspective.execution.PipelineExecutionViewer;
description = "Update the toolbar icons we add in the Beam GUI plugin")
public final class PipelineExecutionViewerUpdateXP
implements IExtensionPoint<PipelineExecutionViewer> {
+
+ private static final Class<?> PKG = HopBeamGuiPlugin.class;
+
@Override
public void callExtensionPoint(
ILogChannel log, IVariables variables, PipelineExecutionViewer viewer)
throws HopException {
@@ -50,5 +55,18 @@ public final class PipelineExecutionViewerUpdateXP
.enableToolbarItem(
HopBeamGuiPlugin.TOOLBAR_ID_PIPELINE_EXECUTION_VIEWER_VISIT_GCP_DATAFLOW,
StringUtils.isNotEmpty(jobId));
+
+ if (executionState != null
+ && executionState.getDetails() != null
+ && viewer.getInfoView() != null) {
+ String flinkJobId =
+
executionState.getDetails().get(BeamFlinkPipelineEngine.DETAIL_FLINK_JOB_ID);
+ if (StringUtils.isNotEmpty(flinkJobId)) {
+ viewer
+ .getInfoView()
+ .add(BaseMessages.getString(PKG,
"BeamGuiPlugin.FlinkJobId.Label"), flinkJobId);
+ viewer.getInfoView().optimizeTableView();
+ }
+ }
}
}
diff --git
a/plugins/engines/beam/src/main/resources/org/apache/hop/beam/engines/flink/messages/messages_en_US.properties
b/plugins/engines/beam/src/main/resources/org/apache/hop/beam/engines/flink/messages/messages_en_US.properties
index 2f9974e092..daabd7c0e3 100644
---
a/plugins/engines/beam/src/main/resources/org/apache/hop/beam/engines/flink/messages/messages_en_US.properties
+++
b/plugins/engines/beam/src/main/resources/org/apache/hop/beam/engines/flink/messages/messages_en_US.properties
@@ -55,3 +55,4 @@ BeamEnginesFlink.OptionsRetryDelay.Label=Execution retry
delay (ms)
BeamEnginesFlink.OptionsRetryDelay.ToolTip=Sets the delay in milliseconds
between executions. A value of -1 indicates that the default value should be
used.
BeamEnginesFlink.OptionsShutdownSourcesAfterIdleMs.Label=Shutdown sources on
final watermark
BeamEnginesFlink.OptionsShutdownSourcesAfterIdleMs.ToolTip=Shuts down sources
which have been idle for the configured time of milliseconds. Once a source has
been shut down, checkpointing is not possible anymore. Shutting down the
sources eventually leads to pipeline shutdown (=Flink job finishes) once all
input has been processed. Unless explicitly set, this will default to
Long.MAX_VALUE when checkpointing is enabled and to 0 when checkpointing is
disabled. See https://issues.apache.or [...]
+BeamEnginesFlink.JobId.Log=Flink job id: {0}
diff --git
a/plugins/engines/beam/src/main/resources/org/apache/hop/beam/gui/messages/messages_en_US.properties
b/plugins/engines/beam/src/main/resources/org/apache/hop/beam/gui/messages/messages_en_US.properties
index 8491ebfbfb..650883bd69 100644
---
a/plugins/engines/beam/src/main/resources/org/apache/hop/beam/gui/messages/messages_en_US.properties
+++
b/plugins/engines/beam/src/main/resources/org/apache/hop/beam/gui/messages/messages_en_US.properties
@@ -34,3 +34,4 @@ BeamGuiPlugin.Menu.GenerateFatJar.Text=Generate a Hop fat
jar...
BeamGuiPlugin.OpenDataflowJob.Dialog.Header=Open Dataflow Job
BeamGuiPlugin.OpenDataflowJob.Dialog.Message=Open dataflow job by url:
{0},\n{1}
BeamGuiPlugin.VisitDataflow.ToolTip=Visit the pipeline execution in the GCP
Dataflow console
+BeamGuiPlugin.FlinkJobId.Label=Flink job ID
diff --git
a/plugins/engines/beam/src/test/java/org/apache/hop/beam/engines/flink/BeamFlinkPipelineEngineTest.java
b/plugins/engines/beam/src/test/java/org/apache/hop/beam/engines/flink/BeamFlinkPipelineEngineTest.java
index e04b622ce1..6aa008ab70 100644
---
a/plugins/engines/beam/src/test/java/org/apache/hop/beam/engines/flink/BeamFlinkPipelineEngineTest.java
+++
b/plugins/engines/beam/src/test/java/org/apache/hop/beam/engines/flink/BeamFlinkPipelineEngineTest.java
@@ -18,15 +18,28 @@
package org.apache.hop.beam.engines.flink;
import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertFalse;
+import static org.junit.jupiter.api.Assertions.assertNotNull;
+import static org.junit.jupiter.api.Assertions.assertNull;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+import java.nio.file.Path;
import java.util.Arrays;
+import org.apache.beam.runners.flink.FlinkPipelineOptions;
+import org.apache.flink.api.common.JobID;
+import org.apache.flink.configuration.GlobalConfiguration;
+import org.apache.flink.configuration.PipelineOptionsInternal;
import org.apache.hop.beam.engines.BeamBasePipelineEngineTest;
import org.apache.hop.beam.util.BeamPipelineMetaUtil;
import org.apache.hop.core.variables.DescribedVariable;
+import org.apache.hop.execution.ExecutionState;
import org.apache.hop.pipeline.PipelineMeta;
import org.apache.hop.pipeline.config.PipelineRunConfiguration;
import org.apache.hop.pipeline.engine.IPipelineEngine;
+import org.apache.hop.pipeline.engine.PipelineEngineFactory;
import org.junit.jupiter.api.Test;
+import org.junit.jupiter.params.ParameterizedTest;
+import org.junit.jupiter.params.provider.ValueSource;
class BeamFlinkPipelineEngineTest extends BeamBasePipelineEngineTest {
@@ -60,5 +73,92 @@ class BeamFlinkPipelineEngineTest extends
BeamBasePipelineEngineTest {
validateInputOutputEngineMetrics(engine);
assertEquals("flink1", engine.getVariable("VAR1"));
+ assertNull(((BeamFlinkPipelineEngine) engine).getFlinkJobId());
+ }
+
+ @Test
+ void flinkJobIdIsFixedBeforeSubmission() throws Exception {
+ BeamFlinkPipelineRunConfiguration configuration =
+ new BeamFlinkPipelineRunConfiguration("127.0.0.1:9", "1");
+ configuration.setEnginePluginId("BeamFlinkPipelineEngine");
+ configuration.setTempLocation(System.getProperty("java.io.tmpdir"));
+ PipelineRunConfiguration pipelineRunConfiguration =
+ new PipelineRunConfiguration(
+ "flink-job-id",
+ "description",
+ "",
+ Arrays.asList(new DescribedVariable("VAR1", "flink1",
"description1")),
+ configuration,
+ null,
+ false);
+
metadataProvider.getSerializer(PipelineRunConfiguration.class).save(pipelineRunConfiguration);
+
+ PipelineMeta pipelineMeta =
+ BeamPipelineMetaUtil.generateBeamInputOutputPipelineMeta(
+ "flink-job-id", "INPUT", "OUTPUT", metadataProvider);
+ IPipelineEngine<PipelineMeta> engine =
+ PipelineEngineFactory.createPipelineEngine(
+ variables, pipelineRunConfiguration.getName(), metadataProvider,
pipelineMeta);
+ engine.prepareExecution();
+
+ Path confDir = null;
+ try {
+ BeamFlinkPipelineEngine flinkEngine = (BeamFlinkPipelineEngine) engine;
+ String jobId = flinkEngine.getFlinkJobId();
+ assertNotNull(jobId);
+ assertEquals(jobId, JobID.fromHexString(jobId).toHexString());
+
+ FlinkPipelineOptions options =
+
flinkEngine.getBeamPipeline().getOptions().as(FlinkPipelineOptions.class);
+ confDir = Path.of(options.getFlinkConfDir());
+ assertEquals(
+ jobId,
+ GlobalConfiguration.loadConfiguration(confDir.toString())
+ .get(PipelineOptionsInternal.PIPELINE_FIXED_JOB_ID));
+
+ ExecutionState state = flinkEngine.capturePipelineExecutionState();
+ assertEquals(jobId,
state.getDetails().get(BeamFlinkPipelineEngine.DETAIL_FLINK_JOB_ID));
+ } finally {
+ FlinkJobConfiguration.delete(confDir);
+ }
+ }
+
+ @ParameterizedTest
+ @ValueSource(strings = {"[local]", "127.0.0.1:8081",
"flink-jobmanager:8081"})
+ void fixedJobIdForLocalAndRemoteMasters(String master) {
+ assertTrue(BeamFlinkPipelineEngine.acceptsFixedJobId(master));
+ }
+
+ @ParameterizedTest
+ @ValueSource(strings = {"", " ", "[auto]", "[collection]"})
+ void noFixedJobIdWhereFlinkIgnoresTheConfDir(String master) {
+ assertFalse(BeamFlinkPipelineEngine.acceptsFixedJobId(master));
+ }
+
+ @Test
+ void noFlinkJobIdForAutoMaster() throws Exception {
+ BeamFlinkPipelineRunConfiguration configuration =
+ new BeamFlinkPipelineRunConfiguration("[auto]", "1");
+ configuration.setEnginePluginId("BeamFlinkPipelineEngine");
+ configuration.setTempLocation(System.getProperty("java.io.tmpdir"));
+ PipelineRunConfiguration pipelineRunConfiguration =
+ new PipelineRunConfiguration(
+ "flink-auto", "description", "", Arrays.asList(), configuration,
null, false);
+
metadataProvider.getSerializer(PipelineRunConfiguration.class).save(pipelineRunConfiguration);
+
+ PipelineMeta pipelineMeta =
+ BeamPipelineMetaUtil.generateBeamInputOutputPipelineMeta(
+ "flink-auto", "INPUT", "OUTPUT", metadataProvider);
+ IPipelineEngine<PipelineMeta> engine =
+ PipelineEngineFactory.createPipelineEngine(
+ variables, pipelineRunConfiguration.getName(), metadataProvider,
pipelineMeta);
+ engine.prepareExecution();
+
+ BeamFlinkPipelineEngine flinkEngine = (BeamFlinkPipelineEngine) engine;
+ assertNull(flinkEngine.getFlinkJobId());
+ ExecutionState state = flinkEngine.capturePipelineExecutionState();
+ assertTrue(
+ state.getDetails() == null
+ ||
!state.getDetails().containsKey(BeamFlinkPipelineEngine.DETAIL_FLINK_JOB_ID));
}
}
diff --git
a/plugins/engines/beam/src/test/java/org/apache/hop/beam/engines/flink/FlinkJobConfigurationTest.java
b/plugins/engines/beam/src/test/java/org/apache/hop/beam/engines/flink/FlinkJobConfigurationTest.java
new file mode 100644
index 0000000000..5fff09f6da
--- /dev/null
+++
b/plugins/engines/beam/src/test/java/org/apache/hop/beam/engines/flink/FlinkJobConfigurationTest.java
@@ -0,0 +1,118 @@
+/*
+ * 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.hop.beam.engines.flink;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertThrows;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+
+import java.nio.charset.StandardCharsets;
+import java.nio.file.Files;
+import java.nio.file.Path;
+import org.apache.flink.api.common.JobID;
+import org.apache.flink.client.deployment.executors.PipelineExecutorUtils;
+import org.apache.flink.configuration.Configuration;
+import org.apache.flink.configuration.CoreOptions;
+import org.apache.flink.configuration.GlobalConfiguration;
+import org.apache.flink.configuration.PipelineOptionsInternal;
+import org.apache.flink.runtime.jobgraph.JobGraph;
+import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
+import org.apache.flink.streaming.api.graph.StreamGraph;
+import org.apache.hop.core.exception.HopException;
+import org.junit.jupiter.api.Test;
+import org.junit.jupiter.api.io.TempDir;
+
+class FlinkJobConfigurationTest {
+
+ @TempDir Path tempDir;
+
+ @Test
+ void fixedJobIdIsUsedOnTheJobGraph() throws Exception {
+ String jobId = new JobID().toHexString();
+ Path confDir = FlinkJobConfiguration.createFromSource(null, jobId);
+ try {
+ Configuration loaded =
GlobalConfiguration.loadConfiguration(confDir.toString());
+ assertEquals(jobId,
loaded.get(PipelineOptionsInternal.PIPELINE_FIXED_JOB_ID));
+
+ StreamExecutionEnvironment env =
StreamExecutionEnvironment.createLocalEnvironment(1);
+ env.fromElements(1).print();
+ StreamGraph streamGraph = env.getStreamGraph();
+ streamGraph.setJobName("hop-flink-job-id");
+ JobGraph jobGraph =
+ PipelineExecutorUtils.getJobGraph(
+ streamGraph, loaded,
Thread.currentThread().getContextClassLoader());
+ assertEquals(jobId, jobGraph.getJobID().toHexString());
+ } finally {
+ FlinkJobConfiguration.delete(confDir);
+ }
+ }
+
+ @Test
+ void existingStandardConfigIsPreserved() throws Exception {
+ Path source = tempDir.resolve("standard");
+ Files.createDirectory(source);
+ Files.writeString(
+ source.resolve(GlobalConfiguration.FLINK_CONF_FILENAME),
+ "parallelism.default: 7\n\""
+ + PipelineOptionsInternal.PIPELINE_FIXED_JOB_ID.key()
+ + "\": \"00000000000000000000000000000000\"\n",
+ StandardCharsets.UTF_8);
+
+ String jobId = new JobID().toHexString();
+ Path confDir = FlinkJobConfiguration.createFromSource(source, jobId);
+ try {
+ Configuration loaded =
GlobalConfiguration.loadConfiguration(confDir.toString());
+ assertEquals(jobId,
loaded.get(PipelineOptionsInternal.PIPELINE_FIXED_JOB_ID));
+ assertEquals(7, loaded.get(CoreOptions.DEFAULT_PARALLELISM));
+
assertTrue(Files.exists(confDir.resolve(GlobalConfiguration.FLINK_CONF_FILENAME)));
+ } finally {
+ FlinkJobConfiguration.delete(confDir);
+ }
+ }
+
+ @Test
+ void existingLegacyConfigIsPreserved() throws Exception {
+ Path source = tempDir.resolve("legacy");
+ Files.createDirectory(source);
+ Files.writeString(
+ source.resolve(GlobalConfiguration.LEGACY_FLINK_CONF_FILENAME),
+ "parallelism.default: 3\n",
+ StandardCharsets.UTF_8);
+
+ String jobId = new JobID().toHexString();
+ Path confDir = FlinkJobConfiguration.createFromSource(source, jobId);
+ try {
+ Configuration loaded =
GlobalConfiguration.loadConfiguration(confDir.toString());
+ assertEquals(jobId,
loaded.get(PipelineOptionsInternal.PIPELINE_FIXED_JOB_ID));
+ assertEquals(3, loaded.get(CoreOptions.DEFAULT_PARALLELISM));
+ } finally {
+ FlinkJobConfiguration.delete(confDir);
+ }
+ }
+
+ @Test
+ void missingConfigurationDirectoryThrows() {
+ HopException exception =
+ assertThrows(
+ HopException.class,
+ () ->
+ FlinkJobConfiguration.createFromSource(
+ tempDir.resolve("missing"), new JobID().toHexString()));
+ assertTrue(exception.getMessage().contains("does not exist"));
+ }
+}