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 dd96261d0d joinrows hardening and block on beam/spark, fixes #2353
(#8574)
dd96261d0d is described below
commit dd96261d0dc8a16dd5e8516c01cccfc2915d39be
Author: Hans Van Akelyen <[email protected]>
AuthorDate: Thu Sep 24 16:05:06 2026 +0200
joinrows hardening and block on beam/spark, fixes #2353 (#8574)
* joinrows hardening and block on beam/spark, fixes #2353
* fix remarks
---
.../pipeline/beam/getting-started-with-beam.adoc | 1 +
.../native-spark-pipeline-engine.adoc | 1 +
.../ROOT/pages/pipeline/transforms/joinrows.adoc | 35 +-
.../single_threaded/0001-join-rows-validation.hpl | 309 +++++++++++++++++
.../single_threaded/0001-join-rows.hpl | 234 +++++++++++++
.../0002-join-rows-after-transforms-validation.hpl | 301 +++++++++++++++++
.../0002-join-rows-after-transforms.hpl | 286 ++++++++++++++++
.../0003-join-rows-three-inputs-validation.hpl | 373 +++++++++++++++++++++
.../0003-join-rows-three-inputs.hpl | 285 ++++++++++++++++
.../single_threaded/dev-env-config.json | 10 +
integration-tests/single_threaded/hop-config.json | 290 ++++++++++++++++
.../single_threaded/main-0001-join-rows.hwf | 138 ++++++++
.../main-0002-join-rows-after-transforms.hwf | 138 ++++++++
.../main-0003-join-rows-three-inputs.hwf | 138 ++++++++
.../metadata/pipeline-run-configuration/local.json | 17 +
.../single-threaded.json | 11 +
.../metadata/workflow-run-configuration/local.json | 9 +
.../single_threaded/project-config.json | 15 +
plugins/engines/beam/pom.xml | 6 +
.../HopPipelineMetaToBeamPipelineConverter.java | 5 +-
.../engines/BeamPipelineEngineSupportsTest.java | 10 +
plugins/engines/spark/pom.xml | 6 +
.../pipeline/HopPipelineMetaToSparkConverter.java | 4 +-
.../java/org/apache/hop/spark/util/SparkConst.java | 1 +
.../HopPipelineMetaToSparkConverterTest.java | 16 +-
.../hop/pipeline/transforms/joinrows/JoinRows.java | 99 ++++--
.../pipeline/transforms/joinrows/JoinRowsMeta.java | 3 +-
.../joinrows/JoinRowsSingleThreadedTest.java | 260 ++++++++++++++
28 files changed, 2958 insertions(+), 43 deletions(-)
diff --git
a/docs/hop-user-manual/modules/ROOT/pages/pipeline/beam/getting-started-with-beam.adoc
b/docs/hop-user-manual/modules/ROOT/pages/pipeline/beam/getting-started-with-beam.adoc
index 0104fa65f4..5dc4e78d9c 100644
---
a/docs/hop-user-manual/modules/ROOT/pages/pipeline/beam/getting-started-with-beam.adoc
+++
b/docs/hop-user-manual/modules/ROOT/pages/pipeline/beam/getting-started-with-beam.adoc
@@ -191,6 +191,7 @@ Each of these transforms would therefore only ever see the
slice of the data one
* xref:pipeline/transforms/uniquerowsbyhashset.adoc[Unique Rows (HashSet)] :
Use the `Memory Group By` instead
* xref:pipeline/transforms/groupby.adoc[Group By] : Use the `Memory Group By`
instead
* xref:pipeline/transforms/sort.adoc[Sort Rows]
+* xref:pipeline/transforms/joinrows.adoc[Join Rows (cartesian product)] : Add
the same constant field to both inputs and use a `Merge Join` on that field
instead
The Hop GUI knows about this list: when you set the canvas palette filter to a
Beam engine these transforms are hidden, and running a pipeline that still
contains one is blocked with the reason shown above.
diff --git
a/docs/hop-user-manual/modules/ROOT/pages/pipeline/pipeline-run-configurations/native-spark-pipeline-engine.adoc
b/docs/hop-user-manual/modules/ROOT/pages/pipeline/pipeline-run-configurations/native-spark-pipeline-engine.adoc
index 4b0132dc99..8bb210e5b0 100644
---
a/docs/hop-user-manual/modules/ROOT/pages/pipeline/pipeline-run-configurations/native-spark-pipeline-engine.adoc
+++
b/docs/hop-user-manual/modules/ROOT/pages/pipeline/pipeline-run-configurations/native-spark-pipeline-engine.adoc
@@ -184,6 +184,7 @@ These rely on waiting for other transforms, full-stream
ordering, or multi-strea
* xref:pipeline/transforms/blockuntiltransformsfinish.adoc[Blocking until
transforms finish]
* xref:pipeline/transforms/detectemptystream.adoc[Detect Empty Stream]
* xref:pipeline/transforms/identifylastrow.adoc[Identify last row in a stream]
+* xref:pipeline/transforms/joinrows.adoc[Join Rows (cartesian product)] — add
the same constant field to both inputs and use
xref:pipeline/transforms/mergejoin.adoc[Merge Join] on that field instead
* xref:pipeline/transforms/mergerows.adoc[Merge rows (diff)]
* xref:pipeline/transforms/multimerge.adoc[Multiway Merge Join]
* xref:pipeline/transforms/sortedmerge.adoc[Sorted Merge]
diff --git
a/docs/hop-user-manual/modules/ROOT/pages/pipeline/transforms/joinrows.adoc
b/docs/hop-user-manual/modules/ROOT/pages/pipeline/transforms/joinrows.adoc
index dcb1a537a2..0effd1ef39 100644
--- a/docs/hop-user-manual/modules/ROOT/pages/pipeline/transforms/joinrows.adoc
+++ b/docs/hop-user-manual/modules/ROOT/pages/pipeline/transforms/joinrows.adoc
@@ -17,13 +17,38 @@ under the License.
:documentationPath: /pipeline/transforms/
:language: en_US
:description: The Join Rows transform allows you to produce combinations
(Cartesian product) of all rows in the input streams.
-:page-engines: Hop Engine=yes, Single Threaded=yes, Native Spark=yes, Beam
Spark=maybe, Beam Flink=maybe, Beam Dataflow=maybe
+:page-engines: Hop Engine=yes, Single Threaded=yes, Native Spark=no, Beam
Spark=no, Beam Flink=no, Beam Dataflow=no
= image:transforms/icons/joinrows.svg[Join Rows transform Icon,
role="image-doc-icon"] Join Rows
== Description
-The Join Rows (cartesian product) transform allows you to combine/join
multiple input streams
(https://en.wikipedia.org/wiki/Cartesian_product[Cartesian product^]) without
joining on keys. It works best with one row from each stream. You can add a
condition to only join when a condition is met.
+The Join Rows (cartesian product) transform allows you to combine/join
multiple input streams
(https://en.wikipedia.org/wiki/Cartesian_product[Cartesian product^]) without
joining on keys.
+Every row of one stream is combined with every row of the other streams, so
the number of output rows is the product of the number of rows of all streams.
+You can add a condition to only write the combinations that meet it.
+
+The transform works like a lookup:
+
+. It first reads all the rows of the other streams, the ones that are not the
main stream.
+It writes them to temporary files, and keeps a stream in memory as well when
it has no more than *Max. cache size* rows.
+. Then it reads the main stream one row at a time, and writes the combinations
of that row with all the rows of the other streams before it reads the next
main row.
+
+The rows of the main stream are never kept, so make the largest stream the
main stream.
+Since the other streams are read completely before the first main row is
joined, they can't depend on the main stream.
+When one of the streams has no rows, the transform writes no rows at all.
+
+[WARNING]
+====
+This transform is not supported on the distributed pipeline engines: neither
the Beam engines nor the native Spark engine.
+A cartesian product needs every row of every input in one place, and these
engines spread the rows over many workers, bundles or partitions, so
combinations would go missing.
+The Hop GUI hides the transform when the canvas palette filter is set to one
of these engines, and running a pipeline that still contains it is refused,
both in the GUI and with `hop-run`.
+
+To build a cartesian product on these engines, add the same constant field to
both inputs with an xref:pipeline/transforms/addconstant.adoc[Add Constants]
transform and join them with a xref:pipeline/transforms/mergejoin.adoc[Merge
Join] on that field.
+See
xref:pipeline/beam/getting-started-with-beam.adoc#_unsupported_transforms[Getting
started with Beam] and
xref:pipeline/pipeline-run-configurations/native-spark-pipeline-engine.adoc#transforms-not-supported-on-native-spark[Transforms
not supported on Native Spark].
+====
+
+NOTE: The transform gives the same result on the
xref:pipeline/pipeline-run-configurations/single-threaded-pipeline-engine.adoc[single
threaded engine] as on the Hop engine.
+That engine runs the previous transforms first, so all the rows of the other
streams are there when Join Rows starts.
== Options
@@ -31,9 +56,9 @@ The Join Rows (cartesian product) transform allows you to
combine/join multiple
|===
|Option|Description
|Transform name|Name of the transform this name has to be unique in a single
pipeline.
-|Temp directory|Specify the name of the directory where the system stores
temporary files in case you want to combine more then the cached number of rows.
+|Temp directory|Specify the name of the directory where the system stores
temporary files in case you want to combine more than the cached number of rows.
|TMP-file prefix|This is the prefix of the temporary files that will be
generated.
-|Max. cache size|The number of rows to cache before the system reads data from
temporary files; required when you want to combine large row sets that do not
fit into memory.
-|Main transform to read from|Specifies the transform from which to read most
of the data; while the data from other transforms are cached or spooled to
disk, the data from this transform is not.
+|Max. cache size|The maximum number of rows of a stream to keep in memory. A
stream with more rows is read back from its temporary file for every main row;
required when you want to combine large row sets that do not fit into memory.
+|Main transform to read from|The main stream: the transform from which to read
most of the data. The rows of the other transforms are cached or spooled to
disk, the rows of this transform are not. When left empty, the first input
stream is the main stream.
|The Condition(s)|You can enter a complex condition to limit the number of
output row.
|===
diff --git a/integration-tests/single_threaded/0001-join-rows-validation.hpl
b/integration-tests/single_threaded/0001-join-rows-validation.hpl
new file mode 100644
index 0000000000..7ca7b72562
--- /dev/null
+++ b/integration-tests/single_threaded/0001-join-rows-validation.hpl
@@ -0,0 +1,309 @@
+<?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>0001-join-rows-validation</name>
+ <name_sync_with_filename>Y</name_sync_with_filename>
+ <description>Verify the rows written by 0001-join-rows running on the
single threaded engine</description>
+ <extended_description/>
+ <pipeline_version/>
+ <pipeline_type>Normal</pipeline_type>
+ <parameters>
+ </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/09/24 12:00:00.000</created_date>
+ <modified_user>-</modified_user>
+ <modified_date>2026/09/24 12:00:00.000</modified_date>
+ </info>
+ <notepads>
+ </notepads>
+ <order>
+ <hop>
+ <from>Expected</from>
+ <to>Merge rows (diff)</to>
+ <enabled>Y</enabled>
+ </hop>
+ <hop>
+ <from>Output of the test</from>
+ <to>Sort output</to>
+ <enabled>Y</enabled>
+ </hop>
+ <hop>
+ <from>Sort output</from>
+ <to>Merge rows (diff)</to>
+ <enabled>Y</enabled>
+ </hop>
+ <hop>
+ <from>Merge rows (diff)</from>
+ <to>Not identical</to>
+ <enabled>Y</enabled>
+ </hop>
+ <hop>
+ <from>Not identical</from>
+ <to>Abort</to>
+ <enabled>Y</enabled>
+ </hop>
+ </order>
+ <transform>
+ <name>Expected</name>
+ <type>DataGrid</type>
+ <description/>
+ <distribute>Y</distribute>
+ <custom_distribution/>
+ <copies>1</copies>
+ <partitioning>
+ <method>none</method>
+ <schema_name/>
+ </partitioning>
+ <fields>
+ <field>
+ <set_empty_string>N</set_empty_string>
+ <length>-1</length>
+ <name>id</name>
+ <precision>-1</precision>
+ <type>Integer</type>
+ </field>
+ <field>
+ <set_empty_string>N</set_empty_string>
+ <length>-1</length>
+ <name>name</name>
+ <precision>-1</precision>
+ <type>String</type>
+ </field>
+ </fields>
+ <data>
+ <line>
+ <item>1</item>
+ <item>x</item>
+ </line>
+ <line>
+ <item>1</item>
+ <item>y</item>
+ </line>
+ <line>
+ <item>2</item>
+ <item>x</item>
+ </line>
+ <line>
+ <item>2</item>
+ <item>y</item>
+ </line>
+ <line>
+ <item>3</item>
+ <item>x</item>
+ </line>
+ <line>
+ <item>3</item>
+ <item>y</item>
+ </line>
+ </data>
+ <attributes/>
+ <GUI>
+ <xloc>96</xloc>
+ <yloc>64</yloc>
+ </GUI>
+ </transform>
+ <transform>
+ <name>Output of the test</name>
+ <type>CSVInput</type>
+ <description/>
+ <distribute>Y</distribute>
+ <custom_distribution/>
+ <copies>1</copies>
+ <partitioning>
+ <method>none</method>
+ <schema_name/>
+ </partitioning>
+
<filename>${java.io.tmpdir}/hop-it-single-threaded/0001-join-rows.csv</filename>
+ <filename_field/>
+ <rownum_field/>
+ <include_filename>N</include_filename>
+ <separator>;</separator>
+ <enclosure>"</enclosure>
+ <header>Y</header>
+ <buffer_size>50000</buffer_size>
+ <lazy_conversion>N</lazy_conversion>
+ <add_filename_result>N</add_filename_result>
+ <parallel>N</parallel>
+ <newline_possible>N</newline_possible>
+ <encoding>UTF-8</encoding>
+ <fields>
+ <field>
+ <name>id</name>
+ <type>Integer</type>
+ <format>0</format>
+ <currency/>
+ <decimal/>
+ <group/>
+ <length>-1</length>
+ <precision>-1</precision>
+ <trim_type>none</trim_type>
+ </field>
+ <field>
+ <name>name</name>
+ <type>String</type>
+ <format></format>
+ <currency/>
+ <decimal/>
+ <group/>
+ <length>-1</length>
+ <precision>-1</precision>
+ <trim_type>none</trim_type>
+ </field>
+ </fields>
+ <attributes/>
+ <GUI>
+ <xloc>96</xloc>
+ <yloc>224</yloc>
+ </GUI>
+ </transform>
+ <transform>
+ <name>Sort output</name>
+ <type>SortRows</type>
+ <description/>
+ <distribute>Y</distribute>
+ <custom_distribution/>
+ <copies>1</copies>
+ <partitioning>
+ <method>none</method>
+ <schema_name/>
+ </partitioning>
+ <directory>${java.io.tmpdir}</directory>
+ <sort_prefix>out</sort_prefix>
+ <sort_size>1000000</sort_size>
+ <free_memory/>
+ <unique_rows>N</unique_rows>
+ <compress>N</compress>
+ <fields>
+ <field>
+ <name>id</name>
+ <ascending>Y</ascending>
+ <case_sensitive>N</case_sensitive>
+ <collator_enabled>N</collator_enabled>
+ <collator_strength>0</collator_strength>
+ <presorted>N</presorted>
+ </field>
+ <field>
+ <name>name</name>
+ <ascending>Y</ascending>
+ <case_sensitive>N</case_sensitive>
+ <collator_enabled>N</collator_enabled>
+ <collator_strength>0</collator_strength>
+ <presorted>N</presorted>
+ </field>
+ </fields>
+ <attributes/>
+ <GUI>
+ <xloc>256</xloc>
+ <yloc>224</yloc>
+ </GUI>
+ </transform>
+ <transform>
+ <name>Merge rows (diff)</name>
+ <type>MergeRows</type>
+ <description/>
+ <distribute>Y</distribute>
+ <custom_distribution/>
+ <copies>1</copies>
+ <partitioning>
+ <method>none</method>
+ <schema_name/>
+ </partitioning>
+ <keys>
+ <key>id</key>
+ <key>name</key>
+ </keys>
+ <values>
+ <value>id</value>
+ <value>name</value>
+ </values>
+ <flag_field>flagfield</flag_field>
+ <reference>Expected</reference>
+ <compare>Sort output</compare>
+ <attributes/>
+ <GUI>
+ <xloc>416</xloc>
+ <yloc>144</yloc>
+ </GUI>
+ </transform>
+ <transform>
+ <name>Not identical</name>
+ <type>FilterRows</type>
+ <description/>
+ <distribute>Y</distribute>
+ <custom_distribution/>
+ <copies>1</copies>
+ <partitioning>
+ <method>none</method>
+ <schema_name/>
+ </partitioning>
+ <send_true_to/>
+ <send_false_to/>
+ <compare>
+ <condition>
+ <negated>N</negated>
+ <leftvalue>flagfield</leftvalue>
+ <function><></function>
+ <rightvalue/>
+ <value>
+ <name>constant</name>
+ <type>String</type>
+ <text>identical</text>
+ <length>-1</length>
+ <precision>-1</precision>
+ <isnull>N</isnull>
+ <mask/>
+ </value>
+ </condition>
+ </compare>
+ <attributes/>
+ <GUI>
+ <xloc>576</xloc>
+ <yloc>144</yloc>
+ </GUI>
+ </transform>
+ <transform>
+ <name>Abort</name>
+ <type>Abort</type>
+ <description/>
+ <distribute>Y</distribute>
+ <custom_distribution/>
+ <copies>1</copies>
+ <partitioning>
+ <method>none</method>
+ <schema_name/>
+ </partitioning>
+ <abort_option>ABORT_WITH_ERROR</abort_option>
+ <always_log_rows>Y</always_log_rows>
+ <message>0001-join-rows: the rows written by Join Rows are not the
expected ones</message>
+ <row_threshold>0</row_threshold>
+ <attributes/>
+ <GUI>
+ <xloc>736</xloc>
+ <yloc>144</yloc>
+ </GUI>
+ </transform>
+ <transform_error_handling>
+ </transform_error_handling>
+ <attributes/>
+</pipeline>
diff --git a/integration-tests/single_threaded/0001-join-rows.hpl
b/integration-tests/single_threaded/0001-join-rows.hpl
new file mode 100644
index 0000000000..e0b4a5298d
--- /dev/null
+++ b/integration-tests/single_threaded/0001-join-rows.hpl
@@ -0,0 +1,234 @@
+<?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>0001-join-rows</name>
+ <name_sync_with_filename>Y</name_sync_with_filename>
+ <description>Join Rows (cartesian product) of two Data Grids on the single
threaded engine (#2353)</description>
+ <extended_description/>
+ <pipeline_version/>
+ <pipeline_type>Normal</pipeline_type>
+ <parameters>
+ </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/09/24 12:00:00.000</created_date>
+ <modified_user>-</modified_user>
+ <modified_date>2026/09/24 12:00:00.000</modified_date>
+ </info>
+ <notepads>
+ </notepads>
+ <order>
+ <hop>
+ <from>ids</from>
+ <to>Join rows</to>
+ <enabled>Y</enabled>
+ </hop>
+ <hop>
+ <from>names</from>
+ <to>Join rows</to>
+ <enabled>Y</enabled>
+ </hop>
+ <hop>
+ <from>Join rows</from>
+ <to>Output</to>
+ <enabled>Y</enabled>
+ </hop>
+ </order>
+ <transform>
+ <name>ids</name>
+ <type>DataGrid</type>
+ <description/>
+ <distribute>Y</distribute>
+ <custom_distribution/>
+ <copies>1</copies>
+ <partitioning>
+ <method>none</method>
+ <schema_name/>
+ </partitioning>
+ <fields>
+ <field>
+ <set_empty_string>N</set_empty_string>
+ <length>-1</length>
+ <name>id</name>
+ <precision>-1</precision>
+ <type>Integer</type>
+ </field>
+ </fields>
+ <data>
+ <line>
+ <item>1</item>
+ </line>
+ <line>
+ <item>2</item>
+ </line>
+ <line>
+ <item>3</item>
+ </line>
+ </data>
+ <attributes/>
+ <GUI>
+ <xloc>96</xloc>
+ <yloc>64</yloc>
+ </GUI>
+ </transform>
+ <transform>
+ <name>names</name>
+ <type>DataGrid</type>
+ <description/>
+ <distribute>Y</distribute>
+ <custom_distribution/>
+ <copies>1</copies>
+ <partitioning>
+ <method>none</method>
+ <schema_name/>
+ </partitioning>
+ <fields>
+ <field>
+ <set_empty_string>N</set_empty_string>
+ <length>-1</length>
+ <name>name</name>
+ <precision>-1</precision>
+ <type>String</type>
+ </field>
+ </fields>
+ <data>
+ <line>
+ <item>x</item>
+ </line>
+ <line>
+ <item>y</item>
+ </line>
+ </data>
+ <attributes/>
+ <GUI>
+ <xloc>96</xloc>
+ <yloc>192</yloc>
+ </GUI>
+ </transform>
+ <transform>
+ <name>Join rows</name>
+ <type>JoinRows</type>
+ <description/>
+ <distribute>Y</distribute>
+ <custom_distribution/>
+ <copies>1</copies>
+ <partitioning>
+ <method>none</method>
+ <schema_name/>
+ </partitioning>
+ <directory>%%java.io.tmpdir%%</directory>
+ <prefix>out</prefix>
+ <cache_size>500</cache_size>
+ <main>ids</main>
+ <compare>
+ <condition>
+ <negated>N</negated>
+ <operator>-</operator>
+ <leftvalue/>
+ <function>=</function>
+ <rightvalue/>
+ </condition>
+ </compare>
+ <attributes/>
+ <GUI>
+ <xloc>400</xloc>
+ <yloc>192</yloc>
+ </GUI>
+ </transform>
+ <transform>
+ <name>Output</name>
+ <type>TextFileOutput</type>
+ <description/>
+ <distribute>Y</distribute>
+ <custom_distribution/>
+ <copies>1</copies>
+ <partitioning>
+ <method>none</method>
+ <schema_name/>
+ </partitioning>
+ <separator>;</separator>
+ <enclosure/>
+ <enclosure_forced>N</enclosure_forced>
+ <enclosure_fix_disabled>N</enclosure_fix_disabled>
+ <header>Y</header>
+ <footer>N</footer>
+ <format>UNIX</format>
+ <compression>None</compression>
+ <encoding>UTF-8</encoding>
+ <endedLine/>
+ <fileNameInField>N</fileNameInField>
+ <fileNameField/>
+ <create_parent_folder>Y</create_parent_folder>
+ <file>
+ <name>${java.io.tmpdir}/hop-it-single-threaded/0001-join-rows</name>
+ <servlet_output>N</servlet_output>
+ <do_not_open_new_file_init>N</do_not_open_new_file_init>
+ <extention>csv</extention>
+ <append>N</append>
+ <split>N</split>
+ <haspartno>N</haspartno>
+ <add_date>N</add_date>
+ <add_time>N</add_time>
+ <SpecifyFormat>N</SpecifyFormat>
+ <date_time_format/>
+ <add_to_result_filenames>N</add_to_result_filenames>
+ <pad>N</pad>
+ <fast_dump>N</fast_dump>
+ <splitevery/>
+ </file>
+ <fields>
+ <field>
+ <name>id</name>
+ <type>Integer</type>
+ <format>0</format>
+ <currency/>
+ <decimal/>
+ <group/>
+ <nullif/>
+ <trim_type>none</trim_type>
+ <length>-1</length>
+ <precision>-1</precision>
+ </field>
+ <field>
+ <name>name</name>
+ <type>String</type>
+ <format></format>
+ <currency/>
+ <decimal/>
+ <group/>
+ <nullif/>
+ <trim_type>none</trim_type>
+ <length>-1</length>
+ <precision>-1</precision>
+ </field>
+ </fields>
+ <attributes/>
+ <GUI>
+ <xloc>640</xloc>
+ <yloc>192</yloc>
+ </GUI>
+ </transform>
+ <transform_error_handling>
+ </transform_error_handling>
+ <attributes/>
+</pipeline>
diff --git
a/integration-tests/single_threaded/0002-join-rows-after-transforms-validation.hpl
b/integration-tests/single_threaded/0002-join-rows-after-transforms-validation.hpl
new file mode 100644
index 0000000000..30d56cced6
--- /dev/null
+++
b/integration-tests/single_threaded/0002-join-rows-after-transforms-validation.hpl
@@ -0,0 +1,301 @@
+<?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>0002-join-rows-after-transforms-validation</name>
+ <name_sync_with_filename>Y</name_sync_with_filename>
+ <description>Verify the rows written by 0002-join-rows-after-transforms
running on the single threaded engine</description>
+ <extended_description/>
+ <pipeline_version/>
+ <pipeline_type>Normal</pipeline_type>
+ <parameters>
+ </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/09/24 12:00:00.000</created_date>
+ <modified_user>-</modified_user>
+ <modified_date>2026/09/24 12:00:00.000</modified_date>
+ </info>
+ <notepads>
+ </notepads>
+ <order>
+ <hop>
+ <from>Expected</from>
+ <to>Merge rows (diff)</to>
+ <enabled>Y</enabled>
+ </hop>
+ <hop>
+ <from>Output of the test</from>
+ <to>Sort output</to>
+ <enabled>Y</enabled>
+ </hop>
+ <hop>
+ <from>Sort output</from>
+ <to>Merge rows (diff)</to>
+ <enabled>Y</enabled>
+ </hop>
+ <hop>
+ <from>Merge rows (diff)</from>
+ <to>Not identical</to>
+ <enabled>Y</enabled>
+ </hop>
+ <hop>
+ <from>Not identical</from>
+ <to>Abort</to>
+ <enabled>Y</enabled>
+ </hop>
+ </order>
+ <transform>
+ <name>Expected</name>
+ <type>DataGrid</type>
+ <description/>
+ <distribute>Y</distribute>
+ <custom_distribution/>
+ <copies>1</copies>
+ <partitioning>
+ <method>none</method>
+ <schema_name/>
+ </partitioning>
+ <fields>
+ <field>
+ <set_empty_string>N</set_empty_string>
+ <length>-1</length>
+ <name>id</name>
+ <precision>-1</precision>
+ <type>Integer</type>
+ </field>
+ <field>
+ <set_empty_string>N</set_empty_string>
+ <length>-1</length>
+ <name>name</name>
+ <precision>-1</precision>
+ <type>String</type>
+ </field>
+ </fields>
+ <data>
+ <line>
+ <item>1</item>
+ <item>x</item>
+ </line>
+ <line>
+ <item>1</item>
+ <item>y</item>
+ </line>
+ <line>
+ <item>3</item>
+ <item>x</item>
+ </line>
+ <line>
+ <item>3</item>
+ <item>y</item>
+ </line>
+ </data>
+ <attributes/>
+ <GUI>
+ <xloc>96</xloc>
+ <yloc>64</yloc>
+ </GUI>
+ </transform>
+ <transform>
+ <name>Output of the test</name>
+ <type>CSVInput</type>
+ <description/>
+ <distribute>Y</distribute>
+ <custom_distribution/>
+ <copies>1</copies>
+ <partitioning>
+ <method>none</method>
+ <schema_name/>
+ </partitioning>
+
<filename>${java.io.tmpdir}/hop-it-single-threaded/0002-join-rows-after-transforms.csv</filename>
+ <filename_field/>
+ <rownum_field/>
+ <include_filename>N</include_filename>
+ <separator>;</separator>
+ <enclosure>"</enclosure>
+ <header>Y</header>
+ <buffer_size>50000</buffer_size>
+ <lazy_conversion>N</lazy_conversion>
+ <add_filename_result>N</add_filename_result>
+ <parallel>N</parallel>
+ <newline_possible>N</newline_possible>
+ <encoding>UTF-8</encoding>
+ <fields>
+ <field>
+ <name>id</name>
+ <type>Integer</type>
+ <format>0</format>
+ <currency/>
+ <decimal/>
+ <group/>
+ <length>-1</length>
+ <precision>-1</precision>
+ <trim_type>none</trim_type>
+ </field>
+ <field>
+ <name>name</name>
+ <type>String</type>
+ <format></format>
+ <currency/>
+ <decimal/>
+ <group/>
+ <length>-1</length>
+ <precision>-1</precision>
+ <trim_type>none</trim_type>
+ </field>
+ </fields>
+ <attributes/>
+ <GUI>
+ <xloc>96</xloc>
+ <yloc>224</yloc>
+ </GUI>
+ </transform>
+ <transform>
+ <name>Sort output</name>
+ <type>SortRows</type>
+ <description/>
+ <distribute>Y</distribute>
+ <custom_distribution/>
+ <copies>1</copies>
+ <partitioning>
+ <method>none</method>
+ <schema_name/>
+ </partitioning>
+ <directory>${java.io.tmpdir}</directory>
+ <sort_prefix>out</sort_prefix>
+ <sort_size>1000000</sort_size>
+ <free_memory/>
+ <unique_rows>N</unique_rows>
+ <compress>N</compress>
+ <fields>
+ <field>
+ <name>id</name>
+ <ascending>Y</ascending>
+ <case_sensitive>N</case_sensitive>
+ <collator_enabled>N</collator_enabled>
+ <collator_strength>0</collator_strength>
+ <presorted>N</presorted>
+ </field>
+ <field>
+ <name>name</name>
+ <ascending>Y</ascending>
+ <case_sensitive>N</case_sensitive>
+ <collator_enabled>N</collator_enabled>
+ <collator_strength>0</collator_strength>
+ <presorted>N</presorted>
+ </field>
+ </fields>
+ <attributes/>
+ <GUI>
+ <xloc>256</xloc>
+ <yloc>224</yloc>
+ </GUI>
+ </transform>
+ <transform>
+ <name>Merge rows (diff)</name>
+ <type>MergeRows</type>
+ <description/>
+ <distribute>Y</distribute>
+ <custom_distribution/>
+ <copies>1</copies>
+ <partitioning>
+ <method>none</method>
+ <schema_name/>
+ </partitioning>
+ <keys>
+ <key>id</key>
+ <key>name</key>
+ </keys>
+ <values>
+ <value>id</value>
+ <value>name</value>
+ </values>
+ <flag_field>flagfield</flag_field>
+ <reference>Expected</reference>
+ <compare>Sort output</compare>
+ <attributes/>
+ <GUI>
+ <xloc>416</xloc>
+ <yloc>144</yloc>
+ </GUI>
+ </transform>
+ <transform>
+ <name>Not identical</name>
+ <type>FilterRows</type>
+ <description/>
+ <distribute>Y</distribute>
+ <custom_distribution/>
+ <copies>1</copies>
+ <partitioning>
+ <method>none</method>
+ <schema_name/>
+ </partitioning>
+ <send_true_to/>
+ <send_false_to/>
+ <compare>
+ <condition>
+ <negated>N</negated>
+ <leftvalue>flagfield</leftvalue>
+ <function><></function>
+ <rightvalue/>
+ <value>
+ <name>constant</name>
+ <type>String</type>
+ <text>identical</text>
+ <length>-1</length>
+ <precision>-1</precision>
+ <isnull>N</isnull>
+ <mask/>
+ </value>
+ </condition>
+ </compare>
+ <attributes/>
+ <GUI>
+ <xloc>576</xloc>
+ <yloc>144</yloc>
+ </GUI>
+ </transform>
+ <transform>
+ <name>Abort</name>
+ <type>Abort</type>
+ <description/>
+ <distribute>Y</distribute>
+ <custom_distribution/>
+ <copies>1</copies>
+ <partitioning>
+ <method>none</method>
+ <schema_name/>
+ </partitioning>
+ <abort_option>ABORT_WITH_ERROR</abort_option>
+ <always_log_rows>Y</always_log_rows>
+ <message>0002-join-rows-after-transforms: the rows written by Join Rows
are not the expected ones</message>
+ <row_threshold>0</row_threshold>
+ <attributes/>
+ <GUI>
+ <xloc>736</xloc>
+ <yloc>144</yloc>
+ </GUI>
+ </transform>
+ <transform_error_handling>
+ </transform_error_handling>
+ <attributes/>
+</pipeline>
diff --git
a/integration-tests/single_threaded/0002-join-rows-after-transforms.hpl
b/integration-tests/single_threaded/0002-join-rows-after-transforms.hpl
new file mode 100644
index 0000000000..f0e2bdf60c
--- /dev/null
+++ b/integration-tests/single_threaded/0002-join-rows-after-transforms.hpl
@@ -0,0 +1,286 @@
+<?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>0002-join-rows-after-transforms</name>
+ <name_sync_with_filename>Y</name_sync_with_filename>
+ <description>Join Rows fed by other transforms, spilling to a temporary
file, with a condition, on the single threaded engine (#2353)</description>
+ <extended_description/>
+ <pipeline_version/>
+ <pipeline_type>Normal</pipeline_type>
+ <parameters>
+ </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/09/24 12:00:00.000</created_date>
+ <modified_user>-</modified_user>
+ <modified_date>2026/09/24 12:00:00.000</modified_date>
+ </info>
+ <notepads>
+ </notepads>
+ <order>
+ <hop>
+ <from>ids</from>
+ <to>pass ids</to>
+ <enabled>Y</enabled>
+ </hop>
+ <hop>
+ <from>names</from>
+ <to>pass names</to>
+ <enabled>Y</enabled>
+ </hop>
+ <hop>
+ <from>pass ids</from>
+ <to>Join rows</to>
+ <enabled>Y</enabled>
+ </hop>
+ <hop>
+ <from>pass names</from>
+ <to>Join rows</to>
+ <enabled>Y</enabled>
+ </hop>
+ <hop>
+ <from>Join rows</from>
+ <to>Output</to>
+ <enabled>Y</enabled>
+ </hop>
+ </order>
+ <transform>
+ <name>ids</name>
+ <type>DataGrid</type>
+ <description/>
+ <distribute>Y</distribute>
+ <custom_distribution/>
+ <copies>1</copies>
+ <partitioning>
+ <method>none</method>
+ <schema_name/>
+ </partitioning>
+ <fields>
+ <field>
+ <set_empty_string>N</set_empty_string>
+ <length>-1</length>
+ <name>id</name>
+ <precision>-1</precision>
+ <type>Integer</type>
+ </field>
+ </fields>
+ <data>
+ <line>
+ <item>1</item>
+ </line>
+ <line>
+ <item>2</item>
+ </line>
+ <line>
+ <item>3</item>
+ </line>
+ </data>
+ <attributes/>
+ <GUI>
+ <xloc>96</xloc>
+ <yloc>64</yloc>
+ </GUI>
+ </transform>
+ <transform>
+ <name>names</name>
+ <type>DataGrid</type>
+ <description/>
+ <distribute>Y</distribute>
+ <custom_distribution/>
+ <copies>1</copies>
+ <partitioning>
+ <method>none</method>
+ <schema_name/>
+ </partitioning>
+ <fields>
+ <field>
+ <set_empty_string>N</set_empty_string>
+ <length>-1</length>
+ <name>name</name>
+ <precision>-1</precision>
+ <type>String</type>
+ </field>
+ </fields>
+ <data>
+ <line>
+ <item>x</item>
+ </line>
+ <line>
+ <item>y</item>
+ </line>
+ </data>
+ <attributes/>
+ <GUI>
+ <xloc>96</xloc>
+ <yloc>192</yloc>
+ </GUI>
+ </transform>
+ <transform>
+ <name>pass ids</name>
+ <type>Dummy</type>
+ <description/>
+ <distribute>Y</distribute>
+ <custom_distribution/>
+ <copies>1</copies>
+ <partitioning>
+ <method>none</method>
+ <schema_name/>
+ </partitioning>
+ <attributes/>
+ <GUI>
+ <xloc>256</xloc>
+ <yloc>64</yloc>
+ </GUI>
+ </transform>
+ <transform>
+ <name>pass names</name>
+ <type>Dummy</type>
+ <description/>
+ <distribute>Y</distribute>
+ <custom_distribution/>
+ <copies>1</copies>
+ <partitioning>
+ <method>none</method>
+ <schema_name/>
+ </partitioning>
+ <attributes/>
+ <GUI>
+ <xloc>256</xloc>
+ <yloc>192</yloc>
+ </GUI>
+ </transform>
+ <transform>
+ <name>Join rows</name>
+ <type>JoinRows</type>
+ <description/>
+ <distribute>Y</distribute>
+ <custom_distribution/>
+ <copies>1</copies>
+ <partitioning>
+ <method>none</method>
+ <schema_name/>
+ </partitioning>
+ <directory>%%java.io.tmpdir%%</directory>
+ <prefix>out</prefix>
+ <cache_size>1</cache_size>
+ <main>pass ids</main>
+ <compare>
+ <condition>
+ <negated>N</negated>
+ <leftvalue>id</leftvalue>
+ <function><></function>
+ <rightvalue/>
+ <value>
+ <name>constant</name>
+ <type>Integer</type>
+ <text>2</text>
+ <length>-1</length>
+ <precision>0</precision>
+ <isnull>N</isnull>
+ <mask>####0;-####0</mask>
+ </value>
+ </condition>
+ </compare>
+ <attributes/>
+ <GUI>
+ <xloc>400</xloc>
+ <yloc>192</yloc>
+ </GUI>
+ </transform>
+ <transform>
+ <name>Output</name>
+ <type>TextFileOutput</type>
+ <description/>
+ <distribute>Y</distribute>
+ <custom_distribution/>
+ <copies>1</copies>
+ <partitioning>
+ <method>none</method>
+ <schema_name/>
+ </partitioning>
+ <separator>;</separator>
+ <enclosure/>
+ <enclosure_forced>N</enclosure_forced>
+ <enclosure_fix_disabled>N</enclosure_fix_disabled>
+ <header>Y</header>
+ <footer>N</footer>
+ <format>UNIX</format>
+ <compression>None</compression>
+ <encoding>UTF-8</encoding>
+ <endedLine/>
+ <fileNameInField>N</fileNameInField>
+ <fileNameField/>
+ <create_parent_folder>Y</create_parent_folder>
+ <file>
+
<name>${java.io.tmpdir}/hop-it-single-threaded/0002-join-rows-after-transforms</name>
+ <servlet_output>N</servlet_output>
+ <do_not_open_new_file_init>N</do_not_open_new_file_init>
+ <extention>csv</extention>
+ <append>N</append>
+ <split>N</split>
+ <haspartno>N</haspartno>
+ <add_date>N</add_date>
+ <add_time>N</add_time>
+ <SpecifyFormat>N</SpecifyFormat>
+ <date_time_format/>
+ <add_to_result_filenames>N</add_to_result_filenames>
+ <pad>N</pad>
+ <fast_dump>N</fast_dump>
+ <splitevery/>
+ </file>
+ <fields>
+ <field>
+ <name>id</name>
+ <type>Integer</type>
+ <format>0</format>
+ <currency/>
+ <decimal/>
+ <group/>
+ <nullif/>
+ <trim_type>none</trim_type>
+ <length>-1</length>
+ <precision>-1</precision>
+ </field>
+ <field>
+ <name>name</name>
+ <type>String</type>
+ <format></format>
+ <currency/>
+ <decimal/>
+ <group/>
+ <nullif/>
+ <trim_type>none</trim_type>
+ <length>-1</length>
+ <precision>-1</precision>
+ </field>
+ </fields>
+ <attributes/>
+ <GUI>
+ <xloc>640</xloc>
+ <yloc>192</yloc>
+ </GUI>
+ </transform>
+ <transform_error_handling>
+ </transform_error_handling>
+ <attributes/>
+</pipeline>
diff --git
a/integration-tests/single_threaded/0003-join-rows-three-inputs-validation.hpl
b/integration-tests/single_threaded/0003-join-rows-three-inputs-validation.hpl
new file mode 100644
index 0000000000..eec3b16842
--- /dev/null
+++
b/integration-tests/single_threaded/0003-join-rows-three-inputs-validation.hpl
@@ -0,0 +1,373 @@
+<?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>0003-join-rows-three-inputs-validation</name>
+ <name_sync_with_filename>Y</name_sync_with_filename>
+ <description>Verify the rows written by 0003-join-rows-three-inputs
running on the single threaded engine</description>
+ <extended_description/>
+ <pipeline_version/>
+ <pipeline_type>Normal</pipeline_type>
+ <parameters>
+ </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/09/24 12:00:00.000</created_date>
+ <modified_user>-</modified_user>
+ <modified_date>2026/09/24 12:00:00.000</modified_date>
+ </info>
+ <notepads>
+ </notepads>
+ <order>
+ <hop>
+ <from>Expected</from>
+ <to>Merge rows (diff)</to>
+ <enabled>Y</enabled>
+ </hop>
+ <hop>
+ <from>Output of the test</from>
+ <to>Sort output</to>
+ <enabled>Y</enabled>
+ </hop>
+ <hop>
+ <from>Sort output</from>
+ <to>Merge rows (diff)</to>
+ <enabled>Y</enabled>
+ </hop>
+ <hop>
+ <from>Merge rows (diff)</from>
+ <to>Not identical</to>
+ <enabled>Y</enabled>
+ </hop>
+ <hop>
+ <from>Not identical</from>
+ <to>Abort</to>
+ <enabled>Y</enabled>
+ </hop>
+ </order>
+ <transform>
+ <name>Expected</name>
+ <type>DataGrid</type>
+ <description/>
+ <distribute>Y</distribute>
+ <custom_distribution/>
+ <copies>1</copies>
+ <partitioning>
+ <method>none</method>
+ <schema_name/>
+ </partitioning>
+ <fields>
+ <field>
+ <set_empty_string>N</set_empty_string>
+ <length>-1</length>
+ <name>id</name>
+ <precision>-1</precision>
+ <type>Integer</type>
+ </field>
+ <field>
+ <set_empty_string>N</set_empty_string>
+ <length>-1</length>
+ <name>name</name>
+ <precision>-1</precision>
+ <type>String</type>
+ </field>
+ <field>
+ <set_empty_string>N</set_empty_string>
+ <length>-1</length>
+ <name>code</name>
+ <precision>-1</precision>
+ <type>String</type>
+ </field>
+ </fields>
+ <data>
+ <line>
+ <item>1</item>
+ <item>x</item>
+ <item>a</item>
+ </line>
+ <line>
+ <item>1</item>
+ <item>x</item>
+ <item>b</item>
+ </line>
+ <line>
+ <item>1</item>
+ <item>y</item>
+ <item>a</item>
+ </line>
+ <line>
+ <item>1</item>
+ <item>y</item>
+ <item>b</item>
+ </line>
+ <line>
+ <item>2</item>
+ <item>x</item>
+ <item>a</item>
+ </line>
+ <line>
+ <item>2</item>
+ <item>x</item>
+ <item>b</item>
+ </line>
+ <line>
+ <item>2</item>
+ <item>y</item>
+ <item>a</item>
+ </line>
+ <line>
+ <item>2</item>
+ <item>y</item>
+ <item>b</item>
+ </line>
+ <line>
+ <item>3</item>
+ <item>x</item>
+ <item>a</item>
+ </line>
+ <line>
+ <item>3</item>
+ <item>x</item>
+ <item>b</item>
+ </line>
+ <line>
+ <item>3</item>
+ <item>y</item>
+ <item>a</item>
+ </line>
+ <line>
+ <item>3</item>
+ <item>y</item>
+ <item>b</item>
+ </line>
+ </data>
+ <attributes/>
+ <GUI>
+ <xloc>96</xloc>
+ <yloc>64</yloc>
+ </GUI>
+ </transform>
+ <transform>
+ <name>Output of the test</name>
+ <type>CSVInput</type>
+ <description/>
+ <distribute>Y</distribute>
+ <custom_distribution/>
+ <copies>1</copies>
+ <partitioning>
+ <method>none</method>
+ <schema_name/>
+ </partitioning>
+
<filename>${java.io.tmpdir}/hop-it-single-threaded/0003-join-rows-three-inputs.csv</filename>
+ <filename_field/>
+ <rownum_field/>
+ <include_filename>N</include_filename>
+ <separator>;</separator>
+ <enclosure>"</enclosure>
+ <header>Y</header>
+ <buffer_size>50000</buffer_size>
+ <lazy_conversion>N</lazy_conversion>
+ <add_filename_result>N</add_filename_result>
+ <parallel>N</parallel>
+ <newline_possible>N</newline_possible>
+ <encoding>UTF-8</encoding>
+ <fields>
+ <field>
+ <name>id</name>
+ <type>Integer</type>
+ <format>0</format>
+ <currency/>
+ <decimal/>
+ <group/>
+ <length>-1</length>
+ <precision>-1</precision>
+ <trim_type>none</trim_type>
+ </field>
+ <field>
+ <name>name</name>
+ <type>String</type>
+ <format></format>
+ <currency/>
+ <decimal/>
+ <group/>
+ <length>-1</length>
+ <precision>-1</precision>
+ <trim_type>none</trim_type>
+ </field>
+ <field>
+ <name>code</name>
+ <type>String</type>
+ <format></format>
+ <currency/>
+ <decimal/>
+ <group/>
+ <length>-1</length>
+ <precision>-1</precision>
+ <trim_type>none</trim_type>
+ </field>
+ </fields>
+ <attributes/>
+ <GUI>
+ <xloc>96</xloc>
+ <yloc>224</yloc>
+ </GUI>
+ </transform>
+ <transform>
+ <name>Sort output</name>
+ <type>SortRows</type>
+ <description/>
+ <distribute>Y</distribute>
+ <custom_distribution/>
+ <copies>1</copies>
+ <partitioning>
+ <method>none</method>
+ <schema_name/>
+ </partitioning>
+ <directory>${java.io.tmpdir}</directory>
+ <sort_prefix>out</sort_prefix>
+ <sort_size>1000000</sort_size>
+ <free_memory/>
+ <unique_rows>N</unique_rows>
+ <compress>N</compress>
+ <fields>
+ <field>
+ <name>id</name>
+ <ascending>Y</ascending>
+ <case_sensitive>N</case_sensitive>
+ <collator_enabled>N</collator_enabled>
+ <collator_strength>0</collator_strength>
+ <presorted>N</presorted>
+ </field>
+ <field>
+ <name>name</name>
+ <ascending>Y</ascending>
+ <case_sensitive>N</case_sensitive>
+ <collator_enabled>N</collator_enabled>
+ <collator_strength>0</collator_strength>
+ <presorted>N</presorted>
+ </field>
+ <field>
+ <name>code</name>
+ <ascending>Y</ascending>
+ <case_sensitive>N</case_sensitive>
+ <collator_enabled>N</collator_enabled>
+ <collator_strength>0</collator_strength>
+ <presorted>N</presorted>
+ </field>
+ </fields>
+ <attributes/>
+ <GUI>
+ <xloc>256</xloc>
+ <yloc>224</yloc>
+ </GUI>
+ </transform>
+ <transform>
+ <name>Merge rows (diff)</name>
+ <type>MergeRows</type>
+ <description/>
+ <distribute>Y</distribute>
+ <custom_distribution/>
+ <copies>1</copies>
+ <partitioning>
+ <method>none</method>
+ <schema_name/>
+ </partitioning>
+ <keys>
+ <key>id</key>
+ <key>name</key>
+ <key>code</key>
+ </keys>
+ <values>
+ <value>id</value>
+ <value>name</value>
+ <value>code</value>
+ </values>
+ <flag_field>flagfield</flag_field>
+ <reference>Expected</reference>
+ <compare>Sort output</compare>
+ <attributes/>
+ <GUI>
+ <xloc>416</xloc>
+ <yloc>144</yloc>
+ </GUI>
+ </transform>
+ <transform>
+ <name>Not identical</name>
+ <type>FilterRows</type>
+ <description/>
+ <distribute>Y</distribute>
+ <custom_distribution/>
+ <copies>1</copies>
+ <partitioning>
+ <method>none</method>
+ <schema_name/>
+ </partitioning>
+ <send_true_to/>
+ <send_false_to/>
+ <compare>
+ <condition>
+ <negated>N</negated>
+ <leftvalue>flagfield</leftvalue>
+ <function><></function>
+ <rightvalue/>
+ <value>
+ <name>constant</name>
+ <type>String</type>
+ <text>identical</text>
+ <length>-1</length>
+ <precision>-1</precision>
+ <isnull>N</isnull>
+ <mask/>
+ </value>
+ </condition>
+ </compare>
+ <attributes/>
+ <GUI>
+ <xloc>576</xloc>
+ <yloc>144</yloc>
+ </GUI>
+ </transform>
+ <transform>
+ <name>Abort</name>
+ <type>Abort</type>
+ <description/>
+ <distribute>Y</distribute>
+ <custom_distribution/>
+ <copies>1</copies>
+ <partitioning>
+ <method>none</method>
+ <schema_name/>
+ </partitioning>
+ <abort_option>ABORT_WITH_ERROR</abort_option>
+ <always_log_rows>Y</always_log_rows>
+ <message>0003-join-rows-three-inputs: the rows written by Join Rows are
not the expected ones</message>
+ <row_threshold>0</row_threshold>
+ <attributes/>
+ <GUI>
+ <xloc>736</xloc>
+ <yloc>144</yloc>
+ </GUI>
+ </transform>
+ <transform_error_handling>
+ </transform_error_handling>
+ <attributes/>
+</pipeline>
diff --git a/integration-tests/single_threaded/0003-join-rows-three-inputs.hpl
b/integration-tests/single_threaded/0003-join-rows-three-inputs.hpl
new file mode 100644
index 0000000000..02a1451898
--- /dev/null
+++ b/integration-tests/single_threaded/0003-join-rows-three-inputs.hpl
@@ -0,0 +1,285 @@
+<?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>0003-join-rows-three-inputs</name>
+ <name_sync_with_filename>Y</name_sync_with_filename>
+ <description>Join Rows of three inputs on the single threaded engine
(#2353)</description>
+ <extended_description/>
+ <pipeline_version/>
+ <pipeline_type>Normal</pipeline_type>
+ <parameters>
+ </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/09/24 12:00:00.000</created_date>
+ <modified_user>-</modified_user>
+ <modified_date>2026/09/24 12:00:00.000</modified_date>
+ </info>
+ <notepads>
+ </notepads>
+ <order>
+ <hop>
+ <from>names</from>
+ <to>Join rows</to>
+ <enabled>Y</enabled>
+ </hop>
+ <hop>
+ <from>codes</from>
+ <to>Join rows</to>
+ <enabled>Y</enabled>
+ </hop>
+ <hop>
+ <from>ids</from>
+ <to>Join rows</to>
+ <enabled>Y</enabled>
+ </hop>
+ <hop>
+ <from>Join rows</from>
+ <to>Output</to>
+ <enabled>Y</enabled>
+ </hop>
+ </order>
+ <transform>
+ <name>ids</name>
+ <type>DataGrid</type>
+ <description/>
+ <distribute>Y</distribute>
+ <custom_distribution/>
+ <copies>1</copies>
+ <partitioning>
+ <method>none</method>
+ <schema_name/>
+ </partitioning>
+ <fields>
+ <field>
+ <set_empty_string>N</set_empty_string>
+ <length>-1</length>
+ <name>id</name>
+ <precision>-1</precision>
+ <type>Integer</type>
+ </field>
+ </fields>
+ <data>
+ <line>
+ <item>1</item>
+ </line>
+ <line>
+ <item>2</item>
+ </line>
+ <line>
+ <item>3</item>
+ </line>
+ </data>
+ <attributes/>
+ <GUI>
+ <xloc>96</xloc>
+ <yloc>64</yloc>
+ </GUI>
+ </transform>
+ <transform>
+ <name>names</name>
+ <type>DataGrid</type>
+ <description/>
+ <distribute>Y</distribute>
+ <custom_distribution/>
+ <copies>1</copies>
+ <partitioning>
+ <method>none</method>
+ <schema_name/>
+ </partitioning>
+ <fields>
+ <field>
+ <set_empty_string>N</set_empty_string>
+ <length>-1</length>
+ <name>name</name>
+ <precision>-1</precision>
+ <type>String</type>
+ </field>
+ </fields>
+ <data>
+ <line>
+ <item>x</item>
+ </line>
+ <line>
+ <item>y</item>
+ </line>
+ </data>
+ <attributes/>
+ <GUI>
+ <xloc>96</xloc>
+ <yloc>192</yloc>
+ </GUI>
+ </transform>
+ <transform>
+ <name>codes</name>
+ <type>DataGrid</type>
+ <description/>
+ <distribute>Y</distribute>
+ <custom_distribution/>
+ <copies>1</copies>
+ <partitioning>
+ <method>none</method>
+ <schema_name/>
+ </partitioning>
+ <fields>
+ <field>
+ <set_empty_string>N</set_empty_string>
+ <length>-1</length>
+ <name>code</name>
+ <precision>-1</precision>
+ <type>String</type>
+ </field>
+ </fields>
+ <data>
+ <line>
+ <item>a</item>
+ </line>
+ <line>
+ <item>b</item>
+ </line>
+ </data>
+ <attributes/>
+ <GUI>
+ <xloc>96</xloc>
+ <yloc>320</yloc>
+ </GUI>
+ </transform>
+ <transform>
+ <name>Join rows</name>
+ <type>JoinRows</type>
+ <description/>
+ <distribute>Y</distribute>
+ <custom_distribution/>
+ <copies>1</copies>
+ <partitioning>
+ <method>none</method>
+ <schema_name/>
+ </partitioning>
+ <directory>%%java.io.tmpdir%%</directory>
+ <prefix>out</prefix>
+ <cache_size>500</cache_size>
+ <main>ids</main>
+ <compare>
+ <condition>
+ <negated>N</negated>
+ <operator>-</operator>
+ <leftvalue/>
+ <function>=</function>
+ <rightvalue/>
+ </condition>
+ </compare>
+ <attributes/>
+ <GUI>
+ <xloc>400</xloc>
+ <yloc>192</yloc>
+ </GUI>
+ </transform>
+ <transform>
+ <name>Output</name>
+ <type>TextFileOutput</type>
+ <description/>
+ <distribute>Y</distribute>
+ <custom_distribution/>
+ <copies>1</copies>
+ <partitioning>
+ <method>none</method>
+ <schema_name/>
+ </partitioning>
+ <separator>;</separator>
+ <enclosure/>
+ <enclosure_forced>N</enclosure_forced>
+ <enclosure_fix_disabled>N</enclosure_fix_disabled>
+ <header>Y</header>
+ <footer>N</footer>
+ <format>UNIX</format>
+ <compression>None</compression>
+ <encoding>UTF-8</encoding>
+ <endedLine/>
+ <fileNameInField>N</fileNameInField>
+ <fileNameField/>
+ <create_parent_folder>Y</create_parent_folder>
+ <file>
+
<name>${java.io.tmpdir}/hop-it-single-threaded/0003-join-rows-three-inputs</name>
+ <servlet_output>N</servlet_output>
+ <do_not_open_new_file_init>N</do_not_open_new_file_init>
+ <extention>csv</extention>
+ <append>N</append>
+ <split>N</split>
+ <haspartno>N</haspartno>
+ <add_date>N</add_date>
+ <add_time>N</add_time>
+ <SpecifyFormat>N</SpecifyFormat>
+ <date_time_format/>
+ <add_to_result_filenames>N</add_to_result_filenames>
+ <pad>N</pad>
+ <fast_dump>N</fast_dump>
+ <splitevery/>
+ </file>
+ <fields>
+ <field>
+ <name>id</name>
+ <type>Integer</type>
+ <format>0</format>
+ <currency/>
+ <decimal/>
+ <group/>
+ <nullif/>
+ <trim_type>none</trim_type>
+ <length>-1</length>
+ <precision>-1</precision>
+ </field>
+ <field>
+ <name>name</name>
+ <type>String</type>
+ <format></format>
+ <currency/>
+ <decimal/>
+ <group/>
+ <nullif/>
+ <trim_type>none</trim_type>
+ <length>-1</length>
+ <precision>-1</precision>
+ </field>
+ <field>
+ <name>code</name>
+ <type>String</type>
+ <format></format>
+ <currency/>
+ <decimal/>
+ <group/>
+ <nullif/>
+ <trim_type>none</trim_type>
+ <length>-1</length>
+ <precision>-1</precision>
+ </field>
+ </fields>
+ <attributes/>
+ <GUI>
+ <xloc>640</xloc>
+ <yloc>192</yloc>
+ </GUI>
+ </transform>
+ <transform_error_handling>
+ </transform_error_handling>
+ <attributes/>
+</pipeline>
diff --git a/integration-tests/single_threaded/dev-env-config.json
b/integration-tests/single_threaded/dev-env-config.json
new file mode 100644
index 0000000000..c4c624dc16
--- /dev/null
+++ b/integration-tests/single_threaded/dev-env-config.json
@@ -0,0 +1,10 @@
+{
+ "description": "Single threaded pipeline engine integration tests - DEV",
+ "metadataBaseFolder": "${PROJECT_HOME}/metadata",
+ "unitTestsBasePath": "${PROJECT_HOME}",
+ "dataSetsCsvFolder": "${PROJECT_HOME}/datasets",
+ "enforcingExecutionInHome": true,
+ "config": {
+ "variables": []
+ }
+}
diff --git a/integration-tests/single_threaded/hop-config.json
b/integration-tests/single_threaded/hop-config.json
new file mode 100644
index 0000000000..d9e1e6562e
--- /dev/null
+++ b/integration-tests/single_threaded/hop-config.json
@@ -0,0 +1,290 @@
+{
+ "variables": [
+ {
+ "name": "HOP_LENIENT_STRING_TO_NUMBER_CONVERSION",
+ "value": "N",
+ "description": "System wide flag to allow lenient string to number
conversion for backward compatibility. If this setting is set to \"Y\", an
string starting with digits will be converted successfully into a number.
(example: 192.168.1.1 will be converted into 192 or 192.168 or 192168 depending
on the decimal and grouping symbol). The default (N) will be to throw an error
if non-numeric symbols are found in the string."
+ },
+ {
+ "name": "HOP_COMPATIBILITY_DB_IGNORE_TIMEZONE",
+ "value": "N",
+ "description": "System wide flag to ignore timezone while writing
date/timestamp value to the database."
+ },
+ {
+ "name": "HOP_LOG_SIZE_LIMIT",
+ "value": "0",
+ "description": "The log size limit for all pipelines and workflows that
don't have the \"log size limit\" property set in their respective properties."
+ },
+ {
+ "name": "HOP_EMPTY_STRING_DIFFERS_FROM_NULL",
+ "value": "N",
+ "description": "NULL vs Empty String. If this setting is set to Y, an
empty string and null are different. Otherwise they are not."
+ },
+ {
+ "name": "HOP_MAX_LOG_SIZE_IN_LINES",
+ "value": "0",
+ "description": "The maximum number of log lines that are kept internally
by Hop. Set to 0 to keep all rows (default)"
+ },
+ {
+ "name": "HOP_MAX_LOG_TIMEOUT_IN_MINUTES",
+ "value": "1440",
+ "description": "The maximum age (in minutes) of a log line while being
kept internally by Hop. Set to 0 to keep all rows indefinitely (default)"
+ },
+ {
+ "name": "HOP_MAX_WORKFLOW_TRACKER_SIZE",
+ "value": "5000",
+ "description": "The maximum number of workflow trackers kept in memory"
+ },
+ {
+ "name": "HOP_MAX_ACTIONS_LOGGED",
+ "value": "5000",
+ "description": "The maximum number of action results kept in memory for
logging purposes."
+ },
+ {
+ "name": "HOP_MAX_LOGGING_REGISTRY_SIZE",
+ "value": "10000",
+ "description": "The maximum number of logging registry entries kept in
memory for logging purposes."
+ },
+ {
+ "name": "HOP_LOG_TAB_REFRESH_DELAY",
+ "value": "1000",
+ "description": "The hop log tab refresh delay."
+ },
+ {
+ "name": "HOP_LOG_TAB_REFRESH_PERIOD",
+ "value": "1000",
+ "description": "The hop log tab refresh period."
+ },
+ {
+ "name": "HOP_PLUGIN_CLASSES",
+ "value": null,
+ "description": "A comma delimited list of classes to scan for plugin
annotations"
+ },
+ {
+ "name": "HOP_PLUGIN_PACKAGES",
+ "value": null,
+ "description": "A comma delimited list of packages to scan for plugin
annotations (warning: slow!!)"
+ },
+ {
+ "name": "HOP_TRANSFORM_PERFORMANCE_SNAPSHOT_LIMIT",
+ "value": "0",
+ "description": "The maximum number of transform performance snapshots to
keep in memory. Set to 0 to keep all snapshots indefinitely (default)"
+ },
+ {
+ "name": "HOP_ROWSET_GET_TIMEOUT",
+ "value": "50",
+ "description": "The name of the variable that optionally contains an
alternative rowset get timeout (in ms). This only makes a difference for
extremely short lived pipelines."
+ },
+ {
+ "name": "HOP_ROWSET_PUT_TIMEOUT",
+ "value": "50",
+ "description": "The name of the variable that optionally contains an
alternative rowset put timeout (in ms). This only makes a difference for
extremely short lived pipelines."
+ },
+ {
+ "name": "HOP_CORE_TRANSFORMS_FILE",
+ "value": null,
+ "description": "The name of the project variable that will contain the
alternative location of the hop-transforms.xml file. You can use this to
customize the list of available internal transforms outside of the codebase."
+ },
+ {
+ "name": "HOP_CORE_WORKFLOW_ACTIONS_FILE",
+ "value": null,
+ "description": "The name of the project variable that will contain the
alternative location of the hop-workflow-actions.xml file."
+ },
+ {
+ "name": "HOP_SERVER_OBJECT_TIMEOUT_MINUTES",
+ "value": "1440",
+ "description": "This project variable will set a time-out after which
waiting, completed or stopped pipelines and workflows will be automatically
cleaned up. The default value is 1440 (one day)."
+ },
+ {
+ "name": "HOP_PIPELINE_PAN_JVM_EXIT_CODE",
+ "value": null,
+ "description": "Set this variable to an integer that will be returned as
the Pan JVM exit code."
+ },
+ {
+ "name": "HOP_DISABLE_CONSOLE_LOGGING",
+ "value": "N",
+ "description": "Set this variable to Y to disable standard Hop logging
to the console. (stdout)"
+ },
+ {
+ "name": "HOP_REDIRECT_STDERR",
+ "value": "N",
+ "description": "Set this variable to Y to redirect stderr to Hop
logging."
+ },
+ {
+ "name": "HOP_REDIRECT_STDOUT",
+ "value": "N",
+ "description": "Set this variable to Y to redirect stdout to Hop
logging."
+ },
+ {
+ "name": "HOP_DEFAULT_NUMBER_FORMAT",
+ "value": null,
+ "description": "The name of the variable containing an alternative
default number format"
+ },
+ {
+ "name": "HOP_DEFAULT_BIGNUMBER_FORMAT",
+ "value": null,
+ "description": "The name of the variable containing an alternative
default bignumber format"
+ },
+ {
+ "name": "HOP_DEFAULT_INTEGER_FORMAT",
+ "value": null,
+ "description": "The name of the variable containing an alternative
default integer format"
+ },
+ {
+ "name": "HOP_DEFAULT_DATE_FORMAT",
+ "value": null,
+ "description": "The name of the variable containing an alternative
default date format"
+ },
+ {
+ "name": "HOP_DEFAULT_TIMESTAMP_FORMAT",
+ "value": null,
+ "description": "The name of the variable containing an alternative
default timestamp format"
+ },
+ {
+ "name": "HOP_DEFAULT_SERVLET_ENCODING",
+ "value": null,
+ "description": "Defines the default encoding for servlets, leave it
empty to use Java default encoding"
+ },
+ {
+ "name": "HOP_FAIL_ON_LOGGING_ERROR",
+ "value": "N",
+ "description": "Set this variable to Y when you want the
workflow/pipeline fail with an error when the related logging process (e.g. to
a database) fails."
+ },
+ {
+ "name": "HOP_AGGREGATION_MIN_NULL_IS_VALUED",
+ "value": "N",
+ "description": "Set this variable to Y to set the minimum to NULL if
NULL is within an aggregate. Otherwise by default NULL is ignored by the MIN
aggregate and MIN is set to the minimum value that is not NULL. See also the
variable HOP_AGGREGATION_ALL_NULLS_ARE_ZERO."
+ },
+ {
+ "name": "HOP_AGGREGATION_ALL_NULLS_ARE_ZERO",
+ "value": "N",
+ "description": "Set this variable to Y to return 0 when all values
within an aggregate are NULL. Otherwise by default a NULL is returned when all
values are NULL."
+ },
+ {
+ "name": "HOP_COMPATIBILITY_TEXT_FILE_OUTPUT_APPEND_NO_HEADER",
+ "value": "N",
+ "description": "Set this variable to Y for backward compatibility for
the Text File Output transform. Setting this to Ywill add no header row at all
when the append option is enabled, regardless if the file is existing or not."
+ },
+ {
+ "name": "HOP_PASSWORD_ENCODER_PLUGIN",
+ "value": "Hop",
+ "description": "Specifies the password encoder plugin to use by ID (Hop
is the default)."
+ },
+ {
+ "name": "HOP_SYSTEM_HOSTNAME",
+ "value": null,
+ "description": "You can use this variable to speed up hostname lookup.
Hostname lookup is performed by Hop so that it is capable of logging the server
on which a workflow or pipeline is executed."
+ },
+ {
+ "name": "HOP_SERVER_JETTY_ACCEPTORS",
+ "value": null,
+ "description": "A variable to configure jetty option: acceptors for
Carte"
+ },
+ {
+ "name": "HOP_SERVER_JETTY_ACCEPT_QUEUE_SIZE",
+ "value": null,
+ "description": "A variable to configure jetty option: acceptQueueSize
for Carte"
+ },
+ {
+ "name": "HOP_SERVER_JETTY_RES_MAX_IDLE_TIME",
+ "value": null,
+ "description": "A variable to configure jetty option:
lowResourcesMaxIdleTime for Carte"
+ },
+ {
+ "name":
"HOP_COMPATIBILITY_MERGE_ROWS_USE_REFERENCE_STREAM_WHEN_IDENTICAL",
+ "value": "N",
+ "description": "Set this variable to Y for backward compatibility for
the Merge Rows (diff) transform. Setting this to Y will use the data from the
reference stream (instead of the comparison stream) in case the compared rows
are identical."
+ },
+ {
+ "name": "HOP_SPLIT_FIELDS_REMOVE_ENCLOSURE",
+ "value": "false",
+ "description": "Set this variable to false to preserve enclosure symbol
after splitting the string in the Split fields transform. Changing it to true
will remove first and last enclosure symbol from the resulting string chunks."
+ },
+ {
+ "name": "HOP_ALLOW_EMPTY_FIELD_NAMES_AND_TYPES",
+ "value": "false",
+ "description": "Set this variable to TRUE to allow your pipeline to pass
'null' fields and/or empty types."
+ },
+ {
+ "name": "HOP_GLOBAL_LOG_VARIABLES_CLEAR_ON_EXPORT",
+ "value": "false",
+ "description": "Set this variable to false to preserve global log
variables defined in pipeline / workflow Properties -> Log panel. Changing it
to true will clear it when export pipeline / workflow."
+ },
+ {
+ "name": "HOP_FILE_OUTPUT_MAX_STREAM_COUNT",
+ "value": "1024",
+ "description": "This project variable is used by the Text File Output
transform. It defines the max number of simultaneously open files within the
transform. The transform will close/reopen files as necessary to insure the max
is not exceeded"
+ },
+ {
+ "name": "HOP_FILE_OUTPUT_MAX_STREAM_LIFE",
+ "value": "0",
+ "description": "This project variable is used by the Text File Output
transform. It defines the max number of milliseconds between flushes of files
opened by the transform."
+ },
+ {
+ "name": "HOP_USE_NATIVE_FILE_DIALOG",
+ "value": "N",
+ "description": "Set this value to Y if you want to use the system file
open/save dialog when browsing files"
+ },
+ {
+ "name": "HOP_AUTO_CREATE_CONFIG",
+ "value": "Y",
+ "description": "Set this value to N if you don't want to automatically
create a hop configuration file (hop-config.json) when it's missing"
+ }
+ ],
+ "LocaleDefault": "en_BE",
+ "guiProperties": {
+ "FontFixedSize": "13",
+ "MaxUndo": "100",
+ "DarkMode": "Y",
+ "FontNoteSize": "13",
+ "ShowOSLook": "Y",
+ "FontFixedStyle": "0",
+ "FontNoteName": ".AppleSystemUIFont",
+ "FontFixedName": "Monospaced",
+ "FontGraphStyle": "0",
+ "FontDefaultSize": "13",
+ "GraphColorR": "255",
+ "FontGraphSize": "13",
+ "IconSize": "32",
+ "BackgroundColorB": "255",
+ "FontNoteStyle": "0",
+ "FontGraphName": ".AppleSystemUIFont",
+ "FontDefaultName": ".AppleSystemUIFont",
+ "GraphColorG": "255",
+ "UseGlobalFileBookmarks": "Y",
+ "FontDefaultStyle": "0",
+ "GraphColorB": "255",
+ "BackgroundColorR": "255",
+ "BackgroundColorG": "255",
+ "WorkflowDialogStyle": "RESIZE,MAX,MIN",
+ "LineWidth": "1",
+ "ContextDialogShowCategories": "Y"
+ },
+ "projectsConfig": {
+ "enabled": true,
+ "projectMandatory": true,
+ "environmentMandatory": false,
+ "defaultProject": "default",
+ "defaultEnvironment": null,
+ "standardParentProject": "default",
+ "standardProjectsFolder": null,
+ "projectConfigurations": [
+ {
+ "projectName": "default",
+ "projectHome": "${HOP_CONFIG_FOLDER}",
+ "configFilename": "project-config.json"
+ }
+ ],
+ "lifecycleEnvironments": [
+ {
+ "name": "dev",
+ "purpose": "Testing",
+ "projectName": "default",
+ "configurationFiles": [
+ "${PROJECT_HOME}/dev-env-config.json"
+ ]
+ }
+ ],
+ "projectLifecycles": []
+ }
+}
\ No newline at end of file
diff --git a/integration-tests/single_threaded/main-0001-join-rows.hwf
b/integration-tests/single_threaded/main-0001-join-rows.hwf
new file mode 100644
index 0000000000..d4273951e4
--- /dev/null
+++ b/integration-tests/single_threaded/main-0001-join-rows.hwf
@@ -0,0 +1,138 @@
+<?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.
+
+-->
+<workflow>
+ <name>main-0001-join-rows</name>
+ <name_sync_with_filename>Y</name_sync_with_filename>
+ <description/>
+ <extended_description/>
+ <workflow_version/>
+ <created_user>-</created_user>
+ <created_date>2026/09/24 12:00:00.000</created_date>
+ <modified_user>-</modified_user>
+ <modified_date>2026/09/24 12:00:00.000</modified_date>
+ <parameters>
+ </parameters>
+ <actions>
+ <action>
+ <name>Start</name>
+ <description/>
+ <type>SPECIAL</type>
+ <attributes/>
+ <start>Y</start>
+ <dummy>N</dummy>
+ <repeat>N</repeat>
+ <schedulerType>0</schedulerType>
+ <intervalSeconds>0</intervalSeconds>
+ <intervalMinutes>60</intervalMinutes>
+ <hour>12</hour>
+ <minutes>0</minutes>
+ <weekDay>1</weekDay>
+ <DayOfMonth>1</DayOfMonth>
+ <parallel>N</parallel>
+ <nr>0</nr>
+ <xloc>96</xloc>
+ <yloc>96</yloc>
+ <attributes_hac/>
+ </action>
+ <action>
+ <name>0001-join-rows.hpl</name>
+ <description/>
+ <type>PIPELINE</type>
+ <attributes/>
+ <filename>${PROJECT_HOME}/0001-join-rows.hpl</filename>
+ <params_from_previous>N</params_from_previous>
+ <exec_per_row>N</exec_per_row>
+ <clear_rows>N</clear_rows>
+ <clear_files>N</clear_files>
+ <set_logfile>N</set_logfile>
+ <logfile/>
+ <logext/>
+ <add_date>N</add_date>
+ <add_time>N</add_time>
+ <loglevel>Basic</loglevel>
+ <set_append_logfile>N</set_append_logfile>
+ <wait_until_finished>Y</wait_until_finished>
+ <follow_abort_remote>N</follow_abort_remote>
+ <create_parent_folder>N</create_parent_folder>
+ <run_configuration>single-threaded</run_configuration>
+ <parameters>
+ <pass_all_parameters>Y</pass_all_parameters>
+ </parameters>
+ <parallel>N</parallel>
+ <nr>0</nr>
+ <xloc>320</xloc>
+ <yloc>96</yloc>
+ <attributes_hac/>
+ </action>
+ <action>
+ <name>0001-join-rows-validation.hpl</name>
+ <description/>
+ <type>PIPELINE</type>
+ <attributes/>
+ <filename>${PROJECT_HOME}/0001-join-rows-validation.hpl</filename>
+ <params_from_previous>N</params_from_previous>
+ <exec_per_row>N</exec_per_row>
+ <clear_rows>N</clear_rows>
+ <clear_files>N</clear_files>
+ <set_logfile>N</set_logfile>
+ <logfile/>
+ <logext/>
+ <add_date>N</add_date>
+ <add_time>N</add_time>
+ <loglevel>Basic</loglevel>
+ <set_append_logfile>N</set_append_logfile>
+ <wait_until_finished>Y</wait_until_finished>
+ <follow_abort_remote>N</follow_abort_remote>
+ <create_parent_folder>N</create_parent_folder>
+ <run_configuration>local</run_configuration>
+ <parameters>
+ <pass_all_parameters>Y</pass_all_parameters>
+ </parameters>
+ <parallel>N</parallel>
+ <nr>0</nr>
+ <xloc>608</xloc>
+ <yloc>96</yloc>
+ <attributes_hac/>
+ </action>
+ </actions>
+ <hops>
+ <hop>
+ <from>Start</from>
+ <to>0001-join-rows.hpl</to>
+ <from_nr>0</from_nr>
+ <to_nr>0</to_nr>
+ <enabled>Y</enabled>
+ <evaluation>Y</evaluation>
+ <unconditional>Y</unconditional>
+ </hop>
+ <hop>
+ <from>0001-join-rows.hpl</from>
+ <to>0001-join-rows-validation.hpl</to>
+ <from_nr>0</from_nr>
+ <to_nr>0</to_nr>
+ <enabled>Y</enabled>
+ <evaluation>Y</evaluation>
+ <unconditional>N</unconditional>
+ </hop>
+ </hops>
+ <notepads>
+ </notepads>
+ <attributes/>
+</workflow>
diff --git
a/integration-tests/single_threaded/main-0002-join-rows-after-transforms.hwf
b/integration-tests/single_threaded/main-0002-join-rows-after-transforms.hwf
new file mode 100644
index 0000000000..cba93a062f
--- /dev/null
+++ b/integration-tests/single_threaded/main-0002-join-rows-after-transforms.hwf
@@ -0,0 +1,138 @@
+<?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.
+
+-->
+<workflow>
+ <name>main-0002-join-rows-after-transforms</name>
+ <name_sync_with_filename>Y</name_sync_with_filename>
+ <description/>
+ <extended_description/>
+ <workflow_version/>
+ <created_user>-</created_user>
+ <created_date>2026/09/24 12:00:00.000</created_date>
+ <modified_user>-</modified_user>
+ <modified_date>2026/09/24 12:00:00.000</modified_date>
+ <parameters>
+ </parameters>
+ <actions>
+ <action>
+ <name>Start</name>
+ <description/>
+ <type>SPECIAL</type>
+ <attributes/>
+ <start>Y</start>
+ <dummy>N</dummy>
+ <repeat>N</repeat>
+ <schedulerType>0</schedulerType>
+ <intervalSeconds>0</intervalSeconds>
+ <intervalMinutes>60</intervalMinutes>
+ <hour>12</hour>
+ <minutes>0</minutes>
+ <weekDay>1</weekDay>
+ <DayOfMonth>1</DayOfMonth>
+ <parallel>N</parallel>
+ <nr>0</nr>
+ <xloc>96</xloc>
+ <yloc>96</yloc>
+ <attributes_hac/>
+ </action>
+ <action>
+ <name>0002-join-rows-after-transforms.hpl</name>
+ <description/>
+ <type>PIPELINE</type>
+ <attributes/>
+ <filename>${PROJECT_HOME}/0002-join-rows-after-transforms.hpl</filename>
+ <params_from_previous>N</params_from_previous>
+ <exec_per_row>N</exec_per_row>
+ <clear_rows>N</clear_rows>
+ <clear_files>N</clear_files>
+ <set_logfile>N</set_logfile>
+ <logfile/>
+ <logext/>
+ <add_date>N</add_date>
+ <add_time>N</add_time>
+ <loglevel>Basic</loglevel>
+ <set_append_logfile>N</set_append_logfile>
+ <wait_until_finished>Y</wait_until_finished>
+ <follow_abort_remote>N</follow_abort_remote>
+ <create_parent_folder>N</create_parent_folder>
+ <run_configuration>single-threaded</run_configuration>
+ <parameters>
+ <pass_all_parameters>Y</pass_all_parameters>
+ </parameters>
+ <parallel>N</parallel>
+ <nr>0</nr>
+ <xloc>320</xloc>
+ <yloc>96</yloc>
+ <attributes_hac/>
+ </action>
+ <action>
+ <name>0002-join-rows-after-transforms-validation.hpl</name>
+ <description/>
+ <type>PIPELINE</type>
+ <attributes/>
+
<filename>${PROJECT_HOME}/0002-join-rows-after-transforms-validation.hpl</filename>
+ <params_from_previous>N</params_from_previous>
+ <exec_per_row>N</exec_per_row>
+ <clear_rows>N</clear_rows>
+ <clear_files>N</clear_files>
+ <set_logfile>N</set_logfile>
+ <logfile/>
+ <logext/>
+ <add_date>N</add_date>
+ <add_time>N</add_time>
+ <loglevel>Basic</loglevel>
+ <set_append_logfile>N</set_append_logfile>
+ <wait_until_finished>Y</wait_until_finished>
+ <follow_abort_remote>N</follow_abort_remote>
+ <create_parent_folder>N</create_parent_folder>
+ <run_configuration>local</run_configuration>
+ <parameters>
+ <pass_all_parameters>Y</pass_all_parameters>
+ </parameters>
+ <parallel>N</parallel>
+ <nr>0</nr>
+ <xloc>608</xloc>
+ <yloc>96</yloc>
+ <attributes_hac/>
+ </action>
+ </actions>
+ <hops>
+ <hop>
+ <from>Start</from>
+ <to>0002-join-rows-after-transforms.hpl</to>
+ <from_nr>0</from_nr>
+ <to_nr>0</to_nr>
+ <enabled>Y</enabled>
+ <evaluation>Y</evaluation>
+ <unconditional>Y</unconditional>
+ </hop>
+ <hop>
+ <from>0002-join-rows-after-transforms.hpl</from>
+ <to>0002-join-rows-after-transforms-validation.hpl</to>
+ <from_nr>0</from_nr>
+ <to_nr>0</to_nr>
+ <enabled>Y</enabled>
+ <evaluation>Y</evaluation>
+ <unconditional>N</unconditional>
+ </hop>
+ </hops>
+ <notepads>
+ </notepads>
+ <attributes/>
+</workflow>
diff --git
a/integration-tests/single_threaded/main-0003-join-rows-three-inputs.hwf
b/integration-tests/single_threaded/main-0003-join-rows-three-inputs.hwf
new file mode 100644
index 0000000000..c5ceeec75b
--- /dev/null
+++ b/integration-tests/single_threaded/main-0003-join-rows-three-inputs.hwf
@@ -0,0 +1,138 @@
+<?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.
+
+-->
+<workflow>
+ <name>main-0003-join-rows-three-inputs</name>
+ <name_sync_with_filename>Y</name_sync_with_filename>
+ <description/>
+ <extended_description/>
+ <workflow_version/>
+ <created_user>-</created_user>
+ <created_date>2026/09/24 12:00:00.000</created_date>
+ <modified_user>-</modified_user>
+ <modified_date>2026/09/24 12:00:00.000</modified_date>
+ <parameters>
+ </parameters>
+ <actions>
+ <action>
+ <name>Start</name>
+ <description/>
+ <type>SPECIAL</type>
+ <attributes/>
+ <start>Y</start>
+ <dummy>N</dummy>
+ <repeat>N</repeat>
+ <schedulerType>0</schedulerType>
+ <intervalSeconds>0</intervalSeconds>
+ <intervalMinutes>60</intervalMinutes>
+ <hour>12</hour>
+ <minutes>0</minutes>
+ <weekDay>1</weekDay>
+ <DayOfMonth>1</DayOfMonth>
+ <parallel>N</parallel>
+ <nr>0</nr>
+ <xloc>96</xloc>
+ <yloc>96</yloc>
+ <attributes_hac/>
+ </action>
+ <action>
+ <name>0003-join-rows-three-inputs.hpl</name>
+ <description/>
+ <type>PIPELINE</type>
+ <attributes/>
+ <filename>${PROJECT_HOME}/0003-join-rows-three-inputs.hpl</filename>
+ <params_from_previous>N</params_from_previous>
+ <exec_per_row>N</exec_per_row>
+ <clear_rows>N</clear_rows>
+ <clear_files>N</clear_files>
+ <set_logfile>N</set_logfile>
+ <logfile/>
+ <logext/>
+ <add_date>N</add_date>
+ <add_time>N</add_time>
+ <loglevel>Basic</loglevel>
+ <set_append_logfile>N</set_append_logfile>
+ <wait_until_finished>Y</wait_until_finished>
+ <follow_abort_remote>N</follow_abort_remote>
+ <create_parent_folder>N</create_parent_folder>
+ <run_configuration>single-threaded</run_configuration>
+ <parameters>
+ <pass_all_parameters>Y</pass_all_parameters>
+ </parameters>
+ <parallel>N</parallel>
+ <nr>0</nr>
+ <xloc>320</xloc>
+ <yloc>96</yloc>
+ <attributes_hac/>
+ </action>
+ <action>
+ <name>0003-join-rows-three-inputs-validation.hpl</name>
+ <description/>
+ <type>PIPELINE</type>
+ <attributes/>
+
<filename>${PROJECT_HOME}/0003-join-rows-three-inputs-validation.hpl</filename>
+ <params_from_previous>N</params_from_previous>
+ <exec_per_row>N</exec_per_row>
+ <clear_rows>N</clear_rows>
+ <clear_files>N</clear_files>
+ <set_logfile>N</set_logfile>
+ <logfile/>
+ <logext/>
+ <add_date>N</add_date>
+ <add_time>N</add_time>
+ <loglevel>Basic</loglevel>
+ <set_append_logfile>N</set_append_logfile>
+ <wait_until_finished>Y</wait_until_finished>
+ <follow_abort_remote>N</follow_abort_remote>
+ <create_parent_folder>N</create_parent_folder>
+ <run_configuration>local</run_configuration>
+ <parameters>
+ <pass_all_parameters>Y</pass_all_parameters>
+ </parameters>
+ <parallel>N</parallel>
+ <nr>0</nr>
+ <xloc>608</xloc>
+ <yloc>96</yloc>
+ <attributes_hac/>
+ </action>
+ </actions>
+ <hops>
+ <hop>
+ <from>Start</from>
+ <to>0003-join-rows-three-inputs.hpl</to>
+ <from_nr>0</from_nr>
+ <to_nr>0</to_nr>
+ <enabled>Y</enabled>
+ <evaluation>Y</evaluation>
+ <unconditional>Y</unconditional>
+ </hop>
+ <hop>
+ <from>0003-join-rows-three-inputs.hpl</from>
+ <to>0003-join-rows-three-inputs-validation.hpl</to>
+ <from_nr>0</from_nr>
+ <to_nr>0</to_nr>
+ <enabled>Y</enabled>
+ <evaluation>Y</evaluation>
+ <unconditional>N</unconditional>
+ </hop>
+ </hops>
+ <notepads>
+ </notepads>
+ <attributes/>
+</workflow>
diff --git
a/integration-tests/single_threaded/metadata/pipeline-run-configuration/local.json
b/integration-tests/single_threaded/metadata/pipeline-run-configuration/local.json
new file mode 100644
index 0000000000..63794efcaf
--- /dev/null
+++
b/integration-tests/single_threaded/metadata/pipeline-run-configuration/local.json
@@ -0,0 +1,17 @@
+{
+ "engineRunConfiguration": {
+ "Local": {
+ "feedback_size": "50000",
+ "sample_size": "100",
+ "sample_type_in_gui": "Last",
+ "rowset_size": "10000",
+ "safe_mode": false,
+ "show_feedback": false,
+ "topo_sort": false,
+ "gather_metrics": false
+ }
+ },
+ "configurationVariables": [],
+ "name": "local",
+ "description": "Runs your pipelines locally with the standard local Hop
pipeline engine"
+}
\ No newline at end of file
diff --git
a/integration-tests/single_threaded/metadata/pipeline-run-configuration/single-threaded.json
b/integration-tests/single_threaded/metadata/pipeline-run-configuration/single-threaded.json
new file mode 100644
index 0000000000..36ab0c30fe
--- /dev/null
+++
b/integration-tests/single_threaded/metadata/pipeline-run-configuration/single-threaded.json
@@ -0,0 +1,11 @@
+{
+ "engineRunConfiguration": {
+ "LocalSingle": {
+ "sample_type_in_gui": "Last",
+ "sample_size": "100"
+ }
+ },
+ "configurationVariables": [],
+ "name": "single-threaded",
+ "description": "Runs your pipelines locally with the single threaded
pipeline engine"
+}
diff --git
a/integration-tests/single_threaded/metadata/workflow-run-configuration/local.json
b/integration-tests/single_threaded/metadata/workflow-run-configuration/local.json
new file mode 100644
index 0000000000..e37a93039a
--- /dev/null
+++
b/integration-tests/single_threaded/metadata/workflow-run-configuration/local.json
@@ -0,0 +1,9 @@
+{
+ "engineRunConfiguration": {
+ "Local": {
+ "safe_mode": false
+ }
+ },
+ "name": "local",
+ "description": "Runs your workflows locally with the standard local Hop
workflow engine"
+}
\ No newline at end of file
diff --git a/integration-tests/single_threaded/project-config.json
b/integration-tests/single_threaded/project-config.json
new file mode 100644
index 0000000000..899c6c1092
--- /dev/null
+++ b/integration-tests/single_threaded/project-config.json
@@ -0,0 +1,15 @@
+{
+ "metadataBaseFolder": "${PROJECT_HOME}/metadata",
+ "unitTestsBasePath": "${PROJECT_HOME}",
+ "dataSetsCsvFolder": "${PROJECT_HOME}/datasets",
+ "enforcingExecutionInHome": true,
+ "config": {
+ "variables": [
+ {
+ "name": "HOP_LICENSE_HEADER_FILE",
+ "value": "${PROJECT_HOME}/../asf-header.txt",
+ "description": "This will automatically serialize the ASF license
header into pipelines and workflows in the integration test projects"
+ }
+ ]
+ }
+}
\ No newline at end of file
diff --git a/plugins/engines/beam/pom.xml b/plugins/engines/beam/pom.xml
index ee0aa70aeb..902ee473ab 100644
--- a/plugins/engines/beam/pom.xml
+++ b/plugins/engines/beam/pom.xml
@@ -1605,6 +1605,12 @@
<version>${project.version}</version>
<scope>provided</scope>
</dependency>
+ <dependency>
+ <groupId>org.apache.hop</groupId>
+ <artifactId>hop-transform-joinrows</artifactId>
+ <version>${project.version}</version>
+ <scope>provided</scope>
+ </dependency>
<dependency>
<groupId>org.apache.hop</groupId>
<artifactId>hop-transform-memgroupby</artifactId>
diff --git
a/plugins/engines/beam/src/main/java/org/apache/hop/beam/pipeline/HopPipelineMetaToBeamPipelineConverter.java
b/plugins/engines/beam/src/main/java/org/apache/hop/beam/pipeline/HopPipelineMetaToBeamPipelineConverter.java
index 0253b28aef..4e8b2a43d4 100644
---
a/plugins/engines/beam/src/main/java/org/apache/hop/beam/pipeline/HopPipelineMetaToBeamPipelineConverter.java
+++
b/plugins/engines/beam/src/main/java/org/apache/hop/beam/pipeline/HopPipelineMetaToBeamPipelineConverter.java
@@ -66,6 +66,7 @@ import
org.apache.hop.pipeline.config.PipelineRunConfiguration;
import org.apache.hop.pipeline.transform.ITransformMeta;
import org.apache.hop.pipeline.transform.TransformMeta;
import org.apache.hop.pipeline.transforms.groupby.GroupByMeta;
+import org.apache.hop.pipeline.transforms.joinrows.JoinRowsMeta;
import org.apache.hop.pipeline.transforms.sort.SortRowsMeta;
import org.apache.hop.pipeline.transforms.uniquerows.UniqueRowsMeta;
import
org.apache.hop.pipeline.transforms.uniquerowsbyhashset.UniqueRowsByHashSetMeta;
@@ -103,7 +104,9 @@ public class HopPipelineMetaToBeamPipelineConverter {
SortRowsMeta.class,
"Sort Rows is not supported on Beam. A Beam pipeline re-shuffles
rows across workers to maximize parallelism, so this transform would only order
the rows a single worker happens to hold, not the data set as a whole.",
UniqueRowsByHashSetMeta.class,
- "Unique Rows By Hashset is not supported on Beam. Every worker
keeps its own hash set, so duplicates spread over different workers would
survive. Use a Memory Group By to get distinct rows.");
+ "Unique Rows By Hashset is not supported on Beam. Every worker
keeps its own hash set, so duplicates spread over different workers would
survive. Use a Memory Group By to get distinct rows.",
+ JoinRowsMeta.class,
+ "Join Rows is not supported on Beam. A cartesian product needs
every row of every input in one place, but every worker would only combine the
rows it happens to hold, so combinations would go missing. Add the same
constant field to both inputs and use a Merge Join on that field instead.");
protected final String runConfigName;
protected final PipelineRunConfiguration runConfiguration;
diff --git
a/plugins/engines/beam/src/test/java/org/apache/hop/beam/engines/BeamPipelineEngineSupportsTest.java
b/plugins/engines/beam/src/test/java/org/apache/hop/beam/engines/BeamPipelineEngineSupportsTest.java
index 59cf072568..323c0deb14 100644
---
a/plugins/engines/beam/src/test/java/org/apache/hop/beam/engines/BeamPipelineEngineSupportsTest.java
+++
b/plugins/engines/beam/src/test/java/org/apache/hop/beam/engines/BeamPipelineEngineSupportsTest.java
@@ -31,6 +31,7 @@ import org.apache.hop.core.plugins.EngineCompatibility;
import org.apache.hop.core.plugins.IPlugin;
import org.apache.hop.pipeline.transform.BaseTransformMeta;
import org.apache.hop.pipeline.transforms.groupby.GroupByMeta;
+import org.apache.hop.pipeline.transforms.joinrows.JoinRowsMeta;
import org.apache.hop.pipeline.transforms.sort.SortRowsMeta;
import org.apache.hop.pipeline.transforms.uniquerows.UniqueRowsMeta;
import
org.apache.hop.pipeline.transforms.uniquerowsbyhashset.UniqueRowsByHashSetMeta;
@@ -71,6 +72,15 @@ class BeamPipelineEngineSupportsTest {
assertTrue(engine.supports(pluginWithMainType(UniqueRowsByHashSetMeta.class)).isUnsupported());
}
+ @Test
+ void joinRowsMetaIsHardBanned() {
+ EngineCompatibility verdict =
engine.supports(pluginWithMainType(JoinRowsMeta.class));
+ assertTrue(verdict.isUnsupported(), "JoinRowsMeta should be UNSUPPORTED");
+ assertEquals(
+
HopPipelineMetaToBeamPipelineConverter.HARD_BANNED_META_TYPES.get(JoinRowsMeta.class),
+ verdict.getReason());
+ }
+
@Test
void mainTypeImplementingBeamHandlerIsSupported() {
EngineCompatibility verdict =
engine.supports(pluginWithMainType(StubBeamMeta.class));
diff --git a/plugins/engines/spark/pom.xml b/plugins/engines/spark/pom.xml
index 7c99a7a871..6c70c2576b 100644
--- a/plugins/engines/spark/pom.xml
+++ b/plugins/engines/spark/pom.xml
@@ -229,6 +229,12 @@
<version>${project.version}</version>
<scope>test</scope>
</dependency>
+ <dependency>
+ <groupId>org.apache.hop</groupId>
+ <artifactId>hop-transform-joinrows</artifactId>
+ <version>${project.version}</version>
+ <scope>test</scope>
+ </dependency>
<!-- Stream Lookup for info-stream unit tests (mapPartitions +
broadcast) -->
<dependency>
<groupId>org.apache.hop</groupId>
diff --git
a/plugins/engines/spark/src/main/java/org/apache/hop/spark/pipeline/HopPipelineMetaToSparkConverter.java
b/plugins/engines/spark/src/main/java/org/apache/hop/spark/pipeline/HopPipelineMetaToSparkConverter.java
index 72c548e299..427a3b4e6d 100644
---
a/plugins/engines/spark/src/main/java/org/apache/hop/spark/pipeline/HopPipelineMetaToSparkConverter.java
+++
b/plugins/engines/spark/src/main/java/org/apache/hop/spark/pipeline/HopPipelineMetaToSparkConverter.java
@@ -94,7 +94,9 @@ public class HopPipelineMetaToSparkConverter {
SparkConst.GROUP_BY_PLUGIN_ID,
"Group By is not supported on the native Spark engine. Use Memory
Group By (native Spark shuffle) instead, or run on the Local engine.",
SparkConst.UNIQUE_ROWS_BY_HASH_SET_PLUGIN_ID,
- "Unique Rows By Hashset is not supported on the native Spark engine.
Every partition would keep its own hash set, so duplicates spread over
different partitions would survive. Use Unique Rows (native Spark distinct) or
Memory Group By instead.");
+ "Unique Rows By Hashset is not supported on the native Spark engine.
Every partition would keep its own hash set, so duplicates spread over
different partitions would survive. Use Unique Rows (native Spark distinct) or
Memory Group By instead.",
+ SparkConst.JOIN_ROWS_PLUGIN_ID,
+ "Join Rows is not supported on the native Spark engine. A cartesian
product needs every row of every input in one place, but every partition would
only combine the rows it happens to hold, so combinations would go missing. Add
the same constant field to both inputs and use Merge Join on that field
instead.");
private final IVariables variables;
private final PipelineMeta pipelineMeta;
diff --git
a/plugins/engines/spark/src/main/java/org/apache/hop/spark/util/SparkConst.java
b/plugins/engines/spark/src/main/java/org/apache/hop/spark/util/SparkConst.java
index 83c533e923..9e3620a165 100644
---
a/plugins/engines/spark/src/main/java/org/apache/hop/spark/util/SparkConst.java
+++
b/plugins/engines/spark/src/main/java/org/apache/hop/spark/util/SparkConst.java
@@ -38,6 +38,7 @@ public final class SparkConst {
public static final String UNIQUE_ROWS_BY_HASH_SET_PLUGIN_ID =
"UniqueRowsByHashSet";
public static final String SORT_ROWS_PLUGIN_ID = "SortRows";
public static final String GROUP_BY_PLUGIN_ID = "GroupBy";
+ public static final String JOIN_ROWS_PLUGIN_ID = "JoinRows";
public static final String SPARK_FILE_INPUT_PLUGIN_ID = "SparkFileInput";
public static final String SPARK_FILE_OUTPUT_PLUGIN_ID = "SparkFileOutput";
diff --git
a/plugins/engines/spark/src/test/java/org/apache/hop/spark/pipeline/HopPipelineMetaToSparkConverterTest.java
b/plugins/engines/spark/src/test/java/org/apache/hop/spark/pipeline/HopPipelineMetaToSparkConverterTest.java
index 2eca1106b0..ccf07ad49f 100644
---
a/plugins/engines/spark/src/test/java/org/apache/hop/spark/pipeline/HopPipelineMetaToSparkConverterTest.java
+++
b/plugins/engines/spark/src/test/java/org/apache/hop/spark/pipeline/HopPipelineMetaToSparkConverterTest.java
@@ -38,6 +38,7 @@ import org.apache.hop.pipeline.PipelineMeta;
import org.apache.hop.pipeline.transform.TransformMeta;
import org.apache.hop.pipeline.transforms.dummy.DummyMeta;
import org.apache.hop.pipeline.transforms.groupby.GroupByMeta;
+import org.apache.hop.pipeline.transforms.joinrows.JoinRowsMeta;
import
org.apache.hop.pipeline.transforms.uniquerowsbyhashset.UniqueRowsByHashSetMeta;
import org.apache.hop.spark.engines.SparkPipelineEngine;
import org.apache.hop.spark.util.SparkConst;
@@ -56,6 +57,17 @@ class HopPipelineMetaToSparkConverterTest {
assertTrue(group.getMessage().contains("Group By"));
}
+ @Test
+ void joinRowsIsBanned() {
+ HopException join =
+ assertThrows(
+ HopException.class,
+ () ->
+ HopPipelineMetaToSparkConverter.validateTransformSparkUsage(
+ SparkConst.JOIN_ROWS_PLUGIN_ID, "j1"));
+ assertTrue(join.getMessage().contains("Merge Join"));
+ }
+
@Test
void nativeHandlersAreNotBanned() throws HopException {
HopPipelineMetaToSparkConverter.validateTransformSparkUsage(
@@ -166,7 +178,9 @@ class HopPipelineMetaToSparkConverterTest {
SparkConst.GROUP_BY_PLUGIN_ID,
GroupByMeta.class,
SparkConst.UNIQUE_ROWS_BY_HASH_SET_PLUGIN_ID,
- UniqueRowsByHashSetMeta.class);
+ UniqueRowsByHashSetMeta.class,
+ SparkConst.JOIN_ROWS_PLUGIN_ID,
+ JoinRowsMeta.class);
assertEquals(
HopPipelineMetaToSparkConverter.HARD_BANNED_PLUGIN_IDS.keySet(),
diff --git
a/plugins/transforms/joinrows/src/main/java/org/apache/hop/pipeline/transforms/joinrows/JoinRows.java
b/plugins/transforms/joinrows/src/main/java/org/apache/hop/pipeline/transforms/joinrows/JoinRows.java
index 418b4fb4a3..8db4e80edd 100644
---
a/plugins/transforms/joinrows/src/main/java/org/apache/hop/pipeline/transforms/joinrows/JoinRows.java
+++
b/plugins/transforms/joinrows/src/main/java/org/apache/hop/pipeline/transforms/joinrows/JoinRows.java
@@ -73,7 +73,7 @@ public class JoinRows extends BaseTransform<JoinRowsMeta,
JoinRowsData> {
// See if a main transform is supplied: in that case move the
corresponding rowset to position
// 0
- swapFirstInputRowSetIfExists(meta.getMainTransformName());
+ moveInputRowSetToFront(meta.getMainTransformName());
List<IRowSet> inputRowSets = getInputRowSets();
int rowSetsSize = inputRowSets.size();
@@ -133,7 +133,7 @@ public class JoinRows extends BaseTransform<JoinRowsMeta,
JoinRowsData> {
if (filenr == 0) {
// Rowset 0:
IRowSet rowSet = getFirstInputRowSet();
- rowData = getRowFrom(rowSet);
+ rowData = readRow(rowSet);
if (rowData != null) {
data.fileRowMeta[0] = rowSet.getRowMeta();
}
@@ -262,15 +262,23 @@ public class JoinRows extends BaseTransform<JoinRowsMeta,
JoinRowsData> {
initialize();
}
- if (data.caching) {
+ // Like the lookup rows of a Stream Lookup: read all the rows of the other
input streams first.
+ //
+ while (data.caching) {
if (!cacheInputRow()) {
return false;
}
- } else {
+ }
+
+ // Read one row of the main stream and write all its combinations with the
other streams.
+ // We're back at the main stream (filenr 0) once they are all written.
+ //
+ do {
if (!outputRow()) {
return false;
}
- }
+ } while (data.filenr != 0 && !isStopped());
+
return true;
}
@@ -287,12 +295,7 @@ public class JoinRows extends BaseTransform<JoinRowsMeta,
JoinRowsData> {
// Before we exit we need to make sure the 100 rows in the other streams
are consumed
// though...
//
- while (getRow() != null) {
- // Consume
- if (isStopped()) {
- break;
- }
- }
+ consumeRemainingInput();
setOutputDone();
return false;
@@ -381,7 +384,7 @@ public class JoinRows extends BaseTransform<JoinRowsMeta,
JoinRowsData> {
// Read a line from the appropriate rowset...
IRowSet rowSet = data.rs[data.filenr];
- Object[] rowData = getRowFrom(rowSet);
+ Object[] rowData = readRow(rowSet);
if (rowData != null) {
// We read a row from one of the input streams...
@@ -462,6 +465,56 @@ public class JoinRows extends BaseTransform<JoinRowsMeta,
JoinRowsData> {
return outputRowMeta;
}
+ /**
+ * Move the row set of the given transform to position 0, and keep the other
row sets in hop
+ * order. That's the order in which {@link JoinRowsMeta#getFields} lists
their fields. Swapping it
+ * with position 0 would mix up the fields of the other streams when there
are more than two.
+ */
+ private void moveInputRowSetToFront(String transformName) {
+ List<IRowSet> rowSets = getInputRowSets();
+ for (int i = 1; i < rowSets.size(); i++) {
+ if
(rowSets.get(i).getOriginTransformName().equalsIgnoreCase(transformName)) {
+ rowSets.add(0, rowSets.remove(i));
+ setInputRowSets(rowSets);
+ return;
+ }
+ }
+ }
+
+ /**
+ * Read a row from one of the input row sets. The single threaded executor
runs the previous
+ * transforms first, so the rows waiting on the input row sets are all there
is. It doesn't flag
+ * those row sets as done though: an empty row set marks the end of that
input, reading on would
+ * wait forever.
+ */
+ private Object[] readRow(IRowSet rowSet) throws HopException {
+ if (isSingleThreaded() && rowSet.size() == 0) {
+ return null;
+ }
+ return getRowFrom(rowSet);
+ }
+
+ private void consumeRemainingInput() throws HopException {
+ if (isSingleThreaded()) {
+ for (IRowSet rowSet : new ArrayList<>(getInputRowSets())) {
+ while (!isStopped() && readRow(rowSet) != null) {
+ // Consume
+ }
+ }
+ } else {
+ while (getRow() != null) {
+ // Consume
+ if (isStopped()) {
+ break;
+ }
+ }
+ }
+ }
+
+ private boolean isSingleThreaded() {
+ return getPipeline().getPipelineType() ==
PipelineMeta.PipelineType.SingleThreaded;
+ }
+
@Override
public void dispose() {
@@ -476,26 +529,4 @@ public class JoinRows extends BaseTransform<JoinRowsMeta,
JoinRowsData> {
super.dispose();
}
-
- @Override
- public void batchComplete() throws HopException {
- IRowSet rowSet = getFirstInputRowSet();
- int repeats = 0;
- for (int i = 0; i < data.cache.length; i++) {
- if (repeats == 0) {
- repeats = 1;
- }
- if (data.cache[i] != null) {
- repeats *= data.cache[i].size();
- }
- }
- while (rowSet.size() > 0 && !isStopped()) {
- init();
- }
- // The last row needs to be written too to the account of the number of
input rows.
- //
- for (int i = 0; i < repeats; i++) {
- init();
- }
- }
}
diff --git
a/plugins/transforms/joinrows/src/main/java/org/apache/hop/pipeline/transforms/joinrows/JoinRowsMeta.java
b/plugins/transforms/joinrows/src/main/java/org/apache/hop/pipeline/transforms/joinrows/JoinRowsMeta.java
index 08e6266860..ffbb1b4102 100644
---
a/plugins/transforms/joinrows/src/main/java/org/apache/hop/pipeline/transforms/joinrows/JoinRowsMeta.java
+++
b/plugins/transforms/joinrows/src/main/java/org/apache/hop/pipeline/transforms/joinrows/JoinRowsMeta.java
@@ -52,7 +52,8 @@ import org.apache.hop.pipeline.transform.stream.StreamIcon;
description = "i18n::BaseTransform.TypeTooltipDesc.JoinRows",
categoryDescription =
"i18n:org.apache.hop.pipeline.transform:BaseTransform.Category.Joins",
keywords = "i18n::JoinRowsMeta.keyword",
- documentationUrl = "/pipeline/transforms/joinrows.html")
+ documentationUrl = "/pipeline/transforms/joinrows.html",
+ excludedEngines = {"Beam*", "SparkPipelineEngine"})
@Getter
@Setter
public class JoinRowsMeta extends BaseTransformMeta<JoinRows, JoinRowsData> {
diff --git
a/plugins/transforms/joinrows/src/test/java/org/apache/hop/pipeline/transforms/joinrows/JoinRowsSingleThreadedTest.java
b/plugins/transforms/joinrows/src/test/java/org/apache/hop/pipeline/transforms/joinrows/JoinRowsSingleThreadedTest.java
new file mode 100644
index 0000000000..af51ce2e35
--- /dev/null
+++
b/plugins/transforms/joinrows/src/test/java/org/apache/hop/pipeline/transforms/joinrows/JoinRowsSingleThreadedTest.java
@@ -0,0 +1,260 @@
+/*
+ * 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.transforms.joinrows;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertTimeoutPreemptively;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+
+import java.time.Duration;
+import java.util.ArrayList;
+import java.util.HashMap;
+import java.util.List;
+import java.util.Map;
+import org.apache.hop.core.HopEnvironment;
+import org.apache.hop.core.RowMetaAndData;
+import org.apache.hop.core.annotations.Transform;
+import org.apache.hop.core.plugins.PluginRegistry;
+import org.apache.hop.core.plugins.TransformPluginType;
+import org.apache.hop.core.row.IRowMeta;
+import org.apache.hop.core.row.RowMeta;
+import org.apache.hop.core.row.value.ValueMetaString;
+import org.apache.hop.pipeline.Pipeline;
+import org.apache.hop.pipeline.PipelineHopMeta;
+import org.apache.hop.pipeline.PipelineMeta;
+import org.apache.hop.pipeline.RowProducer;
+import org.apache.hop.pipeline.SingleThreadedPipelineExecutor;
+import org.apache.hop.pipeline.engines.local.LocalPipelineEngine;
+import org.apache.hop.pipeline.transform.ITransform;
+import org.apache.hop.pipeline.transform.ITransformMeta;
+import org.apache.hop.pipeline.transform.TransformMeta;
+import org.apache.hop.pipeline.transforms.dummy.DummyMeta;
+import org.apache.hop.pipeline.transforms.injector.InjectorMeta;
+import org.junit.jupiter.api.BeforeAll;
+import org.junit.jupiter.api.Test;
+
+/**
+ * Runs Join Rows through the {@link SingleThreadedPipelineExecutor}: the
executor behind the single
+ * threaded pipeline engine and single threaded sub-pipelines (#2353).
+ *
+ * <p>The executor calls processRow() once per row waiting on the input row
sets, and for as long as
+ * the main stream (an info stream of Join Rows) holds rows. The rows on the
input row sets are the
+ * complete input, but the row sets are not flagged as done, so any read
beyond them would wait
+ * forever.
+ */
+class JoinRowsSingleThreadedTest {
+
+ private static final Duration TIMEOUT = Duration.ofSeconds(20);
+
+ @BeforeAll
+ static void initHop() throws Exception {
+ HopEnvironment.init();
+ PluginRegistry registry = PluginRegistry.getInstance();
+ for (ITransformMeta meta : List.of(new JoinRowsMeta(), new InjectorMeta(),
new DummyMeta())) {
+ if (registry.getPluginId(TransformPluginType.class, meta) == null) {
+ registry.registerPluginClass(
+ meta.getClass().getName(), TransformPluginType.class,
Transform.class);
+ }
+ }
+ }
+
+ @Test
+ void cartesianProductOfTwoInputs() throws Exception {
+ List<Input> inputs = List.of(new Input("main", 3), new Input("other", 2));
+ List<RowMetaAndData> rows =
+ assertTimeoutPreemptively(TIMEOUT, () -> runJoin(inputs, "main", 500,
false));
+
+ assertEquals(6, rows.size());
+ assertProduct(rows, inputs);
+ }
+
+ @Test
+ void cartesianProductWhenInputRowSetsAreFinished() throws Exception {
+ List<Input> inputs = List.of(new Input("main", 3), new Input("other", 2));
+ List<RowMetaAndData> rows =
+ assertTimeoutPreemptively(TIMEOUT, () -> runJoin(inputs, "main", 500,
true));
+
+ assertEquals(6, rows.size());
+ assertProduct(rows, inputs);
+ }
+
+ @Test
+ void cartesianProductSpillingToTemporaryFile() throws Exception {
+ // A cache size below the number of rows makes the transform read the rows
back from disk.
+ List<Input> inputs = List.of(new Input("main", 4), new Input("other", 5));
+ List<RowMetaAndData> rows =
+ assertTimeoutPreemptively(TIMEOUT, () -> runJoin(inputs, "main", 2,
false));
+
+ assertEquals(20, rows.size());
+ assertProduct(rows, inputs);
+ }
+
+ @Test
+ void cartesianProductOfThreeInputsWithMainStreamLast() throws Exception {
+ // The main stream is the last hop into Join Rows: its row set has to move
to the front while
+ // the other streams keep their hop order, the order in which
JoinRowsMeta.getFields() lists
+ // their fields.
+ List<Input> inputs = List.of(new Input("a", 2), new Input("b", 3), new
Input("main", 2));
+ List<RowMetaAndData> rows =
+ assertTimeoutPreemptively(TIMEOUT, () -> runJoin(inputs, "main", 500,
false));
+
+ assertEquals(12, rows.size());
+ assertProduct(rows, List.of(inputs.get(2), inputs.get(0), inputs.get(1)));
+ }
+
+ @Test
+ void cartesianProductOfThreeInputsWithMainStreamInTheMiddle() throws
Exception {
+ List<Input> inputs = List.of(new Input("a", 2), new Input("main", 3), new
Input("b", 2));
+ List<RowMetaAndData> rows =
+ assertTimeoutPreemptively(TIMEOUT, () -> runJoin(inputs, "main", 1,
false));
+
+ assertEquals(12, rows.size());
+ assertProduct(rows, List.of(inputs.get(1), inputs.get(0), inputs.get(2)));
+ }
+
+ @Test
+ void noOutputWhenOneInputIsEmpty() throws Exception {
+ List<Input> inputs = List.of(new Input("main", 3), new Input("other", 0));
+ List<RowMetaAndData> rows =
+ assertTimeoutPreemptively(TIMEOUT, () -> runJoin(inputs, "main", 500,
false));
+
+ assertTrue(rows.isEmpty());
+ }
+
+ @Test
+ void singleInputPassesRowsThrough() throws Exception {
+ // This is how the Beam engine used to feed the transform: all inputs
flattened into one row
+ // set, one row per iteration. The old batchComplete() hung on the row
that was left behind.
+ List<Input> inputs = List.of(new Input("main", 3));
+ List<RowMetaAndData> rows =
+ assertTimeoutPreemptively(TIMEOUT, () -> runJoin(inputs, "main", 500,
false));
+
+ assertEquals(3, rows.size());
+ assertProduct(rows, inputs);
+ }
+
+ /**
+ * An input of Join Rows: an injector with a single String field named after
the input, holding
+ * the values {@code <name>1 .. <name><nrRows>}.
+ */
+ private record Input(String name, int nrRows) {}
+
+ /**
+ * Asserts the complete cartesian product, with the fields in the given
order of the inputs. The
+ * values of the last input vary fastest.
+ */
+ private static void assertProduct(List<RowMetaAndData> rows, List<Input>
fieldOrder) {
+ List<String> expectedFields =
fieldOrder.stream().map(Input::name).toList();
+ List<String> expected = List.of("");
+ for (Input input : fieldOrder) {
+ List<String> combined = new ArrayList<>();
+ for (String prefix : expected) {
+ for (int n = 1; n <= input.nrRows(); n++) {
+ combined.add(prefix + (prefix.isEmpty() ? "" : "|") + input.name() +
n);
+ }
+ }
+ expected = combined;
+ }
+
+ List<String> actual = new ArrayList<>();
+ for (RowMetaAndData row : rows) {
+ assertEquals(expectedFields, List.of(row.getRowMeta().getFieldNames()));
+ List<String> values = new ArrayList<>();
+ for (int i = 0; i < row.size(); i++) {
+ values.add((String) row.getData()[i]);
+ }
+ actual.add(String.join("|", values));
+ }
+ assertEquals(expected, actual);
+ }
+
+ private static List<RowMetaAndData> runJoin(
+ List<Input> inputs, String mainName, int cacheSize, boolean
finishInputs) throws Exception {
+ PipelineMeta pipelineMeta = new PipelineMeta();
+ pipelineMeta.setName("join-rows-single-threaded");
+
+ List<TransformMeta> inputTransforms = new ArrayList<>();
+ for (Input input : inputs) {
+ inputTransforms.add(addTransform(pipelineMeta, input.name(), new
InjectorMeta()));
+ }
+
+ JoinRowsMeta joinRowsMeta = new JoinRowsMeta();
+ joinRowsMeta.setDefault();
+ joinRowsMeta.setCacheSize(cacheSize);
+ joinRowsMeta.setMainTransformName(mainName);
+ joinRowsMeta.setDirectory(System.getProperty("java.io.tmpdir"));
+ TransformMeta join = addTransform(pipelineMeta, "join", joinRowsMeta);
+
+ TransformMeta output = addTransform(pipelineMeta, "output", new
DummyMeta());
+
+ for (TransformMeta inputTransform : inputTransforms) {
+ pipelineMeta.addPipelineHop(new PipelineHopMeta(inputTransform, join));
+ }
+ pipelineMeta.addPipelineHop(new PipelineHopMeta(join, output));
+
+ // Like loading the pipeline from a file: the main transform is an info
stream of Join Rows,
+ // which the executor feeds with processRow() calls for as long as it
holds rows.
+ joinRowsMeta.searchInfoAndTargetTransforms(pipelineMeta.getTransforms());
+
+ Pipeline pipeline = new LocalPipelineEngine(pipelineMeta);
+ pipeline.setPipelineType(PipelineMeta.PipelineType.SingleThreaded);
+ pipeline.prepareExecution();
+
+ Map<String, RowProducer> producers = new HashMap<>();
+ for (Input input : inputs) {
+ producers.put(input.name(), pipeline.addRowProducer(input.name(), 0));
+ }
+
+ ITransform joinTransform = pipeline.getTransform("join", 0);
+ TransformRowsCollector collector = new TransformRowsCollector();
+ joinTransform.addRowListener(collector);
+
+ pipeline.startThreads();
+
+ for (Input input : inputs) {
+ IRowMeta rowMeta = new RowMeta();
+ rowMeta.addValueMeta(new ValueMetaString(input.name()));
+ RowProducer producer = producers.get(input.name());
+ for (int n = 1; n <= input.nrRows(); n++) {
+ producer.putRow(rowMeta, new Object[] {input.name() + n});
+ }
+ }
+ if (finishInputs) {
+ producers.values().forEach(RowProducer::finished);
+ }
+
+ SingleThreadedPipelineExecutor executor = new
SingleThreadedPipelineExecutor(pipeline);
+ assertTrue(executor.init());
+ try {
+ executor.oneIteration();
+ } finally {
+ executor.dispose();
+ }
+ assertEquals(0, joinTransform.getErrors());
+
+ return collector.getRowsWritten();
+ }
+
+ private static TransformMeta addTransform(
+ PipelineMeta pipelineMeta, String name, ITransformMeta meta) {
+ String pluginId =
PluginRegistry.getInstance().getPluginId(TransformPluginType.class, meta);
+ TransformMeta transformMeta = new TransformMeta(pluginId, name, meta);
+ pipelineMeta.addTransform(transformMeta);
+ return transformMeta;
+ }
+}