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.

Reply via email to