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) {