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

Reply via email to