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, "");