This is an automated email from the ASF dual-hosted git repository.

mattcasters 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 f7db648264 engine type should not be driven by Metadata, fixes #8262 
(#8264)
f7db648264 is described below

commit f7db648264629233a30a72ad3f1665ab2e00185d
Author: Hans Van Akelyen <[email protected]>
AuthorDate: Fri Sep 4 15:08:31 2026 +0200

    engine type should not be driven by Metadata, fixes #8262 (#8264)
---
 .../modules/ROOT/pages/hop-server/rest-api.adoc    |   1 -
 .../java/org/apache/hop/pipeline/Pipeline.java     |  25 ++--
 .../java/org/apache/hop/pipeline/PipelineMeta.java |  28 ++---
 .../org/apache/hop/pipeline/PipelineMetaInfo.java  |   5 -
 .../hop/pipeline/engine/IPipelineEngine.java       |  16 +++
 .../pipeline/PipelineTypeIsEngineDrivenTest.java   | 131 +++++++++++++++++++++
 .../core/transform/TransformBatchTransform.java    |   2 +-
 .../hop/beam/core/transform/TransformFn.java       |   2 +-
 .../SingleThreadedPipelineEngine.java              |   9 +-
 .../apache/hop/spark/core/HopMapPartitionsFn.java  |   2 +-
 .../spark/core/SparkParallelFileContextTest.java   |   2 +-
 .../transforms/eventhubs/listen/AzureListener.java |   2 +-
 .../kafka/consumer/KafkaConsumerInput.java         |   2 +-
 .../pipeline/transforms/mapping/SimpleMapping.java |  31 +++--
 .../streamschemamerge/TestUtilities.java           |   1 -
 15 files changed, 200 insertions(+), 59 deletions(-)

diff --git a/docs/hop-user-manual/modules/ROOT/pages/hop-server/rest-api.adoc 
b/docs/hop-user-manual/modules/ROOT/pages/hop-server/rest-api.adoc
index 666a2dc6e3..13683f0e1a 100644
--- a/docs/hop-user-manual/modules/ROOT/pages/hop-server/rest-api.adoc
+++ b/docs/hop-user-manual/modules/ROOT/pages/hop-server/rest-api.adoc
@@ -397,7 +397,6 @@ with XML payload (example):
     <description/>
     <extended_description/>
     <pipeline_version/>
-    <pipeline_type>Normal</pipeline_type>
     <parameters>
     </parameters>
     <capture_transform_performance>N</capture_transform_performance>
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 609e6a16f6..6e5e6ba7e7 100644
--- a/engine/src/main/java/org/apache/hop/pipeline/Pipeline.java
+++ b/engine/src/main/java/org/apache/hop/pipeline/Pipeline.java
@@ -195,6 +195,19 @@ public abstract class Pipeline
   /** The pipeline metadata to execute. */
   protected PipelineMeta pipelineMeta;
 
+  /**
+   * The way this engine drives the transforms. {@link 
PipelineMeta.PipelineType#Normal} gives every
+   * transform its own thread and blocking row sets; {@link
+   * PipelineMeta.PipelineType#SingleThreaded} leaves the transforms to be 
driven one iteration at a
+   * time by a {@link SingleThreadedPipelineExecutor}.
+   *
+   * <p>This belongs to the engine, never to the pipeline metadata: the same 
pipeline can be run
+   * either way and nothing about the run may be written back into the 
design-time metadata.
+   * Transforms that embed a sub-pipeline (Simple Mapping, Kafka Consumer, the 
Beam and Spark
+   * workers) set this on the child engine they create.
+   */
+  @Getter @Setter private PipelineMeta.PipelineType pipelineType = 
PipelineMeta.PipelineType.Normal;
+
   /** The MetaStore to use */
   protected IHopMetadataProvider metadataProvider;
 
@@ -842,7 +855,7 @@ public abstract class Pipeline
         if (dispatchType != TYPE_DISP_N_M) {
           for (int c = 0; c < nrCopies; c++) {
             IRowSet rowSet;
-            switch (pipelineMeta.getPipelineType()) {
+            switch (getPipelineType()) {
               case Normal:
                 // This is a temporary patch until the batching rowset has 
proven
                 // to be working in all situations.
@@ -867,8 +880,7 @@ public abstract class Pipeline
                 break;
 
               default:
-                throw new HopException(
-                    "Unhandled pipeline type: " + 
pipelineMeta.getPipelineType());
+                throw new HopException("Unhandled pipeline type: " + 
getPipelineType());
             }
 
             switch (dispatchType) {
@@ -1478,7 +1490,7 @@ public abstract class Pipeline
 
     setRunning(true);
 
-    switch (pipelineMeta.getPipelineType()) {
+    switch (getPipelineType()) {
       case Normal:
 
         // Now start all the threads...
@@ -2234,11 +2246,10 @@ public abstract class Pipeline
 
     // We are going to add an extra IRowSet to this iTransform.
     IRowSet rowSet =
-        switch (pipelineMeta.getPipelineType()) {
+        switch (getPipelineType()) {
           case Normal -> new BlockingRowSet(rowSetSize);
           case SingleThreaded -> new QueueRowSet();
-          default ->
-              throw new HopException("Unhandled pipeline type: " + 
pipelineMeta.getPipelineType());
+          default -> throw new HopException("Unhandled pipeline type: " + 
getPipelineType());
         };
 
     // Add this rowset to the list of active rowsets for the selected transform
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 f0f0834355..fdd91de3b9 100644
--- a/engine/src/main/java/org/apache/hop/pipeline/PipelineMeta.java
+++ b/engine/src/main/java/org/apache/hop/pipeline/PipelineMeta.java
@@ -3345,24 +3345,6 @@ public class PipelineMeta extends AbstractMeta
     previousTransformCache.clear();
   }
 
-  /**
-   * Gets the pipeline type.
-   *
-   * @return the pipelineType
-   */
-  public PipelineType getPipelineType() {
-    return info.getPipelineType();
-  }
-
-  /**
-   * Sets the pipeline type.
-   *
-   * @param pipelineType the pipelineType to set
-   */
-  public void setPipelineType(PipelineType pipelineType) {
-    this.info.setPipelineType(pipelineType);
-  }
-
   public void addTransformChangeListener(ITransformMetaChangeListener 
listener) {
     transformChangeListeners.add(listener);
   }
@@ -3437,8 +3419,14 @@ public class PipelineMeta extends AbstractMeta
   }
 
   /**
-   * The PipelineType enum describes the various types of pipelines in terms 
of execution, including
-   * Normal, Serial Single-Threaded, and Single-Threaded.
+   * Describes how an engine drives the transforms of a pipeline. This is a 
property of the engine
+   * that executes the pipeline, not of the pipeline itself: the very same 
pipeline runs under
+   * either type, so it is never stored in the .hpl file. See {@link
+   * org.apache.hop.pipeline.engine.IPipelineEngine#getPipelineType()}.
+   *
+   * <p>Transforms use it in {@link
+   * 
org.apache.hop.pipeline.transform.BaseTransformMeta#getSupportedPipelineTypes()}
 to declare
+   * which of these execution models they can cope with.
    */
   @SuppressWarnings("java:S115")
   @Getter
diff --git a/engine/src/main/java/org/apache/hop/pipeline/PipelineMetaInfo.java 
b/engine/src/main/java/org/apache/hop/pipeline/PipelineMetaInfo.java
index ee4b26c052..58db96e4f4 100644
--- a/engine/src/main/java/org/apache/hop/pipeline/PipelineMetaInfo.java
+++ b/engine/src/main/java/org/apache/hop/pipeline/PipelineMetaInfo.java
@@ -42,10 +42,6 @@ public class PipelineMetaInfo extends AbstractMetaInfo {
   @HopMetadataProperty(key = "transform_performance_capturing_size_limit")
   protected String transformPerformanceCapturingSizeLimit;
 
-  /** The pipeline type. */
-  @HopMetadataProperty(key = "pipeline_type", storeWithCode = true)
-  protected PipelineMeta.PipelineType pipelineType;
-
   /** The status of the pipeline. */
   @HopMetadataProperty(key = "pipeline_status")
   protected int pipelineStatus;
@@ -58,6 +54,5 @@ public class PipelineMetaInfo extends AbstractMetaInfo {
     this.capturingTransformPerformanceSnapShots = false;
     this.transformPerformanceCapturingDelay = 1000; // every 1 seconds
     this.transformPerformanceCapturingSizeLimit = "100"; // maximum 100 data 
points
-    this.pipelineType = PipelineMeta.PipelineType.Normal;
   }
 }
diff --git 
a/engine/src/main/java/org/apache/hop/pipeline/engine/IPipelineEngine.java 
b/engine/src/main/java/org/apache/hop/pipeline/engine/IPipelineEngine.java
index 16234eb4d4..e723456aa8 100644
--- a/engine/src/main/java/org/apache/hop/pipeline/engine/IPipelineEngine.java
+++ b/engine/src/main/java/org/apache/hop/pipeline/engine/IPipelineEngine.java
@@ -80,6 +80,22 @@ public interface IPipelineEngine<T extends PipelineMeta>
    */
   PipelineEngineCapabilities getEngineCapabilities();
 
+  /**
+   * The way in which this engine drives the transforms of the pipeline. 
Engines that give every
+   * transform its own thread report {@link PipelineMeta.PipelineType#Normal}; 
engines whose
+   * transforms are driven one iteration at a time from a single thread report 
{@link
+   * PipelineMeta.PipelineType#SingleThreaded}.
+   *
+   * <p>This is a property of the engine, not of the pipeline: the same 
pipeline runs under either
+   * type and an engine must never write its choice back into the pipeline 
metadata.
+   *
+   * @return The execution model of this engine, {@link 
PipelineMeta.PipelineType#Normal} by
+   *     default.
+   */
+  default PipelineMeta.PipelineType getPipelineType() {
+    return PipelineMeta.PipelineType.Normal;
+  }
+
   /**
    * Engine's compatibility verdict for a transform plugin. Default is UNKNOWN 
("no opinion, fall
    * back to annotation"). Engines override to surface SUPPORTED / UNSUPPORTED 
authoritatively.
diff --git 
a/engine/src/test/java/org/apache/hop/pipeline/PipelineTypeIsEngineDrivenTest.java
 
b/engine/src/test/java/org/apache/hop/pipeline/PipelineTypeIsEngineDrivenTest.java
new file mode 100644
index 0000000000..8b9f8e2313
--- /dev/null
+++ 
b/engine/src/test/java/org/apache/hop/pipeline/PipelineTypeIsEngineDrivenTest.java
@@ -0,0 +1,131 @@
+/*
+ * 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;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertFalse;
+
+import java.io.ByteArrayInputStream;
+import java.nio.charset.StandardCharsets;
+import org.apache.hop.core.logging.LoggingObject;
+import org.apache.hop.core.variables.IVariables;
+import org.apache.hop.core.variables.Variables;
+import org.apache.hop.junit.rules.RestoreHopEngineEnvironmentExtension;
+import org.apache.hop.metadata.api.IHopMetadataProvider;
+import org.apache.hop.metadata.serializer.memory.MemoryMetadataProvider;
+import org.apache.hop.pipeline.engines.local.LocalPipelineEngine;
+import org.junit.jupiter.api.Test;
+import org.junit.jupiter.api.extension.ExtendWith;
+
+/**
+ * How the transforms of a pipeline are driven - every transform in its own 
thread, or all of them
+ * one iteration at a time - is decided by the engine the pipeline is run 
with. It is not a property
+ * of the pipeline and is therefore never written to the .hpl file.
+ *
+ * <p>It used to be one: engines wrote their choice into the {@link 
PipelineMeta} they were handed,
+ * which in Hop GUI is the very object the editor holds. A single run with the 
single threaded
+ * engine left {@code <pipeline_type>SingleThreaded</pipeline_type>} in the 
file on the next save,
+ * nothing ever set it back, and the local engine then started no threads at 
all for it - issue
+ * #8262, where the pipeline hung after switching the run configuration back.
+ */
+@ExtendWith(RestoreHopEngineEnvironmentExtension.class)
+class PipelineTypeIsEngineDrivenTest {
+
+  /** A pipeline as saved by a Hop version that still wrote the execution 
model into the file. */
+  private static final String PIPELINE_WITH_LEGACY_TYPE =
+      """
+      <pipeline>
+        <info>
+          <name>legacy-pipeline-type</name>
+          <pipeline_type>SingleThreaded</pipeline_type>
+        </info>
+      </pipeline>
+      """;
+
+  private final IVariables variables = new Variables();
+  private final IHopMetadataProvider metadataProvider = new 
MemoryMetadataProvider();
+
+  /** An engine that drives its transforms from a single thread, like the 
LocalSingle engine. */
+  private static class SingleThreadedTestEngine extends LocalPipelineEngine {
+    SingleThreadedTestEngine(PipelineMeta pipelineMeta) {
+      super(pipelineMeta, new Variables(), new LoggingObject("test"));
+    }
+
+    @Override
+    public PipelineMeta.PipelineType getPipelineType() {
+      return PipelineMeta.PipelineType.SingleThreaded;
+    }
+  }
+
+  @Test
+  void executionModelIsNotWrittenToTheFile() throws Exception {
+    PipelineMeta pipelineMeta = new PipelineMeta();
+    pipelineMeta.setName("no-pipeline-type");
+
+    assertFalse(pipelineMeta.getXml(variables).contains("pipeline_type"));
+  }
+
+  @Test
+  void legacyExecutionModelInTheFileIsIgnored() throws Exception {
+    PipelineMeta pipelineMeta =
+        new PipelineMeta(
+            new 
ByteArrayInputStream(PIPELINE_WITH_LEGACY_TYPE.getBytes(StandardCharsets.UTF_8)),
+            metadataProvider,
+            variables);
+
+    // The file still loads, ...
+    assertEquals("legacy-pipeline-type", pipelineMeta.getName());
+    // ... the stale element no longer makes the local engine skip starting 
its threads, ...
+    assertEquals(
+        PipelineMeta.PipelineType.Normal,
+        new LocalPipelineEngine(pipelineMeta).getPipelineType(),
+        "a pipeline_type left in the file by an older release must not affect 
the engine");
+    // ... and saving the pipeline again drops it.
+    assertFalse(pipelineMeta.getXml(variables).contains("pipeline_type"));
+  }
+
+  @Test
+  void eachEngineCarriesItsOwnExecutionModel() {
+    PipelineMeta pipelineMeta = new PipelineMeta();
+    pipelineMeta.setName("shared-metadata");
+
+    // Hop GUI hands the same PipelineMeta instance to every engine it creates 
for the tab.
+    LocalPipelineEngine local = new LocalPipelineEngine(pipelineMeta);
+    SingleThreadedTestEngine singleThreaded = new 
SingleThreadedTestEngine(pipelineMeta);
+
+    assertEquals(PipelineMeta.PipelineType.SingleThreaded, 
singleThreaded.getPipelineType());
+    assertEquals(
+        PipelineMeta.PipelineType.Normal,
+        local.getPipelineType(),
+        "one engine's execution model must not leak into another engine over 
the same pipeline");
+  }
+
+  @Test
+  void anEngineCanBeDrivenSingleThreadedWithoutTouchingTheMetadata() throws 
Exception {
+    PipelineMeta pipelineMeta = new PipelineMeta();
+    pipelineMeta.setName("embedded-sub-pipeline");
+
+    // Transforms that embed a sub-pipeline (Simple Mapping, Kafka Consumer, 
the Beam and Spark
+    // workers) push rows through it one batch at a time and say so on the 
engine they create.
+    LocalPipelineEngine subPipeline = new LocalPipelineEngine(pipelineMeta);
+    subPipeline.setPipelineType(PipelineMeta.PipelineType.SingleThreaded);
+
+    assertEquals(PipelineMeta.PipelineType.SingleThreaded, 
subPipeline.getPipelineType());
+    assertFalse(pipelineMeta.getXml(variables).contains("pipeline_type"));
+  }
+}
diff --git 
a/plugins/engines/beam/src/main/java/org/apache/hop/beam/core/transform/TransformBatchTransform.java
 
b/plugins/engines/beam/src/main/java/org/apache/hop/beam/core/transform/TransformBatchTransform.java
index 7de4835801..feb9d1bc1d 100644
--- 
a/plugins/engines/beam/src/main/java/org/apache/hop/beam/core/transform/TransformBatchTransform.java
+++ 
b/plugins/engines/beam/src/main/java/org/apache/hop/beam/core/transform/TransformBatchTransform.java
@@ -376,7 +376,6 @@ public class TransformBatchTransform extends 
TransformTransform {
           //
           pipelineMeta = new PipelineMeta();
           pipelineMeta.setName(transformName);
-          
pipelineMeta.setPipelineType(PipelineMeta.PipelineType.SingleThreaded);
           pipelineMeta.setMetadataProvider(metadataProvider);
 
           // When the first row ends up in the buffer we start the timer.
@@ -493,6 +492,7 @@ public class TransformBatchTransform extends 
TransformTransform {
           pipeline =
               new LocalPipelineEngine(
                   pipelineMeta, variables, new 
LoggingObject("apache-beam-transform"));
+          pipeline.setPipelineType(PipelineMeta.PipelineType.SingleThreaded);
           pipeline.setLogLevel(
               
context.getPipelineOptions().as(HopPipelineExecutionOptions.class).getLogLevel());
           pipeline.setMetadataProvider(pipelineMeta.getMetadataProvider());
diff --git 
a/plugins/engines/beam/src/main/java/org/apache/hop/beam/core/transform/TransformFn.java
 
b/plugins/engines/beam/src/main/java/org/apache/hop/beam/core/transform/TransformFn.java
index 3500aa7895..756f85177d 100644
--- 
a/plugins/engines/beam/src/main/java/org/apache/hop/beam/core/transform/TransformFn.java
+++ 
b/plugins/engines/beam/src/main/java/org/apache/hop/beam/core/transform/TransformFn.java
@@ -226,7 +226,6 @@ public class TransformFn extends TransformBaseFn {
     //
     pipelineMeta = new PipelineMeta();
     pipelineMeta.setName(transformName);
-    pipelineMeta.setPipelineType(PipelineMeta.PipelineType.SingleThreaded);
     pipelineMeta.setMetadataProvider(metadataProvider);
 
     // Input row metadata...
@@ -346,6 +345,7 @@ public class TransformFn extends TransformBaseFn {
       parentLoggingObject = new LoggingObject("apache-beam-transform");
     }
     pipeline = new LocalPipelineEngine(pipelineMeta, variables, 
parentLoggingObject);
+    pipeline.setPipelineType(PipelineMeta.PipelineType.SingleThreaded);
     pipeline.setLogLevel(
         
context.getPipelineOptions().as(HopPipelineExecutionOptions.class).getLogLevel());
     pipeline.setMetadataProvider(pipelineMeta.getMetadataProvider());
diff --git 
a/plugins/engines/single-threaded/src/main/java/org/apache/hop/pipeline/engines/singlethreaded/SingleThreadedPipelineEngine.java
 
b/plugins/engines/single-threaded/src/main/java/org/apache/hop/pipeline/engines/singlethreaded/SingleThreadedPipelineEngine.java
index 0deb7de0b2..7ce863d4ef 100644
--- 
a/plugins/engines/single-threaded/src/main/java/org/apache/hop/pipeline/engines/singlethreaded/SingleThreadedPipelineEngine.java
+++ 
b/plugins/engines/single-threaded/src/main/java/org/apache/hop/pipeline/engines/singlethreaded/SingleThreadedPipelineEngine.java
@@ -66,10 +66,13 @@ public class SingleThreadedPipelineEngine extends Pipeline
     return new PipelineEngineCapabilities(true, true, true, true);
   }
 
+  /**
+   * This engine drives every transform from a single thread. Reported here 
rather than written into
+   * the pipeline metadata: the .hpl being executed says nothing about the 
engine it runs on.
+   */
   @Override
-  public void prepareExecution() throws HopException {
-    pipelineMeta.setPipelineType(PipelineMeta.PipelineType.SingleThreaded);
-    super.prepareExecution();
+  public PipelineMeta.PipelineType getPipelineType() {
+    return PipelineMeta.PipelineType.SingleThreaded;
   }
 
   @Override
diff --git 
a/plugins/engines/spark/src/main/java/org/apache/hop/spark/core/HopMapPartitionsFn.java
 
b/plugins/engines/spark/src/main/java/org/apache/hop/spark/core/HopMapPartitionsFn.java
index e310fd2034..f4c0377789 100644
--- 
a/plugins/engines/spark/src/main/java/org/apache/hop/spark/core/HopMapPartitionsFn.java
+++ 
b/plugins/engines/spark/src/main/java/org/apache/hop/spark/core/HopMapPartitionsFn.java
@@ -306,7 +306,6 @@ public class HopMapPartitionsFn implements 
MapPartitionsFunction<Row, Row>, Seri
 
       PipelineMeta pipelineMeta = new PipelineMeta();
       pipelineMeta.setName(transformName);
-      pipelineMeta.setPipelineType(PipelineMeta.PipelineType.SingleThreaded);
       pipelineMeta.setMetadataProvider(metadataProvider);
 
       if (infoTransforms.size() != infoRowMetaJsons.size()
@@ -404,6 +403,7 @@ public class HopMapPartitionsFn implements 
MapPartitionsFunction<Row, Row>, Seri
       LocalPipelineEngine pipeline =
           new LocalPipelineEngine(
               pipelineMeta, variables, new 
LoggingObject("apache-spark-transform"));
+      pipeline.setPipelineType(PipelineMeta.PipelineType.SingleThreaded);
       pipeline.setMetadataProvider(metadataProvider);
       pipeline
           .getPipelineRunConfiguration()
diff --git 
a/plugins/engines/spark/src/test/java/org/apache/hop/spark/core/SparkParallelFileContextTest.java
 
b/plugins/engines/spark/src/test/java/org/apache/hop/spark/core/SparkParallelFileContextTest.java
index 7ae9c78e96..f1694093bf 100644
--- 
a/plugins/engines/spark/src/test/java/org/apache/hop/spark/core/SparkParallelFileContextTest.java
+++ 
b/plugins/engines/spark/src/test/java/org/apache/hop/spark/core/SparkParallelFileContextTest.java
@@ -56,7 +56,6 @@ class SparkParallelFileContextTest {
   void sparkParallelFileContextSetsBeamContextAndOverridesId() throws 
Exception {
     PipelineMeta pipelineMeta = new PipelineMeta();
     pipelineMeta.setName("writer");
-    pipelineMeta.setPipelineType(PipelineMeta.PipelineType.SingleThreaded);
     TransformMeta dummyTm = new TransformMeta("writer", new DummyMeta());
     dummyTm.setTransformPluginId("Dummy");
     pipelineMeta.addTransform(dummyTm);
@@ -65,6 +64,7 @@ class SparkParallelFileContextTest {
     LocalPipelineEngine pipeline =
         new LocalPipelineEngine(
             pipelineMeta, variables, new 
LoggingObject("spark-file-context-test"));
+    pipeline.setPipelineType(PipelineMeta.PipelineType.SingleThreaded);
     pipeline.prepareExecution();
 
     // Simulate init having set UUID-based Internal.Transform.ID
diff --git 
a/plugins/tech/azure/src/main/java/org/apache/hop/pipeline/transforms/eventhubs/listen/AzureListener.java
 
b/plugins/tech/azure/src/main/java/org/apache/hop/pipeline/transforms/eventhubs/listen/AzureListener.java
index fa0992eea4..275c3eb205 100644
--- 
a/plugins/tech/azure/src/main/java/org/apache/hop/pipeline/transforms/eventhubs/listen/AzureListener.java
+++ 
b/plugins/tech/azure/src/main/java/org/apache/hop/pipeline/transforms/eventhubs/listen/AzureListener.java
@@ -116,8 +116,8 @@ public class AzureListener extends 
BaseTransform<AzureListenerMeta, AzureListene
       data.stt = true;
       data.sttMaxWaitTime = Const.toLong(resolve(meta.getBatchMaxWaitTime()), 
-1L);
       data.sttPipelineMeta = AzureListenerMeta.loadBatchPipelineMeta(meta, 
metadataProvider, this);
-      
data.sttPipelineMeta.setPipelineType(PipelineMeta.PipelineType.SingleThreaded);
       data.sttPipeline = new LocalPipelineEngine(data.sttPipelineMeta, this, 
this);
+      
data.sttPipeline.setPipelineType(PipelineMeta.PipelineType.SingleThreaded);
       data.sttPipeline.setParent(getPipeline());
       data.sttPipeline.setParentPipeline(getPipeline());
 
diff --git 
a/plugins/transforms/kafka/src/main/java/org/apache/hop/pipeline/transforms/kafka/consumer/KafkaConsumerInput.java
 
b/plugins/transforms/kafka/src/main/java/org/apache/hop/pipeline/transforms/kafka/consumer/KafkaConsumerInput.java
index 7288f916c5..59970c7610 100644
--- 
a/plugins/transforms/kafka/src/main/java/org/apache/hop/pipeline/transforms/kafka/consumer/KafkaConsumerInput.java
+++ 
b/plugins/transforms/kafka/src/main/java/org/apache/hop/pipeline/transforms/kafka/consumer/KafkaConsumerInput.java
@@ -118,7 +118,6 @@ public class KafkaConsumerInput
       PipelineMeta subTransMeta = new PipelineMeta(realFilename, 
metadataProvider, this);
       subTransMeta.setMetadataProvider(metadataProvider);
       subTransMeta.setFilename(realFilename);
-      subTransMeta.setPipelineType(PipelineMeta.PipelineType.SingleThreaded);
       logDetailed("Loaded sub-pipeline '" + realFilename + "'");
 
       PipelineRunConfiguration runConfiguration =
@@ -132,6 +131,7 @@ public class KafkaConsumerInput
               false);
 
       LocalPipelineEngine kafkaPipeline = new 
LocalPipelineEngine(subTransMeta, this, this);
+      kafkaPipeline.setPipelineType(PipelineMeta.PipelineType.SingleThreaded);
       kafkaPipeline.setParentPipeline(getPipeline());
       kafkaPipeline.setPipelineRunConfiguration(runConfiguration);
       kafkaPipeline.prepareExecution();
diff --git 
a/plugins/transforms/mapping/src/main/java/org/apache/hop/pipeline/transforms/mapping/SimpleMapping.java
 
b/plugins/transforms/mapping/src/main/java/org/apache/hop/pipeline/transforms/mapping/SimpleMapping.java
index 97036a3a52..43c7f3d4d0 100644
--- 
a/plugins/transforms/mapping/src/main/java/org/apache/hop/pipeline/transforms/mapping/SimpleMapping.java
+++ 
b/plugins/transforms/mapping/src/main/java/org/apache/hop/pipeline/transforms/mapping/SimpleMapping.java
@@ -120,15 +120,6 @@ public class SimpleMapping extends 
BaseTransform<SimpleMappingMeta, SimpleMappin
   }
 
   public void prepareMappingExecution() throws HopException {
-    boolean singleThreaded =
-        data.isBeamContext()
-            || (getPipeline() != null
-                && getPipeline().getPipelineMeta() != null
-                && getPipeline().getPipelineMeta().getPipelineType()
-                    == PipelineMeta.PipelineType.SingleThreaded);
-    if (singleThreaded) {
-      
data.mappingPipelineMeta.setPipelineType(PipelineMeta.PipelineType.SingleThreaded);
-    }
     SimpleMappingData simpleMappingData = getData();
     // Resolve pipeline full name in case variables are used and pipeline meta 
is not initialized in
     // advance
@@ -168,6 +159,10 @@ public class SimpleMapping extends 
BaseTransform<SimpleMappingMeta, SimpleMappin
                   simpleMappingData.mappingPipelineMeta);
     }
 
+    if (isSingleThreaded()) {
+      
simpleMappingData.mappingPipeline.setPipelineType(PipelineMeta.PipelineType.SingleThreaded);
+    }
+
     // Copy the parameters over...
     //
     simpleMappingData.mappingPipeline.copyParametersFromDefinitions(
@@ -242,6 +237,16 @@ public class SimpleMapping extends 
BaseTransform<SimpleMappingMeta, SimpleMappin
     getPipeline().addActiveSubPipeline(getTransformName(), 
simpleMappingData.mappingPipeline);
   }
 
+  /**
+   * Are we running in a context where the sub-pipeline has to be driven row 
by row from this
+   * transform instead of by threads of its own?
+   */
+  private boolean isSingleThreaded() {
+    return data.isBeamContext()
+        || (getPipeline() != null
+            && getPipeline().getPipelineType() == 
PipelineMeta.PipelineType.SingleThreaded);
+  }
+
   public static List<MappingInput> findMappingInputs(Pipeline mappingPipeline) 
{
     return MappingTransforms.findMappingInputs(mappingPipeline);
   }
@@ -272,13 +277,7 @@ public class SimpleMapping extends 
BaseTransform<SimpleMappingMeta, SimpleMappin
           // We don't want to process one-row batches in a parallel engine 
where we need to wait for
           // the threads to finish.
           //
-          boolean singleThreaded =
-              data.isBeamContext()
-                  || (getPipeline() != null
-                      && getPipeline().getPipelineMeta() != null
-                      && getPipeline().getPipelineMeta().getPipelineType()
-                          == PipelineMeta.PipelineType.SingleThreaded);
-          if (singleThreaded) {
+          if (isSingleThreaded()) {
             data.executor = new 
SingleThreadedPipelineExecutor(data.mappingPipeline);
           }
 
diff --git 
a/plugins/transforms/streamschemamerge/src/test/java/org/apache/hop/pipeline/transforms/streamschemamerge/TestUtilities.java
 
b/plugins/transforms/streamschemamerge/src/test/java/org/apache/hop/pipeline/transforms/streamschemamerge/TestUtilities.java
index 6ecd8697ae..4ddb1c7c5b 100755
--- 
a/plugins/transforms/streamschemamerge/src/test/java/org/apache/hop/pipeline/transforms/streamschemamerge/TestUtilities.java
+++ 
b/plugins/transforms/streamschemamerge/src/test/java/org/apache/hop/pipeline/transforms/streamschemamerge/TestUtilities.java
@@ -371,7 +371,6 @@ public class TestUtilities {
 
   public static Pipeline loadAndRunPipeline(String path, Object... parameters) 
throws Exception {
     PipelineMeta pipelineMeta = new PipelineMeta();
-    pipelineMeta.setPipelineType(PipelineMeta.PipelineType.Normal);
 
     Pipeline trans = new LocalPipelineEngine(pipelineMeta);
     if (parameters != null) {

Reply via email to