Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 4 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -134,6 +134,8 @@ All notable changes to this project will be documented in this file. It uses the
serializer previously for scalar values (now simplified), preventing
PostgreSQL `IN` lists from producing invalid or mismatched `FixedString`
comparisons. Thanks to @jxom for the PR (#333)!
* Fixed `array_agg()` pushdown dropping `NULL` elements because ClickHouse's
`groupArray()` skips them ([#340]).
* Fixed `array_position()` pushdown to return `NULL`, rather than zero, when
ClickHouse does not find an element. Thanks to Minh Vu for the PR ([#334])!
* Fixed `array_length()` pushdown to preserve empty-array and requested
Expand Down Expand Up @@ -211,6 +213,8 @@ All notable changes to this project will be documented in this file. It uses the
"ClickHouse/pg_clickhouse#335 Reject JSON paths containing NULL elements"
[#336]: https://github.com/ClickHouse/pg_clickhouse/pull/336
"ClickHouse/pg_clickhouse#336 Preserve array_length semantics"
[#340]: https://github.com/ClickHouse/pg_clickhouse/pull/340
"ClickHouse/pg_clickhouse#340 Preserve NULLs in array_agg pushdown"

## [v0.3.2] — 2026-06-16

Expand Down
3 changes: 2 additions & 1 deletion src/custom_types.c
Original file line number Diff line number Diff line change
Expand Up @@ -536,7 +536,8 @@ lookup_builtin_func(Oid funcid, builtin_func_def* def) {
def->ch_name = "quantilesExactLow";
return true;
case F_ARRAY_AGG_ANYNONARRAY:
def->ch_name = "groupArray";
def->cf_type = CF_ARRAY_AGG;
def->ch_name = "\1";
return true;
case F_MD5_BYTEA:
case F_MD5_TEXT:
Expand Down
68 changes: 67 additions & 1 deletion src/deparse.c
Original file line number Diff line number Diff line change
Expand Up @@ -324,6 +324,8 @@ deparseRangeTblRef(
static void
deparseAggref(Aggref* node, deparse_expr_cxt* context);
static void
deparseArrayAggref(Aggref* node, deparse_expr_cxt* context);
static void
deparseWindowFunc(WindowFunc* node, deparse_expr_cxt* context);
static void
appendGroupByClause(List* tlist, deparse_expr_cxt* context);
Expand Down Expand Up @@ -1235,7 +1237,9 @@ foreign_expr_walker(Node* node, foreign_glob_cxt* glob_cxt, ExprTruthCtx ctx) {
}

/* groupConcat has no ORDER BY; block ordered string_agg */
if (agg->aggfnoid == F_STRING_AGG_TEXT_TEXT && agg->aggorder != NIL) {
if ((agg->aggfnoid == F_STRING_AGG_TEXT_TEXT ||
agg->aggfnoid == F_ARRAY_AGG_ANYNONARRAY) &&
agg->aggorder != NIL) {
Comment on lines +1240 to +1242

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Need test rejecting use of ORDER with F_ARRAY_AGG_ANYNONARRAY

return false;
}

Expand Down Expand Up @@ -5855,6 +5859,62 @@ deparsePartialStatArray(Aggref* node, AggPartialKind kind, deparse_expr_cxt* con
pfree(arg);
}

static void
deparseArrayAggref(Aggref* node, deparse_expr_cxt* context) {
StringInfo buf = context->buf;
foreign_glob_cxt glob_cxt;
TargetEntry* arg = NULL;
bool nullable;
ListCell* lc;

foreach (lc, node->args) {
TargetEntry* tle = lfirst_node(TargetEntry, lc);

if (!tle->resjunk) {
arg = tle;
break;
}
}

Assert(arg != NULL);

memset(&glob_cxt, 0, sizeof(glob_cxt));
glob_cxt.root = context->root;
glob_cxt.foreignrel = context->foreignrel;
glob_cxt.relids = context->scanrel->relids;
nullable = !expr_never_null((Expr*)arg->expr, &glob_cxt);

if (nullable) {
appendStringInfoString(buf, "arrayMap(x -> x.1, ");
}
appendStringInfoString(buf, "groupArray");
if (node->aggfilter) {
appendStringInfoString(buf, "If");
}
appendStringInfoString(buf, "(");
if (node->aggdistinct != NIL) {
appendStringInfoString(buf, "DISTINCT ");
}
if (nullable) {
appendStringInfoString(buf, "tuple(");
}
deparseExpr((Expr*)arg->expr, context);
if (nullable) {
appendStringInfoChar(buf, ')');
}

if (node->aggfilter) {
appendStringInfoString(buf, ",((");
deparseExpr((Expr*)node->aggfilter, context);
appendStringInfoString(buf, ") > 0)");
}

appendStringInfoChar(buf, ')');
if (nullable) {
appendStringInfoChar(buf, ')');
}
}

/*
* Deparse an Aggref node.
*/
Expand Down Expand Up @@ -5891,6 +5951,12 @@ deparseAggref(Aggref* node, deparse_expr_cxt* context) {
cdef = context->func;
context->func = appendFunctionName(node->aggfnoid, context);

if (context->func && context->func->cf_type == CF_ARRAY_AGG) {
deparseArrayAggref(node, context);
context->func = cdef;
return;
}

/* 'If' part */
if (context->func && context->func->cf_type == CF_SIGN_COUNT && !node->aggstar) {
sign_count_filter = true;
Expand Down
2 changes: 2 additions & 0 deletions src/include/fdw.h
Original file line number Diff line number Diff line change
Expand Up @@ -406,6 +406,8 @@ typedef enum {
* length(arr)-n) */
CF_ARRAY_SORT_DESC, /* array_sort(arr,desc) →
* arrayReverseSort/arraySort */
CF_ARRAY_AGG, /* array_agg → groupArray; tuple-wrap nullable
* inputs because groupArray skips NULLs */
CF_ARRAY_FILL, /* array_fill → arrayWithConstant,
* swap+extract */
CF_ARRAY_CONTAINS, /* @> → hasAll(left, right) */
Expand Down
136 changes: 134 additions & 2 deletions test/expected/aggregates.out
Original file line number Diff line number Diff line change
Expand Up @@ -27,6 +27,16 @@
clickhouse_raw_query
----------------------

(1 row)

clickhouse_raw_query
----------------------

(1 row)

clickhouse_raw_query
----------------------

(1 row)

Foreign table "agg_bin.agg_numbers"
Expand All @@ -50,6 +60,13 @@ FDW options: (database 'agg_test', table_name 'agg_numbers', engine 'MergeTree')
Server: agg_bin_svr
FDW options: (database 'agg_test', table_name 'hits', engine 'MergeTree')

Foreign table "agg_bin.null_agg"
Column | Type | Collation | Nullable | Default | FDW options
--------+---------+-----------+----------+---------+-------------
v | integer | | | |
Server: agg_bin_svr
FDW options: (database 'agg_test', table_name 'null_agg', engine 'TinyLog')

Foreign table "agg_http.agg_numbers"
Column | Type | Collation | Nullable | Default | FDW options
--------+----------+-----------+----------+---------+-------------
Expand All @@ -71,6 +88,13 @@ FDW options: (database 'agg_test', table_name 'agg_numbers', engine 'MergeTree')
Server: agg_http_svr
FDW options: (database 'agg_test', table_name 'hits', engine 'MergeTree')

Foreign table "agg_http.null_agg"
Column | Type | Collation | Nullable | Default | FDW options
--------+---------+-----------+----------+---------+-------------
v | integer | | | |
Server: agg_http_svr
FDW options: (database 'agg_test', table_name 'null_agg', engine 'TinyLog')

-- AVG(UInt64)
QUERY PLAN
-------------------------------------------------
Expand Down Expand Up @@ -601,6 +625,112 @@ FDW options: (database 'agg_test', table_name 'hits', engine 'MergeTree')
{5.10,1.99,8.14,8.53,3.57,2.47,8.24,4.79,6.75,7.69}
(1 row)

QUERY PLAN
--------------------------------------------------------------------------------------
Foreign Scan
Output: (array_agg(v))
Relations: Aggregate on (null_agg)
Remote SQL: SELECT arrayMap(x -> x.1, groupArray(tuple(v))) FROM agg_test.null_agg
(4 rows)

array_agg
-------------------
{1,NULL,2,1,NULL}
(1 row)

QUERY PLAN
--------------------------------------------------------------------------------------
Foreign Scan
Output: (array_agg(v))
Relations: Aggregate on (null_agg)
Remote SQL: SELECT arrayMap(x -> x.1, groupArray(tuple(v))) FROM agg_test.null_agg
(4 rows)

array_agg
-------------------
{1,NULL,2,1,NULL}
(1 row)

QUERY PLAN
--------------------------------------------------------------------------------------------------------
Foreign Scan
Output: (array_agg(v) FILTER (WHERE (v > 1)))
Relations: Aggregate on (null_agg)
Remote SQL: SELECT arrayMap(x -> x.1, groupArrayIf(tuple(v),(((v > 1)) > 0))) FROM agg_test.null_agg
(4 rows)

array_agg
-----------
{2}
(1 row)

QUERY PLAN
--------------------------------------------------------------------------------------------------------
Foreign Scan
Output: (array_agg(v) FILTER (WHERE (v > 1)))
Relations: Aggregate on (null_agg)
Remote SQL: SELECT arrayMap(x -> x.1, groupArrayIf(tuple(v),(((v > 1)) > 0))) FROM agg_test.null_agg
(4 rows)

array_agg
-----------
{2}
(1 row)

QUERY PLAN
-----------------------------------------------------------------------------------------------
Foreign Scan
Output: (array_agg(DISTINCT v))
Relations: Aggregate on (null_agg)
Remote SQL: SELECT arrayMap(x -> x.1, groupArray(DISTINCT tuple(v))) FROM agg_test.null_agg
(4 rows)

array_agg
------------
{1,NULL,2}
(1 row)

QUERY PLAN
-----------------------------------------------------------------------------------------------
Foreign Scan
Output: (array_agg(DISTINCT v))
Relations: Aggregate on (null_agg)
Remote SQL: SELECT arrayMap(x -> x.1, groupArray(DISTINCT tuple(v))) FROM agg_test.null_agg
(4 rows)

array_agg
------------
{1,NULL,2}
(1 row)

QUERY PLAN
-------------------------------------------------------------------------
Aggregate
Output: array_agg(v ORDER BY v)
-> Foreign Scan on agg_bin.null_agg
Output: v
Remote SQL: SELECT v FROM agg_test.null_agg ORDER BY v ASC NULLS LAST
(5 rows)

array_agg
-------------------
{1,1,2,NULL,NULL}
(1 row)

QUERY PLAN
-------------------------------------------------------------------------
Aggregate
Output: array_agg(v ORDER BY v)
-> Foreign Scan on agg_http.null_agg
Output: v
Remote SQL: SELECT v FROM agg_test.null_agg ORDER BY v ASC NULLS LAST
(5 rows)

array_agg
-------------------
{1,1,2,NULL,NULL}
(1 row)

-- min(UInt64)
QUERY PLAN
-------------------------------------------------
Expand Down Expand Up @@ -1834,9 +1964,11 @@ FDW options: (database 'agg_test', table_name 'hits', engine 'MergeTree')
100 | 100
(4 rows)

NOTICE: drop cascades to 2 other objects
NOTICE: drop cascades to 3 other objects
DETAIL: drop cascades to foreign table agg_bin.agg_numbers
drop cascades to foreign table agg_bin.hits
NOTICE: drop cascades to 2 other objects
drop cascades to foreign table agg_bin.null_agg
NOTICE: drop cascades to 3 other objects
DETAIL: drop cascades to foreign table agg_http.agg_numbers
drop cascades to foreign table agg_http.hits
drop cascades to foreign table agg_http.null_agg
Loading