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 9a7f60f407 Fixes #3360 : add context about using Group By in
distributed environments (#8510)
9a7f60f407 is described below
commit 9a7f60f4071c52949787660673644eb769e2f308
Author: Bart Maertens <[email protected]>
AuthorDate: Mon Sep 21 18:54:11 2026 +0200
Fixes #3360 : add context about using Group By in distributed environments
(#8510)
Group By claimed Beam support in its docs while both the Beam engines and
the native Spark engine refuse it. Correct the page-engines attribute and
explain, on the transform pages themselves, why order-dependent transforms
cannot work on an engine that re-shuffles rows across workers.
Sort Rows and Unique Rows (HashSet) were documented as unsupported on Beam
but ran through the generic handler per bundle, and Unique Rows (HashSet)
did the same per partition on native Spark. Hard-ban them so the palette
filter and run gate refuse them instead of returning silently wrong results.
---
.../pages/pipeline/beam/getting-started-with-beam.adoc | 9 +++++++--
.../modules/ROOT/pages/pipeline/transforms/groupby.adoc | 16 +++++++++++++++-
.../ROOT/pages/pipeline/transforms/memgroupby.adoc | 8 ++++++++
.../modules/ROOT/pages/pipeline/transforms/sort.adoc | 12 ++++++++++++
.../pages/pipeline/transforms/uniquerowsbyhashset.adoc | 11 +++++++++++
plugins/engines/beam/pom.xml | 6 ++++++
.../pipeline/HopPipelineMetaToBeamPipelineConverter.java | 8 +++++++-
.../hop/beam/engines/BeamPipelineEngineSupportsTest.java | 12 ++++++++++++
plugins/engines/spark/pom.xml | 8 +++++++-
.../spark/pipeline/HopPipelineMetaToSparkConverter.java | 4 +++-
.../main/java/org/apache/hop/spark/util/SparkConst.java | 1 +
.../pipeline/HopPipelineMetaToSparkConverterTest.java | 8 +++++++-
.../uniquerowsbyhashset/UniqueRowsByHashSetMeta.java | 4 +++-
13 files changed, 99 insertions(+), 8 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 1c93658037..0104fa65f4 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
@@ -183,12 +183,17 @@ When using the Beam engines it uses
`org.apache.beam.sdk.io.synthetic.SyntheticB
[#_unsupported_transforms]
=== Unsupported transforms
-A few transforms are simply not supported because we haven't found a good way
to do this on Beam yet:
+A few transforms are simply not supported because we haven't found a good way
to do this on Beam yet.
+They all depend on seeing the rows in a particular order, or on holding state
for the whole data set in one place, and a Beam pipeline re-shuffles rows
across workers to maximize parallelism.
+Each of these transforms would therefore only ever see the slice of the data
one worker happens to hold, and would quietly produce a wrong result, so the
Beam engines refuse them with an error instead:
-* xref:pipeline/transforms/uniquerows.adoc[Unique Rows]
+* xref:pipeline/transforms/uniquerows.adoc[Unique Rows] : Use the `Memory
Group By` instead
+* 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]
+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.
+
The xref:pipeline/transforms/rowdenormaliser.adoc[Denormaliser] transform
works technically correct on Apache Beam in release 1.1.0 and later.
Even so you need to consider that the aggregation of the key-value pairs in
that transform (in the general case) only happens on a sub-set of the rows.
That is because in a Beam pipeline the order in which rows arrive is lost
because they are continuously re-shuffled to maximize parallelism.
diff --git
a/docs/hop-user-manual/modules/ROOT/pages/pipeline/transforms/groupby.adoc
b/docs/hop-user-manual/modules/ROOT/pages/pipeline/transforms/groupby.adoc
index 219064e535..2252420db0 100644
--- a/docs/hop-user-manual/modules/ROOT/pages/pipeline/transforms/groupby.adoc
+++ b/docs/hop-user-manual/modules/ROOT/pages/pipeline/transforms/groupby.adoc
@@ -17,7 +17,7 @@ under the License.
:documentationPath: /pipeline/transforms/
:language: en_US
:description: The Group By transform groups rows from a source, based on a
specified field or collection of fields. A new row is generated for each group.
-:page-engines: Hop Engine=yes, Single Threaded=yes, Native Spark=no, Beam
Spark=yes, Beam Flink=yes, Beam Dataflow=yes
+:page-engines: Hop Engine=yes, Single Threaded=yes, Native Spark=no, Beam
Spark=no, Beam Flink=no, Beam Dataflow=no
= image:transforms/icons/groupby.svg[Group By transform Icon,
role="image-doc-icon"] Group By
@@ -37,6 +37,20 @@ If you sort the data outside of Hop, the case sensitivity of
the data in the fie
You can use the xref:pipeline/transforms/memgroupby.adoc[Memory Group By]
transform to handle non-sorted input.
+[WARNING]
+====
+Group By does not run on the distributed pipeline engines.
+Both the Beam engines and the native Spark engine refuse the transform and
stop the pipeline with an error.
+
+The reason is the sorted input this transform relies on.
+A distributed engine spreads the rows over many workers and re-shuffles them
to maximize parallelism, so a sort upstream only ever orders the rows one
worker happens to hold.
+The groups would be formed per worker and the aggregates would be silently
wrong.
+
+Use xref:pipeline/transforms/memgroupby.adoc[Memory Group By] instead.
+It needs no sorted input, and both engines map it onto their own shuffle
(`GroupByKey` on Beam) so the grouping covers the whole data set.
+See
xref:pipeline/beam/getting-started-with-beam.adoc#_unsupported_transforms[Getting
started with Beam] for the other transforms in this situation.
+====
+
== Options
[options="header"]
diff --git
a/docs/hop-user-manual/modules/ROOT/pages/pipeline/transforms/memgroupby.adoc
b/docs/hop-user-manual/modules/ROOT/pages/pipeline/transforms/memgroupby.adoc
index 4ae7c9c8d6..45e036ca6c 100644
---
a/docs/hop-user-manual/modules/ROOT/pages/pipeline/transforms/memgroupby.adoc
+++
b/docs/hop-user-manual/modules/ROOT/pages/pipeline/transforms/memgroupby.adoc
@@ -30,6 +30,14 @@ However, it **does** require all data to fit into memory.
TIP: When the number of rows is too large to fit into memory, use a
combination of xref:pipeline/transforms/sort.adoc[Sort Rows] and
xref:pipeline/transforms/groupby.adoc[Group By] transforms.
+[NOTE]
+====
+That fallback to xref:pipeline/transforms/sort.adoc[Sort Rows] and
xref:pipeline/transforms/groupby.adoc[Group By] only applies to the Hop engines.
+Neither of those transforms runs on the Beam engines or on the native Spark
engine.
+
+On those engines Memory Group By is the aggregation to reach for: it is
translated into the engine's own grouping operation
(`org.apache.beam.sdk.transforms.GroupByKey` on Beam, a Spark shuffle on native
Spark) rather than run as a Hop transform per worker, so neither the sorting
nor the memory limit described above applies.
+====
+
== Options
[options="header"]
diff --git
a/docs/hop-user-manual/modules/ROOT/pages/pipeline/transforms/sort.adoc
b/docs/hop-user-manual/modules/ROOT/pages/pipeline/transforms/sort.adoc
index aecbadd6d0..0e9ab79147 100644
--- a/docs/hop-user-manual/modules/ROOT/pages/pipeline/transforms/sort.adoc
+++ b/docs/hop-user-manual/modules/ROOT/pages/pipeline/transforms/sort.adoc
@@ -29,6 +29,18 @@ The transform optionally only pass unique records, based on
the sort keys.
TIP: You use can use multiple copies of the Sort Rows transform to speed up
large sort operations. Make sure to add a
xref:pipeline/transforms/sortedmerge.adoc[Sorted Merge] transform after the
sort to correctly merge the streams of sorted that are generated by the various
sort copies.
+[WARNING]
+====
+Sort Rows is not supported on the Beam engines.
+A Beam pipeline spreads rows over many workers and re-shuffles them to
maximize parallelism, so this transform can only order the rows a single worker
happens to hold at that moment, never the data set as a whole.
+
+If you need sorted input for a downstream transform, that transform will not
work on Beam either.
+Use an alternative that does not depend on row order, such as
xref:pipeline/transforms/memgroupby.adoc[Memory Group By] in place of
xref:pipeline/transforms/groupby.adoc[Group By].
+See
xref:pipeline/beam/getting-started-with-beam.adoc#_unsupported_transforms[Getting
started with Beam].
+
+The native Spark engine does support this transform: it is mapped onto a Spark
sort across the full data set instead of being run per partition.
+====
+
== Options
[options="header"]
diff --git
a/docs/hop-user-manual/modules/ROOT/pages/pipeline/transforms/uniquerowsbyhashset.adoc
b/docs/hop-user-manual/modules/ROOT/pages/pipeline/transforms/uniquerowsbyhashset.adoc
index 9098a8b228..fb45f70fa5 100644
---
a/docs/hop-user-manual/modules/ROOT/pages/pipeline/transforms/uniquerowsbyhashset.adoc
+++
b/docs/hop-user-manual/modules/ROOT/pages/pipeline/transforms/uniquerowsbyhashset.adoc
@@ -29,6 +29,17 @@ This transform differs from the
xref:pipeline/transforms/uniquerows.adoc[Unique
WARNING: When using this transform on larger sets of data there is a risk of
hash collisions unless "Compare using stored row values" is enabled.
+[WARNING]
+====
+This transform is not supported on the distributed pipeline engines: neither
the Beam engines nor the native Spark engine.
+The hash set it uses to spot duplicates lives in the memory of a single
worker, and these engines spread the rows over many workers, so duplicates that
end up on different workers would survive.
+Both refuse the transform with an error rather than hand you that incomplete
result.
+
+Use a xref:pipeline/transforms/memgroupby.adoc[Memory Group By] to get
distinct rows instead: it is translated into the engine's own grouping
operation and covers the whole data set.
+On the native Spark engine xref:pipeline/transforms/uniquerows.adoc[Unique
Rows] is an option too, because that engine maps it onto its own distinct
operation; on the Beam engines it is not, as it is refused there as well.
+See
xref:pipeline/beam/getting-started-with-beam.adoc#_unsupported_transforms[Getting
started with Beam].
+====
+
== Options
[options="header"]
diff --git a/plugins/engines/beam/pom.xml b/plugins/engines/beam/pom.xml
index cb065352f6..ee0aa70aeb 100644
--- a/plugins/engines/beam/pom.xml
+++ b/plugins/engines/beam/pom.xml
@@ -1647,6 +1647,12 @@
<version>${project.version}</version>
<scope>provided</scope>
</dependency>
+ <dependency>
+ <groupId>org.apache.hop</groupId>
+ <artifactId>hop-transform-uniquerowsbyhashset</artifactId>
+ <version>${project.version}</version>
+ <scope>provided</scope>
+ </dependency>
<dependency>
<groupId>org.apache.hop</groupId>
<artifactId>hop-transform-writetolog</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 edf25b0bac..0253b28aef 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,7 +66,9 @@ 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.sort.SortRowsMeta;
import org.apache.hop.pipeline.transforms.uniquerows.UniqueRowsMeta;
+import
org.apache.hop.pipeline.transforms.uniquerowsbyhashset.UniqueRowsByHashSetMeta;
import org.jboss.jandex.AnnotationInstance;
import org.jboss.jandex.ClassInfo;
import org.jboss.jandex.IndexView;
@@ -97,7 +99,11 @@ public class HopPipelineMetaToBeamPipelineConverter {
GroupByMeta.class,
"Group By is not supported. Use the Memory Group By transform
instead. It comes closest to Beam functionality.",
UniqueRowsMeta.class,
- "The unique rows transform is not yet supported on Beam, for now use
a Memory Group By to get distrinct rows");
+ "The unique rows transform is not yet supported on Beam, for now use
a Memory Group By to get distrinct rows",
+ 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.");
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 b9e40b1db1..59cf072568 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,7 +31,9 @@ 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.sort.SortRowsMeta;
import org.apache.hop.pipeline.transforms.uniquerows.UniqueRowsMeta;
+import
org.apache.hop.pipeline.transforms.uniquerowsbyhashset.UniqueRowsByHashSetMeta;
import org.junit.jupiter.api.Test;
/**
@@ -59,6 +61,16 @@ class BeamPipelineEngineSupportsTest {
assertTrue(engine.supports(pluginWithMainType(UniqueRowsMeta.class)).isUnsupported());
}
+ @Test
+ void sortRowsMetaIsHardBanned() {
+
assertTrue(engine.supports(pluginWithMainType(SortRowsMeta.class)).isUnsupported());
+ }
+
+ @Test
+ void uniqueRowsByHashSetMetaIsHardBanned() {
+
assertTrue(engine.supports(pluginWithMainType(UniqueRowsByHashSetMeta.class)).isUnsupported());
+ }
+
@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 fb80af8ef2..7c99a7a871 100644
--- a/plugins/engines/spark/pom.xml
+++ b/plugins/engines/spark/pom.xml
@@ -222,7 +222,7 @@
<version>${project.version}</version>
<scope>test</scope>
</dependency>
- <!-- Hard-banned transform meta for annotation lockstep tests only -->
+ <!-- Hard-banned transform metas for annotation lockstep tests only -->
<dependency>
<groupId>org.apache.hop</groupId>
<artifactId>hop-transform-groupby</artifactId>
@@ -243,6 +243,12 @@
<version>${project.version}</version>
<scope>test</scope>
</dependency>
+ <dependency>
+ <groupId>org.apache.hop</groupId>
+ <artifactId>hop-transform-uniquerowsbyhashset</artifactId>
+ <version>${project.version}</version>
+ <scope>test</scope>
+ </dependency>
</dependencies>
<build>
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 60bcf454dc..72c548e299 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
@@ -92,7 +92,9 @@ public class HopPipelineMetaToSparkConverter {
public static final Map<String, String> HARD_BANNED_PLUGIN_IDS =
Map.of(
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 Local/Beam.");
+ "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.");
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 1a55480d62..83c533e923 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
@@ -35,6 +35,7 @@ public final class SparkConst {
public static final String MEMORY_GROUP_BY_PLUGIN_ID = "MemoryGroupBy";
public static final String MERGE_JOIN_PLUGIN_ID = "MergeJoin";
public static final String UNIQUE_ROWS_PLUGIN_ID = "Unique";
+ 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";
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 8a6a2815e4..2eca1106b0 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.uniquerowsbyhashset.UniqueRowsByHashSetMeta;
import org.apache.hop.spark.engines.SparkPipelineEngine;
import org.apache.hop.spark.util.SparkConst;
import org.junit.jupiter.api.Test;
@@ -160,7 +161,12 @@ class HopPipelineMetaToSparkConverterTest {
*/
@Test
void hardBannedTransformsExcludeNativeSparkOnAnnotation() {
- Map<String, Class<?>> bannedMetas = Map.of(SparkConst.GROUP_BY_PLUGIN_ID,
GroupByMeta.class);
+ Map<String, Class<?>> bannedMetas =
+ Map.of(
+ SparkConst.GROUP_BY_PLUGIN_ID,
+ GroupByMeta.class,
+ SparkConst.UNIQUE_ROWS_BY_HASH_SET_PLUGIN_ID,
+ UniqueRowsByHashSetMeta.class);
assertEquals(
HopPipelineMetaToSparkConverter.HARD_BANNED_PLUGIN_IDS.keySet(),
diff --git
a/plugins/transforms/uniquerowsbyhashset/src/main/java/org/apache/hop/pipeline/transforms/uniquerowsbyhashset/UniqueRowsByHashSetMeta.java
b/plugins/transforms/uniquerowsbyhashset/src/main/java/org/apache/hop/pipeline/transforms/uniquerowsbyhashset/UniqueRowsByHashSetMeta.java
index 09c420b1bd..3bb5863ae7 100644
---
a/plugins/transforms/uniquerowsbyhashset/src/main/java/org/apache/hop/pipeline/transforms/uniquerowsbyhashset/UniqueRowsByHashSetMeta.java
+++
b/plugins/transforms/uniquerowsbyhashset/src/main/java/org/apache/hop/pipeline/transforms/uniquerowsbyhashset/UniqueRowsByHashSetMeta.java
@@ -40,7 +40,9 @@ import org.apache.hop.pipeline.transform.TransformMeta;
description = "i18n::UniqueRowsByHashSet.Description",
categoryDescription =
"i18n:org.apache.hop.pipeline.transform:BaseTransform.Category.Transform",
keywords = "i18n::UniqueRowsByHashSetMeta.keyword",
- documentationUrl = "/pipeline/transforms/uniquerowsbyhashset.html")
+ documentationUrl = "/pipeline/transforms/uniquerowsbyhashset.html",
+ // Must match SparkConst.PLUGIN_ID — do not import engines-spark from
transform modules.
+ excludedEngines = {"SparkPipelineEngine"})
@Getter
@Setter
public class UniqueRowsByHashSetMeta