This is an automated email from the ASF dual-hosted git repository. damccorm pushed a commit to branch backport-ca4f59e5 in repository https://gitbox.apache.org/repos/asf/beam.git
commit 01464e96da7d61fd6405b44287dc8e5b26f6e235 Author: Danny McCormick <[email protected]> AuthorDate: Mon Jul 27 20:37:50 2026 +0000 Add SINGLE_VALUE aggregate function --- .../impl/transform/BeamBuiltinAggregations.java | 5 +++ .../extensions/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 2800edfbb99..d2da23a8517 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.<String, Function<Schema.FieldType, CombineFn<?, ?, ?>>>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 37243fe7d5f..9374713a92f 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 @@ public class BeamSqlDslAggregationTest extends BeamSqlDslBase { 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<Row> 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<Row> result = inputRows.apply("sql", SqlTransform.query(sql)); + + Map<Integer, List<Integer>> 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<Row> iter = input.iterator(); + while (iter.hasNext()) { + Row row = iter.next(); + List<Integer> 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 {
