From 01464e96da7d61fd6405b44287dc8e5b26f6e235 Mon Sep 17 00:00:00 2001 From: Danny McCormick Date: Mon, 27 Jul 2026 20:37:50 +0000 Subject: [PATCH 1/2] Add SINGLE_VALUE aggregate function --- .../transform/BeamBuiltinAggregations.java | 5 ++ .../sql/BeamSqlDslAggregationTest.java | 47 +++++++++++++++++++ 2 files changed, 52 insertions(+) diff --git a/sdks/java/extensions/sql/src/main/java/org/apache/beam/sdk/extensions/sql/impl/transform/BeamBuiltinAggregations.java b/sdks/java/extensions/sql/src/main/java/org/apache/beam/sdk/extensions/sql/impl/transform/BeamBuiltinAggregations.java index 2800edfbb99a..d2da23a85171 100644 --- a/sdks/java/extensions/sql/src/main/java/org/apache/beam/sdk/extensions/sql/impl/transform/BeamBuiltinAggregations.java +++ b/sdks/java/extensions/sql/src/main/java/org/apache/beam/sdk/extensions/sql/impl/transform/BeamBuiltinAggregations.java @@ -61,6 +61,11 @@ public class BeamBuiltinAggregations { BUILTIN_AGGREGATOR_FACTORIES = ImmutableMap.>>builder() .put("ANY_VALUE", typeName -> Sample.anyValueCombineFn()) + // SINGLE_VALUE is emitted by Calcite to enforce the cardinality of a scalar + // subquery (a subquery used as a scalar must yield exactly one row). The single + // input value is returned as-is; unlike COUNT/SUM it must not drop nulls, so a + // scalar subquery evaluating to NULL surfaces NULL. + .put("SINGLE_VALUE", typeName -> Sample.anyValueCombineFn()) // Drop null elements for these aggregations. .put("COUNT", typeName -> new DropNullFnWithDefault(Count.combineFn())) .put("MAX", typeName -> new DropNullFn(BeamBuiltinAggregations.createMax(typeName))) diff --git a/sdks/java/extensions/sql/src/test/java/org/apache/beam/sdk/extensions/sql/BeamSqlDslAggregationTest.java b/sdks/java/extensions/sql/src/test/java/org/apache/beam/sdk/extensions/sql/BeamSqlDslAggregationTest.java index 37243fe7d5f6..9374713a92fc 100644 --- a/sdks/java/extensions/sql/src/test/java/org/apache/beam/sdk/extensions/sql/BeamSqlDslAggregationTest.java +++ b/sdks/java/extensions/sql/src/test/java/org/apache/beam/sdk/extensions/sql/BeamSqlDslAggregationTest.java @@ -373,6 +373,53 @@ private void runAggregationFunctions(PCollection input) throws Exception { pipeline.run().waitUntilFinish(); } + /** GROUP-BY with the SINGLE_VALUE aggregation function. */ + @Test + public void testSingleValueFunction() throws Exception { + pipeline.enableAbandonedNodeEnforcement(false); + + Schema schema = Schema.builder().addInt32Field("key").addInt32Field("col").build(); + + PCollection inputRows = + pipeline + .apply( + Create.of( + TestUtils.rowsBuilderOf(schema) + .addRows( + 0, 1, + 0, 2, + 1, 3, + 2, 4, + 2, 5) + .getRows())) + .setRowSchema(schema); + + String sql = "SELECT key, SINGLE_VALUE(col) as single_value FROM PCOLLECTION GROUP BY key"; + + PCollection result = inputRows.apply("sql", SqlTransform.query(sql)); + + Map> allowedTuples = new HashMap<>(); + allowedTuples.put(0, Arrays.asList(1, 2)); + allowedTuples.put(1, Arrays.asList(3)); + allowedTuples.put(2, Arrays.asList(4, 5)); + + PAssert.that(result) + .satisfies( + input -> { + Iterator iter = input.iterator(); + while (iter.hasNext()) { + Row row = iter.next(); + List values = allowedTuples.remove(row.getInt32("key")); + assertTrue(values != null); + assertTrue(values.contains(row.getInt32("single_value"))); + } + assertTrue(allowedTuples.isEmpty()); + return null; + }); + + pipeline.run(); + } + /** GROUP-BY with the any_value aggregation function. */ @Test public void testAnyValueFunction() throws Exception { From 8bf095f8fa6a78b90ece5ca7f2b41185a5818272 Mon Sep 17 00:00:00 2001 From: Danny McCormick Date: Mon, 27 Jul 2026 22:52:11 +0000 Subject: [PATCH 2/2] Implement SingleValueCombineFn to enforce cardinality --- .../transform/BeamBuiltinAggregations.java | 53 ++++++++++++++ .../sql/BeamSqlDslAggregationTest.java | 72 +++++++++++++++++++ 2 files changed, 125 insertions(+) diff --git a/sdks/java/extensions/sql/src/main/java/org/apache/beam/sdk/extensions/sql/impl/transform/BeamBuiltinAggregations.java b/sdks/java/extensions/sql/src/main/java/org/apache/beam/sdk/extensions/sql/impl/transform/BeamBuiltinAggregations.java index d2da23a85171..1759ae12034e 100644 --- a/sdks/java/extensions/sql/src/main/java/org/apache/beam/sdk/extensions/sql/impl/transform/BeamBuiltinAggregations.java +++ b/sdks/java/extensions/sql/src/main/java/org/apache/beam/sdk/extensions/sql/impl/transform/BeamBuiltinAggregations.java @@ -60,6 +60,11 @@ public class BeamBuiltinAggregations { public static final Map>> BUILTIN_AGGREGATOR_FACTORIES = ImmutableMap.>>builder() + // SINGLE_VALUE is emitted by Calcite to enforce the cardinality of a scalar + // subquery (a subquery used as a scalar must yield exactly one row). The single + // input value is returned as-is; unlike COUNT/SUM it must not drop nulls, so a + // scalar subquery evaluating to NULL surfaces NULL. + .put("SINGLE_VALUE", typeName -> new SingleValue<>()) .put("ANY_VALUE", typeName -> Sample.anyValueCombineFn()) // SINGLE_VALUE is emitted by Calcite to enforce the cardinality of a scalar // subquery (a subquery used as a scalar must yield exactly one row). The single @@ -698,4 +703,52 @@ public Long extractOutput(BitXOr.Accum accumulator) { return accumulator.bitXOr; } } + + /** + * {@link CombineFn} for SINGLE_VALUE to enforce that a scalar subquery returns exactly one row. + */ + public static class SingleValue extends CombineFn, T> { + @Override + public java.util.List createAccumulator() { + return new java.util.ArrayList<>(); + } + + @Override + public java.util.List addInput(java.util.List accumulator, T input) { + if (accumulator.size() < 2) { + accumulator.add(input); + } + return accumulator; + } + + @Override + public java.util.List mergeAccumulators(Iterable> accumulators) { + java.util.List merged = createAccumulator(); + for (java.util.List accum : accumulators) { + merged.addAll(accum); + if (merged.size() >= 2) { + merged = new java.util.ArrayList<>(merged.subList(0, 2)); + } + } + return merged; + } + + @Override + public T extractOutput(java.util.List accumulator) { + if (accumulator.isEmpty()) { + return null; + } + if (accumulator.size() > 1) { + throw new IllegalArgumentException("Subquery returned more than one row"); + } + return accumulator.get(0); + } + + @Override + public Coder> getAccumulatorCoder( + CoderRegistry registry, Coder inputCoder) { + return org.apache.beam.sdk.coders.ListCoder.of( + org.apache.beam.sdk.coders.NullableCoder.of(inputCoder)); + } + } } diff --git a/sdks/java/extensions/sql/src/test/java/org/apache/beam/sdk/extensions/sql/BeamSqlDslAggregationTest.java b/sdks/java/extensions/sql/src/test/java/org/apache/beam/sdk/extensions/sql/BeamSqlDslAggregationTest.java index 9374713a92fc..0d354c383f45 100644 --- a/sdks/java/extensions/sql/src/test/java/org/apache/beam/sdk/extensions/sql/BeamSqlDslAggregationTest.java +++ b/sdks/java/extensions/sql/src/test/java/org/apache/beam/sdk/extensions/sql/BeamSqlDslAggregationTest.java @@ -467,6 +467,78 @@ public void testAnyValueFunction() throws Exception { pipeline.run(); } + /** GROUP-BY with the SINGLE_VALUE aggregation function. */ + @Test + public void testSingleValueFunction() throws Exception { + pipeline.enableAbandonedNodeEnforcement(false); + + Schema schema = Schema.builder().addInt32Field("key").addInt32Field("col").build(); + + PCollection inputRows = + pipeline + .apply( + Create.of( + TestUtils.rowsBuilderOf(schema) + .addRows( + 0, 1, + 1, 3, + 2, 4) + .getRows())) + .setRowSchema(schema); + + String sql = "SELECT key, SINGLE_VALUE(col) as single_value FROM PCOLLECTION GROUP BY key"; + + PCollection result = inputRows.apply("sql", SqlTransform.query(sql)); + + Map> allowedTuples = new HashMap<>(); + allowedTuples.put(0, Arrays.asList(1)); + allowedTuples.put(1, Arrays.asList(3)); + allowedTuples.put(2, Arrays.asList(4)); + + PAssert.that(result) + .satisfies( + input -> { + Iterator iter = input.iterator(); + while (iter.hasNext()) { + Row row = iter.next(); + List values = allowedTuples.remove(row.getInt32("key")); + assertTrue(values != null); + assertTrue(values.contains(row.getInt32("single_value"))); + } + assertTrue(allowedTuples.isEmpty()); + return null; + }); + + pipeline.run(); + } + + /** GROUP-BY with the SINGLE_VALUE aggregation function fails on multiple values. */ + @Test + public void testSingleValueFunction_throwsOnMultiple() throws Exception { + pipeline.enableAbandonedNodeEnforcement(false); + + Schema schema = Schema.builder().addInt32Field("key").addInt32Field("col").build(); + + PCollection inputRows = + pipeline + .apply( + Create.of( + TestUtils.rowsBuilderOf(schema) + .addRows( + 0, 1, + 0, 2) + .getRows())) + .setRowSchema(schema); + + String sql = "SELECT key, SINGLE_VALUE(col) as single_value FROM PCOLLECTION GROUP BY key"; + + inputRows.apply("sql", SqlTransform.query(sql)); + + thrown.expect(Exception.class); + thrown.expectMessage("Subquery returned more than one row"); + pipeline.run(); + } + @Test public void testBitOrFunction() throws Exception { pipeline.enableAbandonedNodeEnforcement(false);