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 070be4d976 Issue #8701 : Run a Join once when parallel branches reach 
it at the same time (#8797)
070be4d976 is described below

commit 070be4d97688c5fd9f897088347eb56262be68de
Author: Joel Giovinazzo <[email protected]>
AuthorDate: Fri Oct 9 05:46:36 2026 +1100

    Issue #8701 : Run a Join once when parallel branches reach it at the same 
time (#8797)
---
 .../java/org/apache/hop/workflow/Workflow.java     |  75 +++++-
 .../hop/workflow/actions/join/ActionJoinTest.java  | 276 +++++++++++++++++++++
 2 files changed, 346 insertions(+), 5 deletions(-)

diff --git a/engine/src/main/java/org/apache/hop/workflow/Workflow.java 
b/engine/src/main/java/org/apache/hop/workflow/Workflow.java
index 7785afdbe3..8807569d26 100644
--- a/engine/src/main/java/org/apache/hop/workflow/Workflow.java
+++ b/engine/src/main/java/org/apache/hop/workflow/Workflow.java
@@ -28,6 +28,7 @@ import java.util.Map;
 import java.util.Queue;
 import java.util.Set;
 import java.util.UUID;
+import java.util.concurrent.ConcurrentHashMap;
 import java.util.concurrent.ConcurrentLinkedQueue;
 import java.util.concurrent.atomic.AtomicInteger;
 import org.apache.commons.vfs2.FileName;
@@ -166,6 +167,13 @@ public abstract class Workflow extends Variables
 
   protected Set<ActionMeta> activeActions;
 
+  /**
+   * Join actions that one branch has claimed to run. A branch claims a Join 
before it publishes its
+   * own result and the claim is released when the Join's body returns. A 
branch that cannot claim
+   * the Join leaves it to the pending run, which cannot complete without 
seeing that result.
+   */
+  protected Set<ActionMeta> claimedJoins;
+
   /** Parameters of the workflow. */
   protected INamedParameters namedParams = new NamedParameters();
 
@@ -214,6 +222,7 @@ public abstract class Workflow extends Variables
 
     // this map is being modified concurrently and must be thread-safe
     activeActions = Collections.synchronizedSet(new HashSet<>());
+    claimedJoins = ConcurrentHashMap.newKeySet();
 
     extensionDataMap = new HashMap<>();
 
@@ -410,6 +419,7 @@ public abstract class Workflow extends Variables
 
       setFinished(false);
       setStopped(false);
+      claimedJoins.clear();
       HopEnvironment.setExecutionInformation(this);
 
       log.logBasic(BaseMessages.getString(PKG, CONST_WORKFLOW_STARTED));
@@ -583,6 +593,7 @@ public abstract class Workflow extends Variables
     setFinished(false);
     setActive(true);
     setInitialized(true);
+    claimedJoins.clear();
     HopEnvironment.setExecutionInformation(this);
 
     // Where do we start?
@@ -756,6 +767,7 @@ public abstract class Workflow extends Variables
     Result res = null;
 
     if (isStopped()) {
+      releaseJoin(actionMeta);
       res = newResult();
       res.setEntryNr(nr);
       res.setStopped(true);
@@ -771,6 +783,7 @@ public abstract class Workflow extends Variables
     // if we didn't have a previous result, create one, otherwise, copy the 
content...
     //
     final Result newResult;
+    final Set<ActionMeta> joinClaims;
     Result prevResult = null;
     if (previousResult != null) {
       prevResult = previousResult.clone();
@@ -791,6 +804,8 @@ public abstract class Workflow extends Variables
 
     if (!extension.executeAction) {
       newResult = prevResult;
+      releaseJoin(actionMeta);
+      joinClaims = claimJoinsToFollow(actionMeta, newResult);
     } else {
       if (log.isDetailed()) {
         log.logDetailed(
@@ -839,7 +854,11 @@ public abstract class Workflow extends Variables
       activeActions.add(actionMeta.clone());
 
       log.snap(Metrics.METRIC_ACTION_START, cloneAction.toString());
-      newResult = cloneAction.execute(prevResult, nr);
+      try {
+        newResult = cloneAction.execute(prevResult, nr);
+      } finally {
+        releaseJoin(actionMeta);
+      }
       log.snap(Metrics.METRIC_ACTION_STOP, cloneAction.toString());
 
       // Action execution duration
@@ -848,6 +867,9 @@ public abstract class Workflow extends Variables
 
       activeActions.remove(actionMeta);
 
+      // Claim the Joins this action leads to before its result is published 
below
+      joinClaims = claimJoinsToFollow(actionMeta, newResult);
+
       for (IActionListener actionListener : actionListeners) {
         actionListener.afterExecution(this, actionMeta, cloneAction, 
newResult);
       }
@@ -953,11 +975,11 @@ public abstract class Workflow extends Variables
       // If the start point was an evaluation and the link color is correct:
       // green or red, execute the next action...
       //
-      if (hopMeta.isUnconditional()
-          || (actionMeta.isEvaluation() && (hopMeta.isEvaluation() == 
newResult.isResult()))) {
+      if (isHopFollowed(actionMeta, hopMeta, newResult)) {
 
-        // If the next action is a join, only execute once
-        if (nextAction.isJoin() && activeActions.contains(nextAction)) {
+        // Only the branch that claimed a Join runs it. Without a claim, 
another branch has a run
+        // of that Join pending, and that run waits for the result of this 
action.
+        if (nextAction.isJoin() && !joinClaims.remove(nextAction)) {
           continue;
         }
 
@@ -1023,6 +1045,11 @@ public abstract class Workflow extends Variables
       }
     }
 
+    // Release the claims on Joins that were not launched, for example because 
the workflow stopped
+    for (ActionMeta join : joinClaims) {
+      claimedJoins.remove(join);
+    }
+
     // OK, if we run in parallel, we need to wait for all the actions to
     // finish...
     //
@@ -1094,6 +1121,44 @@ public abstract class Workflow extends Variables
     return res;
   }
 
+  /** True when the hop from {@code actionMeta} is followed given the result 
of that action. */
+  private static boolean isHopFollowed(
+      ActionMeta actionMeta, WorkflowHopMeta hopMeta, Result result) {
+    return hopMeta.isUnconditional()
+        || (actionMeta.isEvaluation() && (hopMeta.isEvaluation() == 
result.isResult()));
+  }
+
+  /**
+   * Claims the Joins that {@code actionMeta} will follow given its result. 
Call this before the
+   * result is published to the workflow tracker: a Join run that is still 
pending then cannot
+   * complete without seeing that result, so a branch that fails to claim it 
can safely skip it.
+   *
+   * @return the Joins this branch claimed and must launch
+   */
+  private Set<ActionMeta> claimJoinsToFollow(ActionMeta actionMeta, Result 
result) {
+    Set<ActionMeta> claims = new HashSet<>();
+    int nrNext = workflowMeta.findNrNextActions(actionMeta);
+    for (int i = 0; i < nrNext; i++) {
+      ActionMeta nextAction = workflowMeta.findNextAction(actionMeta, i);
+      if (nextAction.isJoin()
+          && isHopFollowed(actionMeta, 
workflowMeta.findWorkflowHop(actionMeta, nextAction), result)
+          && claimedJoins.add(nextAction)) {
+        claims.add(nextAction);
+      }
+    }
+    return claims;
+  }
+
+  /**
+   * Releases the claim on a Join once its body has returned. A branch that 
finishes after this
+   * point was not seen by that run, so it may claim and run the Join again.
+   */
+  private void releaseJoin(ActionMeta actionMeta) {
+    if (actionMeta.isJoin()) {
+      claimedJoins.remove(actionMeta);
+    }
+  }
+
   /**
    * Get the number of errors that happened in the workflow.
    *
diff --git 
a/plugins/actions/join/src/test/java/org/apache/hop/workflow/actions/join/ActionJoinTest.java
 
b/plugins/actions/join/src/test/java/org/apache/hop/workflow/actions/join/ActionJoinTest.java
index d1ccaf6197..32b22180ee 100644
--- 
a/plugins/actions/join/src/test/java/org/apache/hop/workflow/actions/join/ActionJoinTest.java
+++ 
b/plugins/actions/join/src/test/java/org/apache/hop/workflow/actions/join/ActionJoinTest.java
@@ -27,16 +27,27 @@ import java.time.Duration;
 import java.util.ArrayList;
 import java.util.List;
 import java.util.concurrent.TimeUnit;
+import java.util.concurrent.atomic.AtomicInteger;
 import org.apache.commons.lang3.ThreadUtils;
 import org.apache.hop.core.Result;
+import org.apache.hop.core.extension.ExtensionPointPluginType;
+import org.apache.hop.core.extension.HopExtensionPoint;
+import org.apache.hop.core.extension.IExtensionPoint;
+import org.apache.hop.core.logging.ILogChannel;
 import org.apache.hop.core.logging.LogLevel;
+import org.apache.hop.core.plugins.IPlugin;
+import org.apache.hop.core.plugins.PluginRegistry;
 import org.apache.hop.core.variables.IVariables;
 import org.apache.hop.junit.rules.RestoreHopEngineEnvironmentExtension;
 import org.apache.hop.metadata.api.IHopMetadataProvider;
+import org.apache.hop.workflow.ActionResult;
+import org.apache.hop.workflow.IActionListener;
+import org.apache.hop.workflow.WorkflowExecutionExtension;
 import org.apache.hop.workflow.WorkflowHopMeta;
 import org.apache.hop.workflow.WorkflowMeta;
 import org.apache.hop.workflow.action.ActionBase;
 import org.apache.hop.workflow.action.ActionMeta;
+import org.apache.hop.workflow.action.IAction;
 import org.apache.hop.workflow.actions.dummy.ActionDummy;
 import org.apache.hop.workflow.actions.start.ActionStart;
 import org.apache.hop.workflow.engine.IWorkflowEngine;
@@ -344,6 +355,271 @@ class ActionJoinTest {
         
engine.getWorkflowTracker().findWorkflowTracker(slowMeta).getActionResult().getResult());
   }
 
+  @Test
+  @Timeout(value = 20, unit = TimeUnit.SECONDS)
+  void joinRunsOnceWhenBranchesArriveTogether() {
+    WorkflowMeta meta = new WorkflowMeta();
+    meta.setName("join-branches-arrive-together");
+
+    ActionMeta startMeta = new ActionMeta(new ActionStart("Start"));
+    startMeta.setLaunchingInParallel(true);
+    ActionMeta joinMeta = new ActionMeta(new ActionJoin("Join", ""));
+    ActionMeta afterMeta = new ActionMeta(new ActionDummy("After join"));
+    meta.addAction(startMeta);
+    meta.addAction(joinMeta);
+    meta.addAction(afterMeta);
+
+    for (int i = 0; i < 6; i++) {
+      ActionMeta branchMeta = new ActionMeta(new ActionDummy("Branch " + i));
+      meta.addAction(branchMeta);
+      meta.addWorkflowHop(new WorkflowHopMeta(startMeta, branchMeta));
+      WorkflowHopMeta branchToJoin = new WorkflowHopMeta(branchMeta, joinMeta);
+      branchToJoin.setUnconditional();
+      meta.addWorkflowHop(branchToJoin);
+    }
+    meta.addWorkflowHop(new WorkflowHopMeta(joinMeta, afterMeta));
+
+    LocalWorkflowEngine engine = new LocalWorkflowEngine(meta);
+    engine.setLogLevel(LogLevel.MINIMAL);
+    // Hold the Join between the moment a branch decides to start it and the 
moment it runs, so
+    // every branch reaches the Join while the first run of it is starting.
+    engine.addActionListener(
+        new IActionListener<WorkflowMeta>() {
+          @Override
+          public void beforeExecution(
+              IWorkflowEngine<WorkflowMeta> workflow, ActionMeta actionMeta, 
IAction action) {
+            if (actionMeta.isJoin()) {
+              sleep(300);
+            }
+          }
+
+          @Override
+          public void afterExecution(
+              IWorkflowEngine<WorkflowMeta> workflow,
+              ActionMeta actionMeta,
+              IAction action,
+              Result result) {
+            // Nothing to do
+          }
+        });
+    Result result = engine.startExecution();
+
+    assertTrue(result.isResult());
+    assertEquals(1, countRuns(engine, "Join"));
+    assertEquals(1, countRuns(engine, "After join"));
+  }
+
+  @Test
+  @Timeout(value = 20, unit = TimeUnit.SECONDS)
+  void joinRunsOnceWhenBranchArrivesAfterJoinCompleted() throws Exception {
+    WorkflowMeta meta = new WorkflowMeta();
+    meta.setName("join-branch-arrives-late");
+
+    ActionMeta startMeta = new ActionMeta(new ActionStart("Start"));
+    startMeta.setLaunchingInParallel(true);
+    ActionMeta fastMeta = new ActionMeta(new ActionDummy("Fast branch"));
+    ActionMeta lateMeta = new ActionMeta(new 
ActionDummy(LateArrivalExtension.LATE_BRANCH));
+    ActionMeta joinMeta = new ActionMeta(new ActionJoin("Join", ""));
+    ActionMeta afterMeta = new ActionMeta(new ActionDummy("After join"));
+    meta.addAction(startMeta);
+    meta.addAction(fastMeta);
+    meta.addAction(lateMeta);
+    meta.addAction(joinMeta);
+    meta.addAction(afterMeta);
+
+    meta.addWorkflowHop(new WorkflowHopMeta(startMeta, fastMeta));
+    meta.addWorkflowHop(new WorkflowHopMeta(startMeta, lateMeta));
+    WorkflowHopMeta fastToJoin = new WorkflowHopMeta(fastMeta, joinMeta);
+    fastToJoin.setUnconditional();
+    meta.addWorkflowHop(fastToJoin);
+    WorkflowHopMeta lateToJoin = new WorkflowHopMeta(lateMeta, joinMeta);
+    lateToJoin.setUnconditional();
+    meta.addWorkflowHop(lateToJoin);
+    meta.addWorkflowHop(new WorkflowHopMeta(joinMeta, afterMeta));
+
+    // Hold the late branch after its result is published, until the Join has 
seen that result and
+    // completed, and only then let it decide whether to start the Join.
+    ExtensionPointPluginType.getInstance()
+        .registerCustom(
+            LateArrivalExtension.class,
+            "test",
+            LateArrivalExtension.PLUGIN_ID,
+            HopExtensionPoint.WorkflowAfterActionExecution.id,
+            "Delays the late branch of a Join test",
+            null);
+    try {
+      LocalWorkflowEngine engine = new LocalWorkflowEngine(meta);
+      engine.setLogLevel(LogLevel.MINIMAL);
+      Result result = engine.startExecution();
+
+      assertTrue(result.isResult());
+      assertEquals(1, countRuns(engine, "Join"));
+      assertEquals(1, countRuns(engine, "After join"));
+    } finally {
+      PluginRegistry registry = PluginRegistry.getInstance();
+      IPlugin plugin =
+          registry.getPlugin(ExtensionPointPluginType.class, 
LateArrivalExtension.PLUGIN_ID);
+      if (plugin != null) {
+        registry.removePlugin(ExtensionPointPluginType.class, plugin);
+      }
+    }
+  }
+
+  @Test
+  @Timeout(value = 20, unit = TimeUnit.SECONDS)
+  void joinRunsOncePerLoopIteration() {
+    WorkflowMeta meta = new WorkflowMeta();
+    meta.setName("join-in-loop");
+
+    ActionMeta startMeta = new ActionMeta(new ActionStart("Start"));
+    ActionMeta fanOutMeta = new ActionMeta(new ActionDummy("Fan out"));
+    fanOutMeta.setLaunchingInParallel(true);
+    ActionMeta aMeta = new ActionMeta(new ActionDummy("Branch A"));
+    ActionMeta bMeta = new ActionMeta(new ActionDummy("Branch B"));
+    ActionMeta joinMeta = new ActionMeta(new ActionJoin("Join", ""));
+    ActionMeta counterMeta = new ActionMeta(new CountingEvalAction("Counter", 
3));
+    meta.addAction(startMeta);
+    meta.addAction(fanOutMeta);
+    meta.addAction(aMeta);
+    meta.addAction(bMeta);
+    meta.addAction(joinMeta);
+    meta.addAction(counterMeta);
+
+    meta.addWorkflowHop(new WorkflowHopMeta(startMeta, fanOutMeta));
+    for (ActionMeta branchMeta : List.of(aMeta, bMeta)) {
+      WorkflowHopMeta fanOutToBranch = new WorkflowHopMeta(fanOutMeta, 
branchMeta);
+      fanOutToBranch.setUnconditional();
+      meta.addWorkflowHop(fanOutToBranch);
+      WorkflowHopMeta branchToJoin = new WorkflowHopMeta(branchMeta, joinMeta);
+      branchToJoin.setUnconditional();
+      meta.addWorkflowHop(branchToJoin);
+    }
+    meta.addWorkflowHop(new WorkflowHopMeta(joinMeta, counterMeta));
+    // Loop back to the fan-out while the counter succeeds
+    meta.addWorkflowHop(new WorkflowHopMeta(counterMeta, fanOutMeta));
+
+    CountingEvalAction.RUNS.set(0);
+    LocalWorkflowEngine engine = new LocalWorkflowEngine(meta);
+    engine.setLogLevel(LogLevel.MINIMAL);
+    engine.startExecution();
+
+    assertEquals(3, countRuns(engine, "Join"));
+    assertEquals(3, countRuns(engine, "Counter"));
+  }
+
+  /**
+   * A branch that has another hop before its hop to the Join claims the Join, 
then follows its
+   * other hop first. The Join runs once, after that other hop has finished, 
even though the other
+   * branch arrived at the Join earlier.
+   */
+  @Test
+  @Timeout(value = 20, unit = TimeUnit.SECONDS)
+  void joinWaitsForEarlierHopOfBranchThatClaimedIt() {
+    WorkflowMeta meta = new WorkflowMeta();
+    meta.setName("join-claimed-by-branch-with-earlier-hop");
+
+    ActionMeta startMeta = new ActionMeta(new ActionStart("Start"));
+    startMeta.setLaunchingInParallel(true);
+    ActionMeta claimingMeta = new ActionMeta(new ActionDummy("Claiming 
branch"));
+    // Arrives at the Join after the claiming branch has claimed it
+    ActionMeta otherMeta = new ActionMeta(new SleepingEvalAction("Other 
branch", 200));
+    ActionMeta slowMeta = new ActionMeta(new SleepingEvalAction("Slow task", 
1000));
+    ActionMeta joinMeta = new ActionMeta(new ActionJoin("Join", ""));
+    ActionMeta afterMeta = new ActionMeta(new ActionDummy("After join"));
+    meta.addAction(startMeta);
+    meta.addAction(claimingMeta);
+    meta.addAction(otherMeta);
+    meta.addAction(slowMeta);
+    meta.addAction(joinMeta);
+    meta.addAction(afterMeta);
+
+    meta.addWorkflowHop(new WorkflowHopMeta(startMeta, claimingMeta));
+    meta.addWorkflowHop(new WorkflowHopMeta(startMeta, otherMeta));
+    // The hop to the slow task comes before the hop to the Join, so it is 
followed first
+    WorkflowHopMeta claimingToSlow = new WorkflowHopMeta(claimingMeta, 
slowMeta);
+    claimingToSlow.setUnconditional();
+    meta.addWorkflowHop(claimingToSlow);
+    WorkflowHopMeta claimingToJoin = new WorkflowHopMeta(claimingMeta, 
joinMeta);
+    claimingToJoin.setUnconditional();
+    meta.addWorkflowHop(claimingToJoin);
+    WorkflowHopMeta otherToJoin = new WorkflowHopMeta(otherMeta, joinMeta);
+    otherToJoin.setUnconditional();
+    meta.addWorkflowHop(otherToJoin);
+    meta.addWorkflowHop(new WorkflowHopMeta(joinMeta, afterMeta));
+
+    LocalWorkflowEngine engine = new LocalWorkflowEngine(meta);
+    engine.setLogLevel(LogLevel.MINIMAL);
+    Result result = engine.startExecution();
+
+    assertTrue(result.isResult());
+    assertEquals(1, countRuns(engine, "Join"));
+    assertEquals(1, countRuns(engine, "After join"));
+    assertTrue(
+        indexOfRun(engine, "Slow task") < indexOfRun(engine, "After join"),
+        "The Join runs only after the earlier hop of the branch that claimed 
it");
+  }
+
+  private static int indexOfRun(LocalWorkflowEngine engine, String actionName) 
{
+    List<ActionResult> actionResults = engine.getActionResults();
+    for (int i = 0; i < actionResults.size(); i++) {
+      if (actionName.equals(actionResults.get(i).getActionName())) {
+        return i;
+      }
+    }
+    return -1;
+  }
+
+  private static long countRuns(LocalWorkflowEngine engine, String actionName) 
{
+    return engine.getActionResults().stream()
+        .filter(actionResult -> 
actionName.equals(actionResult.getActionName()))
+        .count();
+  }
+
+  private static void sleep(long millis) {
+    try {
+      ThreadUtils.sleep(Duration.ofMillis(millis));
+    } catch (InterruptedException e) {
+      Thread.currentThread().interrupt();
+    }
+  }
+
+  public static class LateArrivalExtension implements 
IExtensionPoint<WorkflowExecutionExtension> {
+    static final String PLUGIN_ID = "ActionJoinTestLateArrival";
+    static final String LATE_BRANCH = "Late branch";
+
+    @Override
+    public void callExtensionPoint(
+        ILogChannel log, IVariables variables, WorkflowExecutionExtension 
extension) {
+      if (LATE_BRANCH.equals(extension.actionMeta.getName())) {
+        // Longer than the polling interval of the Join
+        sleep(1500);
+      }
+    }
+  }
+
+  /** Succeeds until it has run {@code maxRuns} times. */
+  static class CountingEvalAction extends ActionBase {
+    static final AtomicInteger RUNS = new AtomicInteger();
+    private final int maxRuns;
+
+    CountingEvalAction(String name, int maxRuns) {
+      super(name, "");
+      this.maxRuns = maxRuns;
+    }
+
+    @Override
+    public Result execute(Result result, int nr) {
+      result.setResult(RUNS.incrementAndGet() < maxRuns);
+      result.setNrErrors(0);
+      return result;
+    }
+
+    @Override
+    public boolean isEvaluation() {
+      return true;
+    }
+  }
+
   static class FailingEvalAction extends ActionBase {
     FailingEvalAction(String name) {
       super(name, "");

Reply via email to