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

Reply via email to