This is an automated email from the ASF dual-hosted git repository.
hansva 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 29f09ddad9 Detect and mitigate classic pipeline buffer deadlocks
(#7752) (#7854)
29f09ddad9 is described below
commit 29f09ddad964ce83b4fb6d8ffd5e9ca81720556e
Author: Matt Casters <[email protected]>
AuthorDate: Thu Aug 13 14:54:33 2026 +0200
Detect and mitigate classic pipeline buffer deadlocks (#7752) (#7854)
* Detect and mitigate classic pipeline buffer deadlocks (#7752)
Add static split-rejoin analysis and optional spilling rowsets on the
local multi-threaded engine so Stream Lookup / Merge Join-style hangs
make progress instead of locking up, without extra consumer threads.
* Clean up spilling rowset temp files when a pipeline ends or is stopped
Rowset.clear() already deleted spill segments, but nothing invoked it on
normal finish or stop—only single-threaded clearError. Call cleanupRowSets
from the execution-finished path, cleanup(), and
disposeInitializedTransforms.
* Add Stream Lookup buffer-deadlock integration test (#7752)
Exercise a split-rejoin Stream Lookup with more rows than the local
rowset size. Enable detect/mitigate on the transforms local run
configuration so the classic engine spills and completes instead of
hanging.
* Fix SpillingRowSet javadoc: avoid @link to engine BaseTransform
core cannot resolve engine classes during javadoc generation.
* issue #7785 : Spotless
---
.../java/org/apache/hop/core/SpillingRowSet.java | 374 +++++++++++++++++++++
.../org/apache/hop/core/SpillingRowSetTest.java | 121 +++++++
.../pages/how-to-guides/avoiding-deadlocks.adoc | 19 ++
.../java/org/apache/hop/pipeline/Pipeline.java | 61 +++-
.../java/org/apache/hop/pipeline/PipelineMeta.java | 16 +
.../hop/pipeline/analysis/BufferDeadlockRisk.java | 79 +++++
.../analysis/PipelineBufferDeadlockAnalyzer.java | 190 +++++++++++
.../engines/local/LocalPipelineEngine.java | 77 +++++
.../local/LocalPipelineRunConfiguration.java | 52 +++
.../pipeline/messages/messages_en_US.properties | 1 +
.../PipelineBufferDeadlockAnalyzerTest.java | 136 ++++++++
.../0022-stream-lookup-buffer-deadlock.hpl | 331 ++++++++++++++++++
.../golden-stream-lookup-buffer-deadlock.csv | 3 +
.../transforms/main-0022-stream-lookup.hwf | 3 +
.../golden-stream-lookup-buffer-deadlock.json | 32 ++
.../metadata/pipeline-run-configuration/local.json | 5 +-
.../0022-stream-lookup-buffer-deadlock UNIT.json | 36 ++
.../config/messages/messages_en_US.properties | 6 +
18 files changed, 1540 insertions(+), 2 deletions(-)
diff --git a/core/src/main/java/org/apache/hop/core/SpillingRowSet.java
b/core/src/main/java/org/apache/hop/core/SpillingRowSet.java
new file mode 100644
index 0000000000..ce80bddfe2
--- /dev/null
+++ b/core/src/main/java/org/apache/hop/core/SpillingRowSet.java
@@ -0,0 +1,374 @@
+/*
+ * 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.core;
+
+import java.io.BufferedInputStream;
+import java.io.BufferedOutputStream;
+import java.io.DataInputStream;
+import java.io.DataOutputStream;
+import java.io.File;
+import java.io.IOException;
+import java.io.InputStream;
+import java.io.OutputStream;
+import java.util.ArrayDeque;
+import java.util.ArrayList;
+import java.util.List;
+import java.util.concurrent.TimeUnit;
+import org.apache.commons.vfs2.FileObject;
+import org.apache.hop.core.exception.HopFileException;
+import org.apache.hop.core.row.IRowMeta;
+import org.apache.hop.core.vfs.HopVfs;
+import org.jspecify.annotations.NonNull;
+
+/**
+ * Bounded in-memory rowset that spills excess rows to temporary files instead
of blocking the
+ * producer when the memory capacity is full.
+ *
+ * <p>While the in-memory queue has room and nothing is waiting on disk,
behaviour matches {@link
+ * BlockingRowSet}. When the queue is full (or unread spilled rows already
exist, so FIFO order must
+ * be preserved), further {@link #putRow} calls serialize rows to a temp file
and return {@code
+ * true} without waiting on the consumer. {@link #size()} reports the
in-memory size only so
+ * existing backpressure heuristics on other hops are unchanged.
+ *
+ * <p>Temp files are created under a configurable directory (default: {@code
java.io.tmpdir}) via
+ * {@link HopVfs}.
+ */
+public class SpillingRowSet extends BaseRowSet implements Comparable<IRowSet>,
IRowSet {
+
+ private final int capacity;
+ private final String directory;
+ private final ArrayDeque<Object[]> memory;
+ private final Object lock = new Object();
+
+ /** Rows written to spill that have not yet been read back. */
+ private long unreadSpilled;
+
+ private final List<SpillSegment> segments = new ArrayList<>();
+ private SpillSegment writeSegment;
+ private int readSegmentIndex;
+ private long readRowsInSegment;
+
+ private boolean firstSpillLogged;
+ private volatile boolean spillIoFailed;
+
+ private final int timeoutGet;
+
+ public SpillingRowSet(int maxSize) {
+ this(maxSize, null);
+ }
+
+ /**
+ * @param maxSize in-memory capacity (same role as {@link BlockingRowSet})
+ * @param directory spill directory; null or blank uses {@code
java.io.tmpdir}
+ */
+ public SpillingRowSet(int maxSize, String directory) {
+ super();
+ if (maxSize < 1) {
+ throw new IllegalArgumentException("SpillingRowSet capacity must be >=
1");
+ }
+ this.capacity = maxSize;
+ this.directory =
+ (directory == null || directory.isBlank())
+ ? System.getProperty("java.io.tmpdir")
+ : directory;
+ this.memory = new ArrayDeque<>(Math.min(maxSize, 1024));
+ this.timeoutGet =
+ Const.toInt(System.getProperty(Const.HOP_ROWSET_GET_TIMEOUT),
Const.TIMEOUT_GET_MILLIS);
+ }
+
+ @Override
+ public boolean putRow(IRowMeta rowMeta, Object[] rowData) {
+ return putRowWait(rowMeta, rowData, 0, TimeUnit.MILLISECONDS);
+ }
+
+ @Override
+ public boolean putRowWait(IRowMeta rowMeta, Object[] rowData, long time,
TimeUnit tu) {
+ if (rowMeta == null || rowData == null || spillIoFailed) {
+ return false;
+ }
+ synchronized (lock) {
+ try {
+ this.rowMeta = rowMeta;
+ if (unreadSpilled == 0 && memory.size() < capacity) {
+ memory.addLast(rowData);
+ } else {
+ spillRow(rowMeta, rowData);
+ unreadSpilled++;
+ }
+ lock.notifyAll();
+ return true;
+ } catch (Exception e) {
+ spillIoFailed = true;
+ lock.notifyAll();
+ return false;
+ }
+ }
+ }
+
+ @Override
+ public Object[] getRow() {
+ return getRowWait(timeoutGet, TimeUnit.MILLISECONDS);
+ }
+
+ @Override
+ public Object[] getRowImmediate() {
+ synchronized (lock) {
+ return takeAvailable();
+ }
+ }
+
+ @Override
+ public Object[] getRowWait(long timeout, TimeUnit tu) {
+ long deadlineNanos = System.nanoTime() + tu.toNanos(timeout);
+ synchronized (lock) {
+ while (true) {
+ Object[] row = takeAvailable();
+ if (row != null) {
+ return row;
+ }
+ if (spillIoFailed) {
+ return null;
+ }
+ if (isDone() && memory.isEmpty() && unreadSpilled == 0) {
+ return null;
+ }
+ long remaining = deadlineNanos - System.nanoTime();
+ if (remaining <= 0L) {
+ return null;
+ }
+ try {
+ long ms = remaining / 1_000_000L;
+ int ns = (int) (remaining % 1_000_000L);
+ lock.wait(ms, ns);
+ } catch (InterruptedException e) {
+ Thread.currentThread().interrupt();
+ return null;
+ }
+ }
+ }
+ }
+
+ /**
+ * Reported occupancy for transform flow-control heuristics (not pure
in-memory depth).
+ *
+ * <p>The classic pipeline engine sleeps with {@code Thread.sleep(0, 1)}
when {@code size() >=
+ * 99%} of the rowset capacity (producer) or {@code size() <= 1%}
(consumer). On many JVMs/OSes
+ * that "1 ns" sleep rounds up to about 1 ms. If we reported a full memory
buffer while spilling
+ * (put is non-blocking), every put would sleep ~1 ms — near-zero CPU/disk
and catastrophic
+ * throughput. Likewise, reporting {@code 0} while rows still sit only on
disk would make the
+ * consumer sleep on every get.
+ *
+ * <p>So when there is pending work we advertise a mid-level size (neither
"full" nor "empty").
+ * When idle we report {@code 0}.
+ */
+ @Override
+ public int size() {
+ synchronized (lock) {
+ long pending = memory.size() + unreadSpilled;
+ if (pending <= 0L) {
+ return 0;
+ }
+ // Stay strictly below upperBufferBoundary (0.99 * capacity) and above
lower (0.01 * capacity)
+ // for typical capacities so BaseTransform does not inject per-row
sleeps.
+ int mid = Math.max(1, capacity / 2);
+ if (unreadSpilled == 0 && memory.size() < mid) {
+ return memory.size();
+ }
+ return mid;
+ }
+ }
+
+ @Override
+ public void setDone() {
+ super.setDone();
+ synchronized (lock) {
+ closeWriteSegmentQuietly();
+ lock.notifyAll();
+ }
+ }
+
+ @Override
+ public void clear() {
+ synchronized (lock) {
+ memory.clear();
+ closeWriteSegmentQuietly();
+ closeReadQuietly();
+ deleteAllSegments();
+ segments.clear();
+ writeSegment = null;
+ readSegmentIndex = 0;
+ readRowsInSegment = 0;
+ unreadSpilled = 0;
+ spillIoFailed = false;
+ firstSpillLogged = false;
+ done.set(false);
+ }
+ }
+
+ @Override
+ public boolean isBlocking() {
+ return true;
+ }
+
+ /** For tests: rows still only on disk. */
+ long getUnreadSpilled() {
+ synchronized (lock) {
+ return unreadSpilled;
+ }
+ }
+
+ boolean hasSpilled() {
+ synchronized (lock) {
+ return firstSpillLogged || !segments.isEmpty();
+ }
+ }
+
+ private Object[] takeAvailable() {
+ Object[] row = memory.pollFirst();
+ if (row != null) {
+ return row;
+ }
+ if (unreadSpilled > 0) {
+ try {
+ row = readSpilledRow();
+ if (row != null) {
+ unreadSpilled--;
+ }
+ return row;
+ } catch (Exception e) {
+ spillIoFailed = true;
+ return null;
+ }
+ }
+ return null;
+ }
+
+ private void spillRow(IRowMeta meta, Object[] row) throws HopFileException,
IOException {
+ if (writeSegment == null || writeSegment.closedForWrite) {
+ openNewWriteSegment();
+ }
+ meta.writeData(writeSegment.output, row);
+ writeSegment.rowCount++;
+ firstSpillLogged = true;
+ }
+
+ private void openNewWriteSegment() throws HopFileException, IOException {
+ closeWriteSegmentQuietly();
+ FileObject file = HopVfs.createTempFile("spilling-rowset", ".tmp",
directory);
+ // Belt-and-suspenders if the JVM exits without pipeline cleanup (kill -9
still loses these).
+ try {
+ new File(file.getName().getPath()).deleteOnExit();
+ } catch (Exception e) {
+ // non-local VFS or path mapping issues — pipeline cleanupRowSets still
deletes via VFS
+ }
+ OutputStream os = HopVfs.getOutputStream(file, false);
+ DataOutputStream dos = new DataOutputStream(new BufferedOutputStream(os,
65536));
+ writeSegment = new SpillSegment(file, dos);
+ segments.add(writeSegment);
+ }
+
+ private Object[] readSpilledRow() throws HopFileException, IOException {
+ closeWriteSegmentQuietly();
+
+ while (readSegmentIndex < segments.size()) {
+ SpillSegment segment = segments.get(readSegmentIndex);
+ if (segment.input == null) {
+ InputStream is = HopVfs.getInputStream(segment.file);
+ segment.input = new DataInputStream(new BufferedInputStream(is,
65536));
+ }
+ if (readRowsInSegment < segment.rowCount) {
+ Object[] row = rowMeta.readData(segment.input);
+ readRowsInSegment++;
+ return row;
+ }
+ // Segment exhausted
+ closeSegmentInput(segment);
+ deleteSegmentFile(segment);
+ readSegmentIndex++;
+ readRowsInSegment = 0;
+ }
+ return null;
+ }
+
+ private void closeWriteSegmentQuietly() {
+ if (writeSegment != null && !writeSegment.closedForWrite) {
+ try {
+ writeSegment.output.flush();
+ writeSegment.output.close();
+ } catch (IOException e) {
+ // ignore on close
+ }
+ writeSegment.closedForWrite = true;
+ writeSegment.output = null;
+ }
+ }
+
+ private void closeReadQuietly() {
+ for (SpillSegment segment : segments) {
+ closeSegmentInput(segment);
+ }
+ }
+
+ private static void closeSegmentInput(SpillSegment segment) {
+ if (segment.input != null) {
+ try {
+ segment.input.close();
+ } catch (IOException e) {
+ // ignore
+ }
+ segment.input = null;
+ }
+ }
+
+ private void deleteAllSegments() {
+ for (SpillSegment segment : segments) {
+ deleteSegmentFile(segment);
+ }
+ }
+
+ private static void deleteSegmentFile(SpillSegment segment) {
+ if (segment.file != null) {
+ try {
+ if (segment.file.exists()) {
+ segment.file.delete();
+ }
+ } catch (Exception e) {
+ // best-effort cleanup
+ }
+ segment.file = null;
+ }
+ }
+
+ @Override
+ public int compareTo(@NonNull IRowSet rowSet) {
+ return super.compareTo(rowSet);
+ }
+
+ private static final class SpillSegment {
+ private FileObject file;
+ private DataOutputStream output;
+ private DataInputStream input;
+ private long rowCount;
+ private boolean closedForWrite;
+
+ private SpillSegment(FileObject file, DataOutputStream output) {
+ this.file = file;
+ this.output = output;
+ }
+ }
+}
diff --git a/core/src/test/java/org/apache/hop/core/SpillingRowSetTest.java
b/core/src/test/java/org/apache/hop/core/SpillingRowSetTest.java
new file mode 100644
index 0000000000..e4528d7ead
--- /dev/null
+++ b/core/src/test/java/org/apache/hop/core/SpillingRowSetTest.java
@@ -0,0 +1,121 @@
+/*
+ * 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.core;
+
+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 org.apache.hop.core.row.IRowMeta;
+import org.apache.hop.core.row.RowMeta;
+import org.apache.hop.core.row.value.ValueMetaInteger;
+import org.apache.hop.junit.rules.RestoreHopEnvironmentExtension;
+import org.junit.jupiter.api.Test;
+import org.junit.jupiter.api.extension.ExtendWith;
+
+@ExtendWith(RestoreHopEnvironmentExtension.class)
+class SpillingRowSetTest {
+
+ private static IRowMeta rowMeta() {
+ IRowMeta rm = new RowMeta();
+ rm.addValueMeta(new ValueMetaInteger("ROWNR"));
+ return rm;
+ }
+
+ @Test
+ void memoryOnlyPathBehavesLikeBoundedQueue() {
+ SpillingRowSet set = new SpillingRowSet(3);
+ IRowMeta rm = rowMeta();
+
+ assertTrue(set.putRow(rm, new Object[] {1L}));
+ assertTrue(set.putRow(rm, new Object[] {2L}));
+ assertTrue(set.size() > 0);
+ assertFalse(set.hasSpilled());
+
+ assertEquals(1L, set.getRowImmediate()[0]);
+ assertEquals(2L, set.getRowImmediate()[0]);
+ assertNull(set.getRowImmediate());
+ assertEquals(0, set.size());
+ }
+
+ @Test
+ void spillsWhenFullAndPreservesFifoOrder() {
+ SpillingRowSet set = new SpillingRowSet(2);
+ IRowMeta rm = rowMeta();
+
+ for (long i = 1; i <= 5; i++) {
+ assertTrue(set.putRow(rm, new Object[] {i}), "put " + i);
+ }
+ assertTrue(set.hasSpilled());
+ // size() is a flow-control signal (mid when work pending), not pure
memory count
+ assertTrue(set.size() > 0);
+ assertTrue(set.size() < 2, "must not report full or BaseTransform will
sleep per row");
+ assertEquals(3, set.getUnreadSpilled());
+
+ set.setDone();
+
+ for (long expected = 1; expected <= 5; expected++) {
+ Object[] row = set.getRow();
+ assertNotNull(row, "missing row " + expected);
+ assertEquals(expected, row[0]);
+ }
+ assertNull(set.getRow());
+ }
+
+ @Test
+ void clearDeletesState() {
+ SpillingRowSet set = new SpillingRowSet(1);
+ IRowMeta rm = rowMeta();
+ assertTrue(set.putRow(rm, new Object[] {1L}));
+ assertTrue(set.putRow(rm, new Object[] {2L}));
+ assertTrue(set.hasSpilled());
+
+ set.clear();
+ assertEquals(0, set.size());
+ assertEquals(0, set.getUnreadSpilled());
+ assertFalse(set.isDone());
+ assertNull(set.getRowImmediate());
+ }
+
+ @Test
+ void doneWithEmptyReturnsNull() {
+ SpillingRowSet set = new SpillingRowSet(2);
+ set.setDone();
+ assertNull(set.getRow());
+ }
+
+ @Test
+ void sizeDoesNotAdvertiseFullWhileSpilling() {
+ // BaseTransform uses size() >= 0.99 * capacity to Thread.sleep(0,1)
before put.
+ SpillingRowSet set = new SpillingRowSet(100);
+ IRowMeta rm = rowMeta();
+ for (long i = 0; i < 150; i++) {
+ assertTrue(set.putRow(rm, new Object[] {i}));
+ }
+ assertTrue(set.hasSpilled());
+ int upper = (int) (100 * 0.99);
+ assertTrue(
+ set.size() < upper,
+ "size()=" + set.size() + " must stay below upper boundary " + upper +
" while spilling");
+ // Still not "empty" so consumer low-water sleep is avoided
+ int lower = (int) (100 * 0.01);
+ assertTrue(set.size() > lower, "size() must stay above lower boundary
while work is pending");
+ }
+}
diff --git
a/docs/hop-user-manual/modules/ROOT/pages/how-to-guides/avoiding-deadlocks.adoc
b/docs/hop-user-manual/modules/ROOT/pages/how-to-guides/avoiding-deadlocks.adoc
index 34c32c22a8..493c04fdce 100644
---
a/docs/hop-user-manual/modules/ROOT/pages/how-to-guides/avoiding-deadlocks.adoc
+++
b/docs/hop-user-manual/modules/ROOT/pages/how-to-guides/avoiding-deadlocks.adoc
@@ -78,6 +78,25 @@
image:how-to-guides/deadlocks-stream-lookup/deadlock-stream-lookup-use-blocking-
* Configure the Blocking transform with the `Pass all rows` option to handle
streams in a sequential manner.
* Adjust settings like cache size within the Blocking transform for optimal
performance.
+=== 5. Detect and mitigate on the local engine (automatic)
+
+For the *Local* pipeline engine, Hop can detect split–rejoin topologies that
risk bounded-rowset deadlocks (for example Stream Lookup or Merge Join with a
shared upstream source).
+
+On the local *Pipeline Run Configuration*:
+
+* *Detect buffer deadlocks* (enabled by default) — analyzes the pipeline
during prepare and logs any risks. Pipeline verify also reports them as
warnings.
+* *Mitigate buffer deadlocks (spill to disk)* (enabled by default) — when
risks are found, only the recommended hops use a spilling rowset: rows stay in
memory up to the rowset size, then excess rows spill to temporary files so
producers do not block. Other hops keep normal bounded buffers. Prefer progress
over hanging; turn this off if you want strict bounded memory only.
+* *Buffer deadlock spill directory* — optional temp directory for spill files
(defaults to the system temporary directory).
+
+This avoids redesigning the pipeline for many common cases, at the cost of
disk I/O on the hops that actually fill. Beam and Spark engines are unaffected
(they materialize between stages).
+
+==== How detection and mitigation work
+
+* *Detection* looks for multi-input transforms (main and/or info streams)
whose inbound predecessors share a common ancestor — the classic split–rejoin
shape behind Stream Lookup and Merge Join hangs. Structural hop cycles are
already forbidden; this is a *bounded-buffer wait-for* risk, not a graph cycle.
+* *Mitigation* replaces only the recommended inbound hops into those
transforms with a *spilling rowset*: in-memory up to the configured rowset
size, then excess rows are serialized to temporary files (row *data* only, not
metadata per row). Producers no longer block forever on a full buffer, so the
pipeline can finish and call end-of-stream.
+* Spilling is *not* applied to info hops alone. For Stream Lookup the critical
hop is often the *main* input (the one not drained while the lookup cache
loads); the analyzer therefore targets all risky inbound hops at the
reconvergence.
+* Existing design-time workarounds (separate sources, Blocking transform,
larger rowset size) remain valid. Mitigation is the automatic safety net for
the local multi-threaded engine.
+
=== How the Merge Join Transform Can Cause Deadlocks
Deadlocks can also occur with the
xref:pipeline/transforms/mergejoin.adoc[Merge Join] transform, particularly
when processing large datasets or running pipelines locally. Here’s an example
scenario that demonstrates how deadlocks might arise with the *Merge Join*
transform:
diff --git a/engine/src/main/java/org/apache/hop/pipeline/Pipeline.java
b/engine/src/main/java/org/apache/hop/pipeline/Pipeline.java
index d227687bef..c21372d6d1 100644
--- a/engine/src/main/java/org/apache/hop/pipeline/Pipeline.java
+++ b/engine/src/main/java/org/apache/hop/pipeline/Pipeline.java
@@ -32,6 +32,7 @@ import java.util.Date;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
+import java.util.Set;
import java.util.Timer;
import java.util.TimerTask;
import java.util.UUID;
@@ -55,6 +56,7 @@ import org.apache.hop.core.QueueRowSet;
import org.apache.hop.core.Result;
import org.apache.hop.core.ResultFile;
import org.apache.hop.core.RowMetaAndData;
+import org.apache.hop.core.SpillingRowSet;
import org.apache.hop.core.database.Database;
import org.apache.hop.core.exception.HopException;
import org.apache.hop.core.exception.HopFileException;
@@ -93,6 +95,8 @@ import
org.apache.hop.execution.sampler.IExecutionDataSamplerStore;
import org.apache.hop.i18n.BaseMessages;
import org.apache.hop.metadata.api.IHopMetadataProvider;
import org.apache.hop.partition.PartitionSchema;
+import org.apache.hop.pipeline.analysis.BufferDeadlockRisk.SpillHop;
+import org.apache.hop.pipeline.analysis.PipelineBufferDeadlockAnalyzer;
import org.apache.hop.pipeline.config.IPipelineEngineRunConfiguration;
import org.apache.hop.pipeline.config.PipelineRunConfiguration;
import org.apache.hop.pipeline.engine.EngineCompatibilityChecker;
@@ -375,6 +379,20 @@ public abstract class Pipeline
@Setter protected int feedbackSize;
+ /**
+ * Hops that should use {@link SpillingRowSet} when the local engine
mitigates buffer deadlocks.
+ * Empty means all hops use the normal bounded rowset.
+ */
+ @Getter protected Set<SpillHop> bufferDeadlockSpillHops = Set.of();
+
+ /** Directory for {@link SpillingRowSet} temp files; blank uses the system
temp directory. */
+ @Getter @Setter protected String bufferDeadlockSpillDirectory;
+
+ public void setBufferDeadlockSpillHops(Set<SpillHop>
bufferDeadlockSpillHops) {
+ this.bufferDeadlockSpillHops =
+ bufferDeadlockSpillHops == null ? Set.of() :
Set.copyOf(bufferDeadlockSpillHops);
+ }
+
/** Instantiates a new pipeline. */
public Pipeline() {
@@ -400,6 +418,7 @@ public abstract class Pipeline
extensionDataMap = new HashMap<>();
rowSetSize = Const.ROWS_IN_ROWSET;
+ bufferDeadlockSpillHops = Set.of();
dataSamplers = Collections.synchronizedList(new ArrayList<>());
}
@@ -768,6 +787,9 @@ public abstract class Pipeline
System.getProperty(Const.HOP_BATCHING_ROWSET));
if (batchingRowSet != null && batchingRowSet) {
rowSet = new BlockingBatchingRowSet(rowSetSize);
+ } else if (PipelineBufferDeadlockAnalyzer.shouldSpill(
+ bufferDeadlockSpillHops, thisTransform.getName(),
nextTransform.getName())) {
+ rowSet = new SpillingRowSet(rowSetSize,
bufferDeadlockSpillDirectory);
} else {
rowSet = new BlockingRowSet(rowSetSize);
}
@@ -817,7 +839,13 @@ public abstract class Pipeline
// distribution...
for (int s = 0; s < thisCopies; s++) {
for (int t = 0; t < nextCopies; t++) {
- BlockingRowSet rowSet = new BlockingRowSet(rowSetSize);
+ IRowSet rowSet;
+ if (PipelineBufferDeadlockAnalyzer.shouldSpill(
+ bufferDeadlockSpillHops, thisTransform.getName(),
nextTransform.getName())) {
+ rowSet = new SpillingRowSet(rowSetSize,
bufferDeadlockSpillDirectory);
+ } else {
+ rowSet = new BlockingRowSet(rowSetSize);
+ }
rowSet.setThreadNameFromToCopy(
thisTransform.getName(), s, nextTransform.getName(), t);
rowsets.add(rowSet);
@@ -1211,6 +1239,8 @@ public abstract class Pipeline
? ComponentExecutionStatus.STATUS_STOPPED
: ComponentExecutionStatus.STATUS_HALTED);
}
+ // Rowsets may already exist if prepare allocated them before a later
failure.
+ cleanupRowSets();
}
/**
@@ -1357,6 +1387,10 @@ public abstract class Pipeline
log.snap(Metrics.METRIC_PIPELINE_EXECUTION_STOP);
+ // Drop rowset buffers and any SpillingRowSet temp files (stop
mid-run or normal end).
+ // Safe here: all transform threads have finished before this
listener runs.
+ cleanupRowSets();
+
// release unused vfs connections
HopVfs.freeUnusedResources();
};
@@ -1565,6 +1599,8 @@ public abstract class Pipeline
*/
@Override
public void cleanup() {
+ cleanupRowSets();
+
// Close all open server sockets.
// We can only close these after all processing has been confirmed to be
finished.
//
@@ -1577,6 +1613,29 @@ public abstract class Pipeline
}
}
+ /**
+ * Clear every pipeline rowset. For {@link
org.apache.hop.core.SpillingRowSet} this closes streams
+ * and deletes spill temp files left after stop or unfinished consumption.
Safe to call more than
+ * once; no-op when rowsets were never allocated.
+ */
+ protected void cleanupRowSets() {
+ if (rowsets == null) {
+ return;
+ }
+ for (IRowSet rowSet : rowsets) {
+ if (rowSet == null) {
+ continue;
+ }
+ try {
+ rowSet.clear();
+ } catch (Exception e) {
+ if (log != null) {
+ log.logError("Error clearing rowset " + rowSet, e);
+ }
+ }
+ }
+ }
+
/** Waits until all RunThreads have finished. */
@Override
public void waitUntilFinished() {
diff --git a/engine/src/main/java/org/apache/hop/pipeline/PipelineMeta.java
b/engine/src/main/java/org/apache/hop/pipeline/PipelineMeta.java
index 75717b752c..4546493d80 100644
--- a/engine/src/main/java/org/apache/hop/pipeline/PipelineMeta.java
+++ b/engine/src/main/java/org/apache/hop/pipeline/PipelineMeta.java
@@ -78,6 +78,8 @@ import org.apache.hop.metadata.api.IEnumHasCodeAndDescription;
import org.apache.hop.metadata.api.IHopMetadataProvider;
import org.apache.hop.metadata.serializer.xml.XmlMetadataUtil;
import org.apache.hop.partition.PartitionSchema;
+import org.apache.hop.pipeline.analysis.BufferDeadlockRisk;
+import org.apache.hop.pipeline.analysis.PipelineBufferDeadlockAnalyzer;
import org.apache.hop.pipeline.transform.BaseTransform;
import org.apache.hop.pipeline.transform.ITransformMeta;
import org.apache.hop.pipeline.transform.ITransformMetaChangeListener;
@@ -2719,6 +2721,20 @@ public class PipelineMeta extends AbstractMeta
remarks.add(cr);
}
+ // Bounded-buffer deadlock risk on split–rejoin multi-input transforms
(classic engine)
+ //
+ List<BufferDeadlockRisk> deadlockRisks =
PipelineBufferDeadlockAnalyzer.analyze(this);
+ for (BufferDeadlockRisk risk : deadlockRisks) {
+ remarks.add(
+ new CheckResult(
+ ICheckResult.TYPE_RESULT_WARNING,
+ BaseMessages.getString(
+ PKG,
+
"PipelineMeta.CheckResult.TypeResultWarning.BufferDeadlockRisk.Description",
+ risk.formatMessage()),
+ risk.reconvergence()));
+ }
+
ExtensionPointHandler.callExtensionPoint(
LogChannel.GENERAL,
variables,
diff --git
a/engine/src/main/java/org/apache/hop/pipeline/analysis/BufferDeadlockRisk.java
b/engine/src/main/java/org/apache/hop/pipeline/analysis/BufferDeadlockRisk.java
new file mode 100644
index 0000000000..9ca299e69f
--- /dev/null
+++
b/engine/src/main/java/org/apache/hop/pipeline/analysis/BufferDeadlockRisk.java
@@ -0,0 +1,79 @@
+/*
+ * 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.pipeline.analysis;
+
+import java.util.List;
+import java.util.Objects;
+import java.util.Set;
+import java.util.stream.Collectors;
+import org.apache.hop.pipeline.transform.TransformMeta;
+
+/**
+ * A bounded-buffer deadlock risk found by {@link
PipelineBufferDeadlockAnalyzer}: a multi-input
+ * transform where at least two inbound predecessors share a common ancestor
(split–rejoin).
+ *
+ * @param reconvergence the multi-input transform where streams rejoin
+ * @param commonAncestor a transform that can reach more than one inbound
predecessor
+ * @param inboundPredecessors predecessors of {@code reconvergence} involved
in the risk
+ * @param spillHops hops {@code from → reconvergence} recommended for a
spilling rowset (v1: all
+ * inbound risky hops)
+ */
+public record BufferDeadlockRisk(
+ TransformMeta reconvergence,
+ TransformMeta commonAncestor,
+ List<TransformMeta> inboundPredecessors,
+ Set<SpillHop> spillHops) {
+
+ public BufferDeadlockRisk {
+ Objects.requireNonNull(reconvergence, "reconvergence");
+ Objects.requireNonNull(commonAncestor, "commonAncestor");
+ inboundPredecessors = List.copyOf(inboundPredecessors);
+ spillHops = Set.copyOf(spillHops);
+ }
+
+ /** One pipeline hop identified by transform names (copy-agnostic). */
+ public record SpillHop(String fromTransformName, String toTransformName) {
+ public SpillHop {
+ Objects.requireNonNull(fromTransformName, "fromTransformName");
+ Objects.requireNonNull(toTransformName, "toTransformName");
+ }
+
+ public boolean matches(String from, String to) {
+ return fromTransformName.equalsIgnoreCase(from) &&
toTransformName.equalsIgnoreCase(to);
+ }
+
+ @Override
+ public String toString() {
+ return fromTransformName + " → " + toTransformName;
+ }
+ }
+
+ public String formatMessage() {
+ String preds =
+
inboundPredecessors.stream().map(TransformMeta::getName).collect(Collectors.joining(",
"));
+ String hops =
spillHops.stream().map(SpillHop::toString).collect(Collectors.joining(", "));
+ return "Possible buffer deadlock at transform '"
+ + reconvergence.getName()
+ + "': inputs ["
+ + preds
+ + "] share common ancestor '"
+ + commonAncestor.getName()
+ + "'. Recommended spill hops: "
+ + hops;
+ }
+}
diff --git
a/engine/src/main/java/org/apache/hop/pipeline/analysis/PipelineBufferDeadlockAnalyzer.java
b/engine/src/main/java/org/apache/hop/pipeline/analysis/PipelineBufferDeadlockAnalyzer.java
new file mode 100644
index 0000000000..b65d05c279
--- /dev/null
+++
b/engine/src/main/java/org/apache/hop/pipeline/analysis/PipelineBufferDeadlockAnalyzer.java
@@ -0,0 +1,190 @@
+/*
+ * 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.pipeline.analysis;
+
+import java.util.ArrayDeque;
+import java.util.ArrayList;
+import java.util.Collections;
+import java.util.LinkedHashMap;
+import java.util.LinkedHashSet;
+import java.util.List;
+import java.util.Map;
+import java.util.Set;
+import org.apache.hop.pipeline.PipelineMeta;
+import org.apache.hop.pipeline.analysis.BufferDeadlockRisk.SpillHop;
+import org.apache.hop.pipeline.transform.TransformMeta;
+
+/**
+ * Static analyzer for classic multi-threaded pipeline buffer deadlocks.
+ *
+ * <p>Hop hop graphs are DAGs ({@link PipelineMeta#hasLoop}); the hangs
documented for Stream Lookup
+ * / Merge Join are <em>bounded-buffer wait-for cycles</em> on split–rejoin
topologies. This class
+ * finds multi-input transforms whose inbound predecessors share a common
ancestor and recommends
+ * the minimal set of inbound hops that may need a spilling {@link
org.apache.hop.core.IRowSet}.
+ *
+ * <p>v1 spill recommendation: every hop {@code predecessor → reconvergence}
involved in a risk (not
+ * info-only — the main hop into Stream Lookup must be included).
+ */
+public final class PipelineBufferDeadlockAnalyzer {
+
+ private PipelineBufferDeadlockAnalyzer() {}
+
+ /**
+ * @return risks in pipeline transform order; empty if none or {@code
pipelineMeta} is null
+ */
+ public static List<BufferDeadlockRisk> analyze(PipelineMeta pipelineMeta) {
+ if (pipelineMeta == null) {
+ return Collections.emptyList();
+ }
+
+ List<BufferDeadlockRisk> risks = new ArrayList<>();
+ Map<TransformMeta, Set<TransformMeta>> ancestorCache = new
LinkedHashMap<>();
+
+ for (TransformMeta reconvergence : pipelineMeta.getTransforms()) {
+ if (reconvergence == null) {
+ continue;
+ }
+
+ // Immediate predecessors including info streams
+ List<TransformMeta> preds =
+ new ArrayList<>(pipelineMeta.findPreviousTransforms(reconvergence,
true));
+ // Deduplicate while preserving order
+ preds = new ArrayList<>(new LinkedHashSet<>(preds));
+ if (preds.size() < 2) {
+ continue;
+ }
+
+ // For each pair, find common ancestors; collect all preds that share
ancestry with another
+ Set<TransformMeta> riskyPreds = new LinkedHashSet<>();
+ TransformMeta chosenAncestor = null;
+
+ for (int i = 0; i < preds.size(); i++) {
+ TransformMeta pi = preds.get(i);
+ if (pi == null) {
+ continue;
+ }
+ Set<TransformMeta> ancestorsI = ancestorsIncludingSelf(pipelineMeta,
pi, ancestorCache);
+ for (int j = i + 1; j < preds.size(); j++) {
+ TransformMeta pj = preds.get(j);
+ if (pj == null) {
+ continue;
+ }
+ Set<TransformMeta> ancestorsJ = ancestorsIncludingSelf(pipelineMeta,
pj, ancestorCache);
+ TransformMeta common = firstCommonAncestor(ancestorsI, ancestorsJ);
+ if (common != null) {
+ riskyPreds.add(pi);
+ riskyPreds.add(pj);
+ if (chosenAncestor == null) {
+ chosenAncestor = common;
+ }
+ }
+ }
+ }
+
+ if (riskyPreds.size() < 2 || chosenAncestor == null) {
+ continue;
+ }
+
+ List<TransformMeta> inbound = new ArrayList<>(riskyPreds);
+ Set<SpillHop> spillHops = new LinkedHashSet<>();
+ for (TransformMeta pred : inbound) {
+ spillHops.add(new SpillHop(pred.getName(), reconvergence.getName()));
+ }
+
+ risks.add(new BufferDeadlockRisk(reconvergence, chosenAncestor, inbound,
spillHops));
+ }
+
+ return risks;
+ }
+
+ /**
+ * Union of all recommended spill hops across risks (copy-agnostic from/to
names).
+ *
+ * @param risks analyzer output
+ * @return immutable set of spill hops
+ */
+ public static Set<SpillHop> collectSpillHops(List<BufferDeadlockRisk> risks)
{
+ if (risks == null || risks.isEmpty()) {
+ return Collections.emptySet();
+ }
+ Set<SpillHop> hops = new LinkedHashSet<>();
+ for (BufferDeadlockRisk risk : risks) {
+ hops.addAll(risk.spillHops());
+ }
+ return Collections.unmodifiableSet(hops);
+ }
+
+ /**
+ * Whether the hop from {@code fromTransform} to {@code toTransform} is in
the spill set.
+ *
+ * @param spillHops recommended hops
+ * @param fromTransform source transform name
+ * @param toTransform target transform name
+ * @return true if this hop should use a spilling rowset
+ */
+ public static boolean shouldSpill(
+ Set<SpillHop> spillHops, String fromTransform, String toTransform) {
+ if (spillHops == null || spillHops.isEmpty() || fromTransform == null ||
toTransform == null) {
+ return false;
+ }
+ for (SpillHop hop : spillHops) {
+ if (hop.matches(fromTransform, toTransform)) {
+ return true;
+ }
+ }
+ return false;
+ }
+
+ private static Set<TransformMeta> ancestorsIncludingSelf(
+ PipelineMeta meta, TransformMeta transform, Map<TransformMeta,
Set<TransformMeta>> cache) {
+ Set<TransformMeta> cached = cache.get(transform);
+ if (cached != null) {
+ return cached;
+ }
+
+ Set<TransformMeta> result = new LinkedHashSet<>();
+ ArrayDeque<TransformMeta> stack = new ArrayDeque<>();
+ result.add(transform);
+ stack.push(transform);
+
+ while (!stack.isEmpty()) {
+ TransformMeta current = stack.pop();
+ for (TransformMeta prev : meta.findPreviousTransforms(current, true)) {
+ if (prev != null && result.add(prev)) {
+ stack.push(prev);
+ }
+ }
+ }
+
+ cache.put(transform, result);
+ return result;
+ }
+
+ /**
+ * Picks a common ancestor for messaging. Walks {@code a} in self-first BFS
order and returns the
+ * first member also in {@code b} (nearest shared node on that walk).
+ */
+ private static TransformMeta firstCommonAncestor(Set<TransformMeta> a,
Set<TransformMeta> b) {
+ for (TransformMeta candidate : a) {
+ if (b.contains(candidate)) {
+ return candidate;
+ }
+ }
+ return null;
+ }
+}
diff --git
a/engine/src/main/java/org/apache/hop/pipeline/engines/local/LocalPipelineEngine.java
b/engine/src/main/java/org/apache/hop/pipeline/engines/local/LocalPipelineEngine.java
index 5ce26d1d7e..1ebefcc73b 100644
---
a/engine/src/main/java/org/apache/hop/pipeline/engines/local/LocalPipelineEngine.java
+++
b/engine/src/main/java/org/apache/hop/pipeline/engines/local/LocalPipelineEngine.java
@@ -50,6 +50,8 @@ import
org.apache.hop.execution.sampler.IExecutionDataSamplerStore;
import org.apache.hop.metadata.api.IHopMetadataProvider;
import org.apache.hop.pipeline.Pipeline;
import org.apache.hop.pipeline.PipelineMeta;
+import org.apache.hop.pipeline.analysis.BufferDeadlockRisk;
+import org.apache.hop.pipeline.analysis.PipelineBufferDeadlockAnalyzer;
import org.apache.hop.pipeline.config.IPipelineEngineRunConfiguration;
import org.apache.hop.pipeline.config.PipelineRunConfiguration;
import org.apache.hop.pipeline.engine.IEngineComponent;
@@ -113,6 +115,71 @@ public class LocalPipelineEngine extends Pipeline
implements IPipelineEngine<Pip
false));
}
+ /**
+ * Risks found during configure; logged after {@code
super.prepareExecution()} creates the log.
+ */
+ private List<BufferDeadlockRisk> pendingBufferDeadlockRisks = List.of();
+
+ private boolean pendingBufferDeadlockMitigated;
+
+ /**
+ * Runs the buffer-deadlock analyzer when detection is enabled (default) and
optionally configures
+ * minimal spill hops when mitigation is enabled. Must run before {@code
super.prepareExecution()}
+ * so rowset allocation sees the spill hop set. Logging is deferred until
the pipeline log channel
+ * exists.
+ */
+ private void configureBufferDeadlockHandling(LocalPipelineRunConfiguration
config) {
+ // Reset so nested/reused engines do not keep a previous spill set
+ setBufferDeadlockSpillHops(java.util.Set.of());
+ setBufferDeadlockSpillDirectory(null);
+ pendingBufferDeadlockRisks = List.of();
+ pendingBufferDeadlockMitigated = false;
+
+ boolean detect = config.isDetectBufferDeadlocks();
+ boolean mitigate = config.isMitigateBufferDeadlocks();
+ if (!detect && !mitigate) {
+ return;
+ }
+
+ List<BufferDeadlockRisk> risks =
PipelineBufferDeadlockAnalyzer.analyze(getPipelineMeta());
+ if (detect) {
+ pendingBufferDeadlockRisks = risks;
+ }
+ pendingBufferDeadlockMitigated = mitigate && !risks.isEmpty();
+
+ if (mitigate && !risks.isEmpty()) {
+
setBufferDeadlockSpillHops(PipelineBufferDeadlockAnalyzer.collectSpillHops(risks));
+ String dir = resolve(Const.NVL(config.getBufferDeadlockSpillDirectory(),
""));
+ setBufferDeadlockSpillDirectory(dir);
+ }
+ }
+
+ private void logPendingBufferDeadlockRisks(LocalPipelineRunConfiguration
config) {
+ if (!pendingBufferDeadlockRisks.isEmpty()) {
+ log.logMinimal(
+ "Detected "
+ + pendingBufferDeadlockRisks.size()
+ + " possible buffer deadlock risk(s) (split–rejoin with shared
ancestors).");
+ for (BufferDeadlockRisk risk : pendingBufferDeadlockRisks) {
+ log.logMinimal(risk.formatMessage());
+ }
+ if (!config.isMitigateBufferDeadlocks()) {
+ log.logMinimal(
+ "Mitigation is disabled: this pipeline may hang on large data.
Enable 'Mitigate "
+ + "buffer deadlocks' on the local run configuration to spill
recommended hops to "
+ + "disk automatically.");
+ }
+ }
+ if (pendingBufferDeadlockMitigated) {
+ log.logMinimal(
+ "Buffer deadlock mitigation: using spilling rowsets for "
+ + getBufferDeadlockSpillHops().size()
+ + " hop(s).");
+ }
+ pendingBufferDeadlockRisks = List.of();
+ pendingBufferDeadlockMitigated = false;
+ }
+
@Override
public void prepareExecution() throws HopException {
@@ -131,6 +198,10 @@ public class LocalPipelineEngine extends Pipeline
implements IPipelineEngine<Pip
setFeedbackShown(config.isFeedbackShown());
setFeedbackSize(Const.toInt(resolve(config.getFeedbackSize()),
Const.ROWS_UPDATE));
+ // Buffer deadlock detection / optional minimal spill (classic local
engine only)
+ //
+ configureBufferDeadlockHandling(config);
+
// See if we need to enable transactions...
//
IExtensionData parentExtensionData = getParentPipeline();
@@ -235,6 +306,12 @@ public class LocalPipelineEngine extends Pipeline
implements IPipelineEngine<Pip
super.prepareExecution();
+ // Log buffer-deadlock analysis now that Pipeline.prepareExecution created
the log channel
+ if (pipelineRunConfiguration.getEngineRunConfiguration()
+ instanceof LocalPipelineRunConfiguration localConfig) {
+ logPendingBufferDeadlockRisks(localConfig);
+ }
+
// The transforms are initialized at this point, which means they can hold
on to database
// connections, files, ... If anything below fails we never get to start
the transform threads,
// so nothing would ever dispose them. Clean up explicitly to avoid
leaking those resources.
diff --git
a/engine/src/main/java/org/apache/hop/pipeline/engines/local/LocalPipelineRunConfiguration.java
b/engine/src/main/java/org/apache/hop/pipeline/engines/local/LocalPipelineRunConfiguration.java
index 7b1ef1b413..e8e1395878 100644
---
a/engine/src/main/java/org/apache/hop/pipeline/engines/local/LocalPipelineRunConfiguration.java
+++
b/engine/src/main/java/org/apache/hop/pipeline/engines/local/LocalPipelineRunConfiguration.java
@@ -148,6 +148,52 @@ public class LocalPipelineRunConfiguration extends
EmptyPipelineRunConfiguration
@HopMetadataProperty(key = "transactional")
protected boolean transactional;
+ /**
+ * When enabled, analyze the pipeline for split–rejoin buffer deadlock risk
and log findings
+ * during prepare. Enabled by default.
+ */
+ @GuiWidgetElement(
+ id = "detectBufferDeadlocks",
+ order = "110",
+ parentId = PipelineRunConfiguration.GUI_PLUGIN_ELEMENT_PARENT_ID,
+ type = GuiElementType.CHECKBOX,
+ label =
+
"i18n:org.apache.hop.ui.pipeline.config:PipelineRunConfigurationDialog.DetectBufferDeadlocks.Label",
+ toolTip =
+
"i18n:org.apache.hop.ui.pipeline.config:PipelineRunConfigurationDialog.DetectBufferDeadlocks.ToolTip")
+ @HopMetadataProperty(key = "detect_buffer_deadlocks")
+ protected boolean detectBufferDeadlocks;
+
+ /**
+ * When enabled, allocate spilling rowsets for hops recommended by the
buffer-deadlock analyzer.
+ * On by default so at-risk topologies make progress (spill to disk) instead
of hanging; detection
+ * logging already surfaces when this engages.
+ */
+ @GuiWidgetElement(
+ id = "mitigateBufferDeadlocks",
+ order = "120",
+ parentId = PipelineRunConfiguration.GUI_PLUGIN_ELEMENT_PARENT_ID,
+ type = GuiElementType.CHECKBOX,
+ label =
+
"i18n:org.apache.hop.ui.pipeline.config:PipelineRunConfigurationDialog.MitigateBufferDeadlocks.Label",
+ toolTip =
+
"i18n:org.apache.hop.ui.pipeline.config:PipelineRunConfigurationDialog.MitigateBufferDeadlocks.ToolTip")
+ @HopMetadataProperty(key = "mitigate_buffer_deadlocks")
+ protected boolean mitigateBufferDeadlocks;
+
+ /** Directory for spilled rowset temp files; empty uses the system temporary
directory. */
+ @GuiWidgetElement(
+ id = "bufferDeadlockSpillDirectory",
+ order = "130",
+ parentId = PipelineRunConfiguration.GUI_PLUGIN_ELEMENT_PARENT_ID,
+ type = GuiElementType.TEXT,
+ label =
+
"i18n:org.apache.hop.ui.pipeline.config:PipelineRunConfigurationDialog.BufferDeadlockSpillDirectory.Label",
+ toolTip =
+
"i18n:org.apache.hop.ui.pipeline.config:PipelineRunConfigurationDialog.BufferDeadlockSpillDirectory.ToolTip")
+ @HopMetadataProperty(key = "buffer_deadlock_spill_directory")
+ protected String bufferDeadlockSpillDirectory;
+
@SuppressWarnings("java:S115")
public enum SampleType {
None,
@@ -165,6 +211,9 @@ public class LocalPipelineRunConfiguration extends
EmptyPipelineRunConfiguration
this.sampleTypeInGui = SampleType.Last.name();
this.sampleSize = "100";
this.transactional = false;
+ this.detectBufferDeadlocks = true;
+ this.mitigateBufferDeadlocks = true;
+ this.bufferDeadlockSpillDirectory = "";
}
public LocalPipelineRunConfiguration(LocalPipelineRunConfiguration config) {
@@ -179,6 +228,9 @@ public class LocalPipelineRunConfiguration extends
EmptyPipelineRunConfiguration
this.sampleTypeInGui = config.sampleTypeInGui;
this.sampleSize = config.sampleSize;
this.transactional = config.transactional;
+ this.detectBufferDeadlocks = config.detectBufferDeadlocks;
+ this.mitigateBufferDeadlocks = config.mitigateBufferDeadlocks;
+ this.bufferDeadlockSpillDirectory = config.bufferDeadlockSpillDirectory;
}
@Override
diff --git
a/engine/src/main/resources/org/apache/hop/pipeline/messages/messages_en_US.properties
b/engine/src/main/resources/org/apache/hop/pipeline/messages/messages_en_US.properties
index c8979ed1d1..650c1c291f 100644
---
a/engine/src/main/resources/org/apache/hop/pipeline/messages/messages_en_US.properties
+++
b/engine/src/main/resources/org/apache/hop/pipeline/messages/messages_en_US.properties
@@ -82,6 +82,7 @@ PipelineMeta.CheckResult.TypeResultWarning.Description=Field
[{0}] \: {1} in tra
PipelineMeta.CheckResult.TypeResultWarning.DeprecatedTransformPlugin.Description=This
transform is deprecated.
PipelineMeta.CheckResult.TypeResultWarning.HaveTheSameNameField.Description=I
found input fields that have the same name [{0}]
PipelineMeta.CheckResult.TypeResultWarning.TransformIsNotUsed.Description=This
transform is not used in the pipeline.
+PipelineMeta.CheckResult.TypeResultWarning.BufferDeadlockRisk.Description={0}
PipelineMeta.Exception.ErrorOfSortingTransforms=Exception sorting transforms\:
PipelineMeta.Exception.ErrorOpeningOrValidatingTheXMLFile=Error
opening/validating the XML file ''{0}''\!
PipelineMeta.Exception.ErrorReadingPipeline=Error reading object from XML file
diff --git
a/engine/src/test/java/org/apache/hop/pipeline/analysis/PipelineBufferDeadlockAnalyzerTest.java
b/engine/src/test/java/org/apache/hop/pipeline/analysis/PipelineBufferDeadlockAnalyzerTest.java
new file mode 100644
index 0000000000..2ed425aa54
--- /dev/null
+++
b/engine/src/test/java/org/apache/hop/pipeline/analysis/PipelineBufferDeadlockAnalyzerTest.java
@@ -0,0 +1,136 @@
+/*
+ * 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.pipeline.analysis;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertFalse;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+
+import java.util.List;
+import java.util.Set;
+import org.apache.hop.junit.rules.RestoreHopEngineEnvironmentExtension;
+import org.apache.hop.pipeline.PipelineHopMeta;
+import org.apache.hop.pipeline.PipelineMeta;
+import org.apache.hop.pipeline.analysis.BufferDeadlockRisk.SpillHop;
+import org.apache.hop.pipeline.transform.TransformMeta;
+import org.apache.hop.pipeline.transforms.dummy.DummyMeta;
+import org.junit.jupiter.api.Test;
+import org.junit.jupiter.api.extension.ExtendWith;
+
+@ExtendWith(RestoreHopEngineEnvironmentExtension.class)
+class PipelineBufferDeadlockAnalyzerTest {
+
+ @Test
+ void linearChainHasNoRisk() {
+ PipelineMeta meta = new PipelineMeta();
+ TransformMeta a = transform("A");
+ TransformMeta b = transform("B");
+ TransformMeta c = transform("C");
+ meta.addTransform(a);
+ meta.addTransform(b);
+ meta.addTransform(c);
+ meta.addPipelineHop(new PipelineHopMeta(a, b));
+ meta.addPipelineHop(new PipelineHopMeta(b, c));
+
+ assertTrue(PipelineBufferDeadlockAnalyzer.analyze(meta).isEmpty());
+ }
+
+ @Test
+ void splitRejoinReportsRiskAndSpillHops() {
+ // Source ──┬──► Left ──► Join
+ // └──► Right ──► Join
+ PipelineMeta meta = new PipelineMeta();
+ TransformMeta source = transform("Source");
+ TransformMeta left = transform("Left");
+ TransformMeta right = transform("Right");
+ TransformMeta join = transform("Join");
+ meta.addTransform(source);
+ meta.addTransform(left);
+ meta.addTransform(right);
+ meta.addTransform(join);
+ meta.addPipelineHop(new PipelineHopMeta(source, left));
+ meta.addPipelineHop(new PipelineHopMeta(source, right));
+ meta.addPipelineHop(new PipelineHopMeta(left, join));
+ meta.addPipelineHop(new PipelineHopMeta(right, join));
+
+ List<BufferDeadlockRisk> risks =
PipelineBufferDeadlockAnalyzer.analyze(meta);
+ assertEquals(1, risks.size());
+ BufferDeadlockRisk risk = risks.getFirst();
+ assertEquals("Join", risk.reconvergence().getName());
+ assertEquals("Source", risk.commonAncestor().getName());
+ assertEquals(2, risk.inboundPredecessors().size());
+ assertEquals(2, risk.spillHops().size());
+ assertTrue(PipelineBufferDeadlockAnalyzer.shouldSpill(risk.spillHops(),
"Left", "Join"));
+ assertTrue(PipelineBufferDeadlockAnalyzer.shouldSpill(risk.spillHops(),
"Right", "Join"));
+ assertFalse(PipelineBufferDeadlockAnalyzer.shouldSpill(risk.spillHops(),
"Source", "Left"));
+ }
+
+ @Test
+ void independentSourcesNoRisk() {
+ // Src1 → Left ──► Join
+ // Src2 → Right ──► Join
+ PipelineMeta meta = new PipelineMeta();
+ TransformMeta src1 = transform("Src1");
+ TransformMeta src2 = transform("Src2");
+ TransformMeta left = transform("Left");
+ TransformMeta right = transform("Right");
+ TransformMeta join = transform("Join");
+ meta.addTransform(src1);
+ meta.addTransform(src2);
+ meta.addTransform(left);
+ meta.addTransform(right);
+ meta.addTransform(join);
+ meta.addPipelineHop(new PipelineHopMeta(src1, left));
+ meta.addPipelineHop(new PipelineHopMeta(src2, right));
+ meta.addPipelineHop(new PipelineHopMeta(left, join));
+ meta.addPipelineHop(new PipelineHopMeta(right, join));
+
+ assertTrue(PipelineBufferDeadlockAnalyzer.analyze(meta).isEmpty());
+ }
+
+ @Test
+ void collectSpillHopsUnionsRisks() {
+ PipelineMeta meta = new PipelineMeta();
+ TransformMeta source = transform("Source");
+ TransformMeta left = transform("Left");
+ TransformMeta right = transform("Right");
+ TransformMeta join = transform("Join");
+ meta.addTransform(source);
+ meta.addTransform(left);
+ meta.addTransform(right);
+ meta.addTransform(join);
+ meta.addPipelineHop(new PipelineHopMeta(source, left));
+ meta.addPipelineHop(new PipelineHopMeta(source, right));
+ meta.addPipelineHop(new PipelineHopMeta(left, join));
+ meta.addPipelineHop(new PipelineHopMeta(right, join));
+
+ Set<SpillHop> hops =
+ PipelineBufferDeadlockAnalyzer.collectSpillHops(
+ PipelineBufferDeadlockAnalyzer.analyze(meta));
+ assertEquals(2, hops.size());
+ }
+
+ @Test
+ void nullMetaReturnsEmpty() {
+ assertTrue(PipelineBufferDeadlockAnalyzer.analyze(null).isEmpty());
+ }
+
+ private static TransformMeta transform(String name) {
+ return new TransformMeta(name, new DummyMeta());
+ }
+}
diff --git
a/integration-tests/transforms/0022-stream-lookup-buffer-deadlock.hpl
b/integration-tests/transforms/0022-stream-lookup-buffer-deadlock.hpl
new file mode 100644
index 0000000000..0a49fed7ab
--- /dev/null
+++ b/integration-tests/transforms/0022-stream-lookup-buffer-deadlock.hpl
@@ -0,0 +1,331 @@
+<?xml version="1.0" encoding="UTF-8"?>
+<!--
+
+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.
+
+-->
+<pipeline>
+ <info>
+ <name>0022-stream-lookup-buffer-deadlock</name>
+ <name_sync_with_filename>Y</name_sync_with_filename>
+ <description>Split-rejoin Stream Lookup that exceeds the local rowset size
(10000). Without buffer-deadlock mitigation the classic engine hangs; with
mitigation the pipeline completes.</description>
+ <extended_description/>
+ <pipeline_version/>
+ <pipeline_type>Normal</pipeline_type>
+ <parameters/>
+ <capture_transform_performance>N</capture_transform_performance>
+
<transform_performance_capturing_delay>1000</transform_performance_capturing_delay>
+
<transform_performance_capturing_size_limit>100</transform_performance_capturing_size_limit>
+ <created_user>-</created_user>
+ <created_date>2026/08/08 20:25:03.550</created_date>
+ <modified_user>-</modified_user>
+ <modified_date>2026/08/08 20:25:03.550</modified_date>
+ </info>
+ <notepads/>
+ <order>
+ <hop>
+ <from>A</from>
+ <to>Collect</to>
+ <enabled>Y</enabled>
+ </hop>
+ <hop>
+ <from>B</from>
+ <to>Collect</to>
+ <enabled>Y</enabled>
+ </hop>
+ <hop>
+ <from>Collect</from>
+ <to>group-value</to>
+ <enabled>Y</enabled>
+ </hop>
+ <hop>
+ <from>Collect</from>
+ <to>Stream lookup</to>
+ <enabled>Y</enabled>
+ </hop>
+ <hop>
+ <from>group-value</from>
+ <to>Stream lookup</to>
+ <enabled>Y</enabled>
+ </hop>
+ <hop>
+ <from>Stream lookup</from>
+ <to>results</to>
+ <enabled>Y</enabled>
+ </hop>
+ <hop>
+ <from>results</from>
+ <to>Summarize</to>
+ <enabled>Y</enabled>
+ </hop>
+ <hop>
+ <from>Summarize</from>
+ <to>Verify</to>
+ <enabled>Y</enabled>
+ </hop>
+ </order>
+ <transform>
+ <name>A</name>
+ <type>RowGenerator</type>
+ <description/>
+ <distribute>Y</distribute>
+ <custom_distribution/>
+ <copies>1</copies>
+ <partitioning>
+ <method>none</method>
+ <schema_name/>
+ </partitioning>
+ <never_ending>N</never_ending>
+ <interval_in_ms>5000</interval_in_ms>
+ <row_time_field>now</row_time_field>
+ <last_time_field>FiveSecondsAgo</last_time_field>
+ <limit>15000</limit>
+ <fields>
+ <field>
+ <name>group</name>
+ <type>String</type>
+ <format/>
+ <length>-1</length>
+ <precision>-1</precision>
+ <currency/>
+ <decimal/>
+ <group/>
+ <nullif>A</nullif>
+ <set_empty_string>N</set_empty_string>
+ </field>
+ <field>
+ <name>value</name>
+ <type>String</type>
+ <format/>
+ <length>-1</length>
+ <precision>-1</precision>
+ <currency/>
+ <decimal/>
+ <group/>
+ <nullif>aaaaa</nullif>
+ <set_empty_string>N</set_empty_string>
+ </field>
+ </fields>
+ <attributes/>
+ <GUI>
+ <xloc>144</xloc>
+ <yloc>144</yloc>
+ </GUI>
+ </transform>
+ <transform>
+ <name>B</name>
+ <type>RowGenerator</type>
+ <description/>
+ <distribute>Y</distribute>
+ <custom_distribution/>
+ <copies>1</copies>
+ <partitioning>
+ <method>none</method>
+ <schema_name/>
+ </partitioning>
+ <never_ending>N</never_ending>
+ <interval_in_ms>5000</interval_in_ms>
+ <row_time_field>now</row_time_field>
+ <last_time_field>FiveSecondsAgo</last_time_field>
+ <limit>15000</limit>
+ <fields>
+ <field>
+ <name>group</name>
+ <type>String</type>
+ <format/>
+ <length>-1</length>
+ <precision>-1</precision>
+ <currency/>
+ <decimal/>
+ <group/>
+ <nullif>B</nullif>
+ <set_empty_string>N</set_empty_string>
+ </field>
+ <field>
+ <name>value</name>
+ <type>String</type>
+ <format/>
+ <length>-1</length>
+ <precision>-1</precision>
+ <currency/>
+ <decimal/>
+ <group/>
+ <nullif>bbbbb</nullif>
+ <set_empty_string>N</set_empty_string>
+ </field>
+ </fields>
+ <attributes/>
+ <GUI>
+ <xloc>144</xloc>
+ <yloc>64</yloc>
+ </GUI>
+ </transform>
+ <transform>
+ <name>Collect</name>
+ <type>Dummy</type>
+ <description>Copies rows to main Stream Lookup hop and to the info path
(group-value).</description>
+ <distribute>N</distribute>
+ <custom_distribution/>
+ <copies>1</copies>
+ <partitioning>
+ <method>none</method>
+ <schema_name/>
+ </partitioning>
+ <attributes/>
+ <GUI>
+ <xloc>288</xloc>
+ <yloc>144</yloc>
+ </GUI>
+ </transform>
+ <transform>
+ <name>group-value</name>
+ <type>MemoryGroupBy</type>
+ <description>Info stream for Stream Lookup (emits after all Collect
rows).</description>
+ <distribute>Y</distribute>
+ <custom_distribution/>
+ <copies>1</copies>
+ <partitioning>
+ <method>none</method>
+ <schema_name/>
+ </partitioning>
+ <group>
+ <field>
+ <name>group</name>
+ </field>
+ </group>
+ <fields>
+ <field>
+ <aggregate>value</aggregate>
+ <subject>value</subject>
+ <type>FIRST</type>
+ <valuefield/>
+ </field>
+ </fields>
+ <give_back_row>N</give_back_row>
+ <attributes/>
+ <GUI>
+ <xloc>416</xloc>
+ <yloc>64</yloc>
+ </GUI>
+ </transform>
+ <transform>
+ <name>Stream lookup</name>
+ <type>StreamLookup</type>
+ <description/>
+ <distribute>Y</distribute>
+ <custom_distribution/>
+ <copies>1</copies>
+ <partitioning>
+ <method>none</method>
+ <schema_name/>
+ </partitioning>
+ <from>group-value</from>
+ <input_sorted>N</input_sorted>
+ <preserve_memory>Y</preserve_memory>
+ <sorted_list>N</sorted_list>
+ <integer_pair>N</integer_pair>
+ <lookup>
+ <key>
+ <name>group</name>
+ <field>group</field>
+ </key>
+ <value>
+ <name>value</name>
+ <rename>value2</rename>
+ <default/>
+ <type>String</type>
+ </value>
+ </lookup>
+ <attributes/>
+ <GUI>
+ <xloc>560</xloc>
+ <yloc>144</yloc>
+ </GUI>
+ </transform>
+ <transform>
+ <name>results</name>
+ <type>Dummy</type>
+ <description/>
+ <distribute>Y</distribute>
+ <custom_distribution/>
+ <copies>1</copies>
+ <partitioning>
+ <method>none</method>
+ <schema_name/>
+ </partitioning>
+ <attributes/>
+ <GUI>
+ <xloc>704</xloc>
+ <yloc>144</yloc>
+ </GUI>
+ </transform>
+ <transform>
+ <name>Summarize</name>
+ <type>MemoryGroupBy</type>
+ <description>Summarize 30000 main rows into two rows (golden is attached
to Verify Dummy, which replaces the marked transform).</description>
+ <distribute>Y</distribute>
+ <custom_distribution/>
+ <copies>1</copies>
+ <partitioning>
+ <method>none</method>
+ <schema_name/>
+ </partitioning>
+ <group>
+ <field>
+ <name>group</name>
+ </field>
+ </group>
+ <fields>
+ <field>
+ <aggregate>row_count</aggregate>
+ <subject>group</subject>
+ <type>COUNT_ALL</type>
+ <valuefield/>
+ </field>
+ <field>
+ <aggregate>value2</aggregate>
+ <subject>value2</subject>
+ <type>FIRST</type>
+ <valuefield/>
+ </field>
+ </fields>
+ <give_back_row>N</give_back_row>
+ <attributes/>
+ <GUI>
+ <xloc>848</xloc>
+ <yloc>144</yloc>
+ </GUI>
+ </transform>
+ <transform>
+ <name>Verify</name>
+ <type>Dummy</type>
+ <description>Golden data set is attached here (unit tests replace this
with Dummy and capture these rows).</description>
+ <distribute>Y</distribute>
+ <custom_distribution/>
+ <copies>1</copies>
+ <partitioning>
+ <method>none</method>
+ <schema_name/>
+ </partitioning>
+ <attributes/>
+ <GUI>
+ <xloc>992</xloc>
+ <yloc>144</yloc>
+ </GUI>
+ </transform>
+ <attributes/>
+ <transform_error_handling/>
+</pipeline>
diff --git
a/integration-tests/transforms/datasets/golden-stream-lookup-buffer-deadlock.csv
b/integration-tests/transforms/datasets/golden-stream-lookup-buffer-deadlock.csv
new file mode 100644
index 0000000000..36a04cbf6b
--- /dev/null
+++
b/integration-tests/transforms/datasets/golden-stream-lookup-buffer-deadlock.csv
@@ -0,0 +1,3 @@
+group,row_count,value2
+A,15000,aaaaa
+B,15000,bbbbb
diff --git a/integration-tests/transforms/main-0022-stream-lookup.hwf
b/integration-tests/transforms/main-0022-stream-lookup.hwf
index 4a159fe323..f6033b97a2 100644
--- a/integration-tests/transforms/main-0022-stream-lookup.hwf
+++ b/integration-tests/transforms/main-0022-stream-lookup.hwf
@@ -61,6 +61,9 @@ limitations under the License.
<test_name>
<name>0022-stream-lookup-default-variables UNIT</name>
</test_name>
+ <test_name>
+ <name>0022-stream-lookup-buffer-deadlock UNIT</name>
+ </test_name>
</test_names>
<parallel>N</parallel>
<xloc>288</xloc>
diff --git
a/integration-tests/transforms/metadata/dataset/golden-stream-lookup-buffer-deadlock.json
b/integration-tests/transforms/metadata/dataset/golden-stream-lookup-buffer-deadlock.json
new file mode 100644
index 0000000000..597e22a032
--- /dev/null
+++
b/integration-tests/transforms/metadata/dataset/golden-stream-lookup-buffer-deadlock.json
@@ -0,0 +1,32 @@
+{
+ "base_filename": "golden-stream-lookup-buffer-deadlock.csv",
+ "name": "golden-stream-lookup-buffer-deadlock",
+ "description": "Summary of Stream Lookup after split-rejoin with row volume
above local rowset size",
+ "dataset_fields": [
+ {
+ "field_comment": "",
+ "field_length": -1,
+ "field_type": 2,
+ "field_precision": -1,
+ "field_format": "",
+ "field_name": "group"
+ },
+ {
+ "field_comment": "",
+ "field_length": -1,
+ "field_type": 5,
+ "field_precision": 0,
+ "field_format": "####0;-####0",
+ "field_name": "row_count"
+ },
+ {
+ "field_comment": "",
+ "field_length": -1,
+ "field_type": 2,
+ "field_precision": -1,
+ "field_format": "",
+ "field_name": "value2"
+ }
+ ],
+ "folder_name": ""
+}
diff --git
a/integration-tests/transforms/metadata/pipeline-run-configuration/local.json
b/integration-tests/transforms/metadata/pipeline-run-configuration/local.json
index 9d16ecd8de..3c0b770bdd 100644
---
a/integration-tests/transforms/metadata/pipeline-run-configuration/local.json
+++
b/integration-tests/transforms/metadata/pipeline-run-configuration/local.json
@@ -10,7 +10,10 @@
"show_feedback": false,
"topo_sort": false,
"gather_metrics": false,
- "transactional": false
+ "transactional": false,
+ "detect_buffer_deadlocks": true,
+ "mitigate_buffer_deadlocks": true,
+ "buffer_deadlock_spill_directory": ""
}
},
"defaultSelection": true,
diff --git
a/integration-tests/transforms/metadata/unit-test/0022-stream-lookup-buffer-deadlock
UNIT.json
b/integration-tests/transforms/metadata/unit-test/0022-stream-lookup-buffer-deadlock
UNIT.json
new file mode 100644
index 0000000000..9fa6253d7e
--- /dev/null
+++
b/integration-tests/transforms/metadata/unit-test/0022-stream-lookup-buffer-deadlock
UNIT.json
@@ -0,0 +1,36 @@
+{
+ "variableValues": [],
+ "database_replacements": [],
+ "autoOpening": true,
+ "basePath": "",
+ "golden_data_sets": [
+ {
+ "field_mappings": [
+ {
+ "transform_field": "group",
+ "data_set_field": "group"
+ },
+ {
+ "transform_field": "row_count",
+ "data_set_field": "row_count"
+ },
+ {
+ "transform_field": "value2",
+ "data_set_field": "value2"
+ }
+ ],
+ "field_order": [
+ "group"
+ ],
+ "transform_name": "Verify",
+ "data_set_name": "golden-stream-lookup-buffer-deadlock"
+ }
+ ],
+ "input_data_sets": [],
+ "name": "0022-stream-lookup-buffer-deadlock UNIT",
+ "description": "Stream Lookup split-rejoin with more rows than local rowset
size; requires buffer-deadlock mitigation to complete",
+ "trans_test_tweaks": [],
+ "persist_filename": "",
+ "pipeline_filename": "./0022-stream-lookup-buffer-deadlock.hpl",
+ "test_type": "UNIT_TEST"
+}
diff --git
a/ui/src/main/resources/org/apache/hop/ui/pipeline/config/messages/messages_en_US.properties
b/ui/src/main/resources/org/apache/hop/ui/pipeline/config/messages/messages_en_US.properties
index dd7b6eebb3..f6177a0635 100644
---
a/ui/src/main/resources/org/apache/hop/ui/pipeline/config/messages/messages_en_US.properties
+++
b/ui/src/main/resources/org/apache/hop/ui/pipeline/config/messages/messages_en_US.properties
@@ -49,3 +49,9 @@ PipelineRunConfigurationDialog.Variables.Column.Name=Variable
name
PipelineRunConfigurationDialog.Variables.Column.Value=Value
PipelineRunConfigurationDialog.VariablesTab.TabTitle=Variables
PipelineRunConfigurationDialog.WaitTime.Label=Wait time for buffer check (ms)
+PipelineRunConfigurationDialog.DetectBufferDeadlocks.Label=Detect buffer
deadlocks
+PipelineRunConfigurationDialog.DetectBufferDeadlocks.ToolTip=Analyze the
pipeline for split-rejoin topologies that can deadlock on bounded rowsets
(classic local engine). Findings are logged during prepare. Enabled by default.
+PipelineRunConfigurationDialog.MitigateBufferDeadlocks.Label=Mitigate buffer
deadlocks (spill to disk)
+PipelineRunConfigurationDialog.MitigateBufferDeadlocks.ToolTip=When risks are
detected, use spilling rowsets only on the recommended hops so producers do not
block when a buffer fills. Enabled by default so pipelines make progress
instead of hanging; excess rows spill to disk (I/O cost). Disable if you prefer
strict bounded memory only.
+PipelineRunConfigurationDialog.BufferDeadlockSpillDirectory.Label=Buffer
deadlock spill directory
+PipelineRunConfigurationDialog.BufferDeadlockSpillDirectory.ToolTip=Directory
for spilled rowset temp files. Leave empty to use the system temporary
directory.