This is an automated email from the ASF dual-hosted git repository.
github-merge-queue[bot] pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/texera.git
The following commit(s) were added to refs/heads/main by this push:
new 7c135ee095 feat(workflow-operator): support COUNT(*) in the Aggregate
operator (#5896)
7c135ee095 is described below
commit 7c135ee0955d539de0f7905b88d8bc4e86cca5a4
Author: Tanishq Gandhi <[email protected]>
AuthorDate: Tue Jul 7 13:58:37 2026 -0700
feat(workflow-operator): support COUNT(*) in the Aggregate operator (#5896)
### What changes were proposed in this PR?
This PR adds `COUNT(*)` support to the **Aggregate** operator, so users
can count all rows without selecting a column — useful when there's no
natural column to count.
The existing **`count`** function now handles both cases via an optional
attribute:
- **`count` with an empty Attribute** → `COUNT(*)` — counts every row
(including rows with nulls)
- **`count` with a column** → counts that column's non-null values
(unchanged)
The Attribute stays **required for every other function** (`sum`,
`average`, `min`,
`max`, `concat`), enforced by a conditional JSON-schema rule and
surfaced in the UI
with the usual required marker.
**Backend**
- `count` counts all rows when the attribute is empty/blank, otherwise
counts non-null values of the selected column.
- Conditional-required rule: attribute is required unless the function
is `count` (validated by Ajv, so an empty attribute on other functions
marks the operator invalid).
**Frontend**
- The Attribute is a normal optional field for `count`; the required
marker (red `*`) shows for all other functions. Description hints that
an empty attribute counts all rows.
**Docs**
- Updated the Aggregate operator reference.
#### Screenshots
<img width="1371" height="603" alt="image"
src="https://github.com/user-attachments/assets/b02cca86-aa69-4aff-a197-52ce1c683298"
/>
<img width="1318" height="618"
alt="2C3CB4D9-7FA6-4098-B7F7-293651EEA7D9"
src="https://github.com/user-attachments/assets/dd6c9330-c06f-48ed-bf11-70044d432aa6"
/>
### Any related issues, documentation, discussions?
Closes #3142.
### How was this PR tested?
**Automated** (`AggregateOpSpec`, `AggregateOpDescSpec`):
- `count` with an empty (or null) attribute counts all rows including
nulls.
- `count` with a column counts only non-null values.
- `getAggregationAttribute` / schema propagation / executor tolerate an
empty attribute (no input-column lookup; result typed `INTEGER`).
- Existing `sum`/`average`/`min`/`max`/`concat`/`getFinal` behavior
unchanged.
All Aggregate tests pass: `sbt "WorkflowOperator/testOnly *aggregate*"`.
**Manual (UI):** verified `count` + empty Attribute counts all rows
(with and without Group By), `count` + a column counts non-null values,
and non-count functions with an empty Attribute are invalid (Run
disabled).
### Was this PR authored or co-authored using generative AI tooling?
Generated-by: Claude Code (Claude Opus 4.8)
---------
Co-authored-by: Meng Wang <[email protected]>
---
.../amber/operator/aggregate/AggregateOpDesc.scala | 15 ++++--
.../amber/operator/aggregate/AggregateOpExec.scala | 15 ++++--
.../operator/aggregate/AggregationOperation.scala | 35 +++++++++----
.../operator/aggregate/AggregateOpDescSpec.scala | 46 ++++++++++++++++-
.../amber/operator/aggregate/AggregateOpSpec.scala | 59 +++++++++++++++++++++-
.../operators/data-cleaning/aggregate/aggregate.md | 6 ++-
.../operator-property-edit-frame.component.spec.ts | 18 ++++++-
.../operator-property-edit-frame.component.ts | 22 ++++++++
8 files changed, 195 insertions(+), 21 deletions(-)
diff --git
a/common/workflow-operator/src/main/scala/org/apache/texera/amber/operator/aggregate/AggregateOpDesc.scala
b/common/workflow-operator/src/main/scala/org/apache/texera/amber/operator/aggregate/AggregateOpDesc.scala
index 7e76b3ce7e..1bace652c8 100644
---
a/common/workflow-operator/src/main/scala/org/apache/texera/amber/operator/aggregate/AggregateOpDesc.scala
+++
b/common/workflow-operator/src/main/scala/org/apache/texera/amber/operator/aggregate/AggregateOpDesc.scala
@@ -78,9 +78,18 @@ class AggregateOpDesc extends LogicalOp {
val inputSchema = inputSchemas(operatorInfo.inputPorts.head.id)
val outputSchema = Schema(
groupByKeys.map(key => inputSchema.getAttribute(key)) ++
- localAggregations.map(agg =>
-
agg.getAggregationAttribute(inputSchema.getAttribute(agg.attribute).getType)
- )
+ localAggregations.map { agg =>
+ // Only COUNT with an empty attribute (COUNT(*)) skips the
column lookup:
+ // its result type is INTEGER regardless. Every other function
resolves
+ // the input attribute (failing fast if it is missing/invalid).
+ val attrType =
+ if (
+ agg.aggFunction == AggregationFunction.COUNT &&
+ (agg.attribute == null || agg.attribute.trim.isEmpty)
+ ) null
+ else inputSchema.getAttribute(agg.attribute).getType
+ agg.getAggregationAttribute(attrType)
+ }
)
Map(PortIdentity(internal = true) -> outputSchema)
})
diff --git
a/common/workflow-operator/src/main/scala/org/apache/texera/amber/operator/aggregate/AggregateOpExec.scala
b/common/workflow-operator/src/main/scala/org/apache/texera/amber/operator/aggregate/AggregateOpExec.scala
index a61703b28c..3c26709116 100644
---
a/common/workflow-operator/src/main/scala/org/apache/texera/amber/operator/aggregate/AggregateOpExec.scala
+++
b/common/workflow-operator/src/main/scala/org/apache/texera/amber/operator/aggregate/AggregateOpExec.scala
@@ -47,9 +47,18 @@ class AggregateOpExec(descString: String) extends
OperatorExecutor {
// Initialize distributedAggregations if it's not yet initialized
if (distributedAggregations == null) {
- distributedAggregations = desc.aggregations.map(agg =>
- agg.getAggFunc(tuple.getSchema.getAttribute(agg.attribute).getType)
- )
+ distributedAggregations = desc.aggregations.map { agg =>
+ // Only COUNT with an empty attribute (COUNT(*)) skips the column
lookup; its
+ // result does not depend on any input attribute. Every other function
resolves
+ // the input attribute (failing fast if it is missing/invalid).
+ val attrType =
+ if (
+ agg.aggFunction == AggregationFunction.COUNT &&
+ (agg.attribute == null || agg.attribute.trim.isEmpty)
+ ) null
+ else tuple.getSchema.getAttribute(agg.attribute).getType
+ agg.getAggFunc(attrType)
+ }
}
// Construct the group key
diff --git
a/common/workflow-operator/src/main/scala/org/apache/texera/amber/operator/aggregate/AggregationOperation.scala
b/common/workflow-operator/src/main/scala/org/apache/texera/amber/operator/aggregate/AggregationOperation.scala
index 70105de9ef..78ee70579c 100644
---
a/common/workflow-operator/src/main/scala/org/apache/texera/amber/operator/aggregate/AggregationOperation.scala
+++
b/common/workflow-operator/src/main/scala/org/apache/texera/amber/operator/aggregate/AggregationOperation.scala
@@ -55,7 +55,23 @@ case class AveragePartialObj(sum: Double, count: Double)
extends Serializable {}
}
]
}
- }
+ },
+ "allOf": [
+ {
+ "if": {
+ "properties": {
+ "aggFunction": { "const": "count" }
+ }
+ },
+ "then": {},
+ "else": {
+ "required": ["attribute"],
+ "properties": {
+ "attribute": { "pattern": "\\S" }
+ }
+ }
+ }
+ ]
}
""")
class AggregationOperation {
@@ -64,8 +80,8 @@ class AggregationOperation {
@JsonPropertyDescription("sum, count, average, min, max, or concat")
var aggFunction: AggregationFunction = _
- @JsonProperty(value = "attribute", required = true)
- @JsonPropertyDescription("column to calculate average value")
+ @JsonProperty(value = "attribute")
+ @JsonPropertyDescription("column to aggregate on")
@AutofillAttributeName
var attribute: String = _
@@ -106,6 +122,7 @@ class AggregationOperation {
@JsonIgnore
def getFinal: AggregationOperation = {
val newAggFunc = aggFunction match {
+ // COUNT emits partial counts locally; the global stage sums them.
case AggregationFunction.COUNT => AggregationFunction.SUM
case a: AggregationFunction => a
}
@@ -139,15 +156,13 @@ class AggregationOperation {
}
private def countAgg(): DistributedAggregation[Integer] = {
+ // An empty attribute means COUNT(*): count every row. Otherwise count
only the
+ // rows whose attribute value is non-null (COUNT(column)).
+ val countAllRows = attribute == null || attribute.trim.isEmpty
new DistributedAggregation[Integer](
() => 0,
- (partial, tuple) => {
- val inc =
- if (attribute == null) 1
- else if (tuple.getField(attribute) != null) 1
- else 0
- partial + inc
- },
+ (partial, tuple) =>
+ partial + (if (countAllRows || tuple.getField(attribute) != null) 1
else 0),
(partial1, partial2) => partial1 + partial2,
partial => partial
)
diff --git
a/common/workflow-operator/src/test/scala/org/apache/texera/amber/operator/aggregate/AggregateOpDescSpec.scala
b/common/workflow-operator/src/test/scala/org/apache/texera/amber/operator/aggregate/AggregateOpDescSpec.scala
index 681a1aa2be..0f0e372ad6 100644
---
a/common/workflow-operator/src/test/scala/org/apache/texera/amber/operator/aggregate/AggregateOpDescSpec.scala
+++
b/common/workflow-operator/src/test/scala/org/apache/texera/amber/operator/aggregate/AggregateOpDescSpec.scala
@@ -22,10 +22,12 @@ package org.apache.texera.amber.operator.aggregate
import org.apache.texera.amber.core.tuple.{AttributeType, Schema}
import org.apache.texera.amber.core.virtualidentity.{ExecutionIdentity,
WorkflowIdentity}
import org.apache.texera.amber.core.workflow.PortIdentity
-import org.apache.texera.amber.operator.metadata.OperatorGroupConstants
+import org.apache.texera.amber.operator.metadata.{OperatorGroupConstants,
OperatorMetadataGenerator}
import org.scalatest.flatspec.AnyFlatSpec
import org.scalatest.matchers.should.Matchers
+import scala.jdk.CollectionConverters._
+
class AggregateOpDescSpec extends AnyFlatSpec with Matchers {
private val workflowId = WorkflowIdentity(1L)
@@ -87,4 +89,46 @@ class AggregateOpDescSpec extends AnyFlatSpec with Matchers {
.getExternalOutputSchemas(Map(PortIdentity() -> input)) shouldBe
Map(PortIdentity() -> Schema().add("avg", AttributeType.DOUBLE))
}
+
+ it should "type a COUNT(*) (empty attribute) result as INTEGER without
looking up an input column" in {
+ // An empty attribute means COUNT(*); schema propagation must not
dereference a column.
+ val input = Schema().add("v", AttributeType.LONG)
+ descWith(List.empty, aggOp(AggregationFunction.COUNT, "", "row_count"))
+ .getExternalOutputSchemas(Map(PortIdentity() -> input)) shouldBe
+ Map(PortIdentity() -> Schema().add("row_count", AttributeType.INTEGER))
+ }
+
+ it should "fail fast for a non-COUNT function with an empty attribute (only
COUNT allows it)" in {
+ // Only COUNT tolerates a blank attribute; SUM/etc. must resolve the
column and fail
+ // fast rather than propagate a null-typed output.
+ val input = Schema().add("v", AttributeType.LONG)
+ assertThrows[Exception] {
+ descWith(List.empty, aggOp(AggregationFunction.SUM, "", "total"))
+ .getExternalOutputSchemas(Map(PortIdentity() -> input))
+ }
+ }
+
+ "AggregateOpDesc JSON schema" should
+ "make the attribute optional only for count and required for every other
function" in {
+ val aggDef = OperatorMetadataGenerator
+ .generateOperatorJsonSchema(classOf[AggregateOpDesc])
+ .get("definitions")
+ .get("AggregationOperation")
+
+ // attribute is not unconditionally required (aggFunction still is)
+ val baseRequired =
aggDef.get("required").elements().asScala.map(_.asText()).toSet
+ baseRequired should contain("aggFunction")
+ baseRequired should not contain "attribute"
+
+ // conditional rule: count -> no attribute requirement; any other function
-> attribute required
+ val rule = aggDef
+ .get("allOf")
+ .elements()
+ .asScala
+ .find(node => node.has("if") && node.has("else"))
+ .getOrElse(fail("expected a conditional if/else rule in the
AggregationOperation schema"))
+ rule.get("if").get("properties").get("aggFunction").get("const").asText()
shouldBe "count"
+ val elseRequired =
rule.get("else").get("required").elements().asScala.map(_.asText()).toList
+ elseRequired should contain("attribute")
+ }
}
diff --git
a/common/workflow-operator/src/test/scala/org/apache/texera/amber/operator/aggregate/AggregateOpSpec.scala
b/common/workflow-operator/src/test/scala/org/apache/texera/amber/operator/aggregate/AggregateOpSpec.scala
index cb7925ec41..3730181d9d 100644
---
a/common/workflow-operator/src/test/scala/org/apache/texera/amber/operator/aggregate/AggregateOpSpec.scala
+++
b/common/workflow-operator/src/test/scala/org/apache/texera/amber/operator/aggregate/AggregateOpSpec.scala
@@ -61,6 +61,16 @@ class AggregateOpSpec extends AnyFunSuite {
assert(attr.getType == AttributeType.INTEGER)
}
+ test("getAggregationAttribute maps COUNT result to INTEGER even with a null
input type") {
+ // COUNT(*) (empty attribute) has no input column, so schema propagation
passes a
+ // null attrType; it must still resolve to INTEGER without dereferencing
it.
+ val operation = makeAggregationOp(AggregationFunction.COUNT, "",
"row_count")
+ val attr = operation.getAggregationAttribute(null)
+
+ assert(attr.getName == "row_count")
+ assert(attr.getType == AttributeType.INTEGER)
+ }
+
test("getAggregationAttribute maps CONCAT result type to STRING") {
val operation = makeAggregationOp(AggregationFunction.CONCAT, "tag",
"all_tags")
val attr = operation.getAggregationAttribute(AttributeType.INTEGER)
@@ -142,7 +152,27 @@ class AggregateOpSpec extends AnyFunSuite {
assert(math.abs(result - 4.0) < 1e-6)
}
- test("COUNT aggregation with attribute == null counts all rows") {
+ test("COUNT with an empty attribute (COUNT(*)) counts all rows regardless of
nulls") {
+ // An empty attribute means COUNT(*); the GUI sends "" when no column is
selected.
+ val schema = makeSchema("points" -> AttributeType.INTEGER)
+ val tuple1 = makeTuple(schema, 10)
+ val tuple2 = makeTuple(schema, null)
+ val tuple3 = makeTuple(schema, 20)
+
+ val operation = makeAggregationOp(AggregationFunction.COUNT, "",
"row_count")
+ val agg = operation.getAggFunc(AttributeType.INTEGER)
+
+ var partial = agg.init()
+ partial = agg.iterate(partial, tuple1)
+ partial = agg.iterate(partial, tuple2)
+ partial = agg.iterate(partial, tuple3)
+
+ val result = agg.finalAgg(partial).asInstanceOf[Number].intValue()
+ assert(result == 3)
+ }
+
+ test("COUNT with a null attribute also counts all rows") {
+ // A null attribute is treated the same as empty (COUNT(*)).
val schema = makeSchema("points" -> AttributeType.INTEGER)
val tuple1 = makeTuple(schema, 10)
val tuple2 = makeTuple(schema, null)
@@ -535,4 +565,31 @@ class AggregateOpSpec extends AnyFunSuite {
assert(totalRevenue == 350)
assert(rowCount == 3)
}
+
+ test("AggregateOpExec computes COUNT(*) over every row (including nulls)
end-to-end") {
+ // region (ignored), revenue (one null). An empty attribute means
COUNT(*), so the
+ // executor must skip the input-column lookup and still count all 3 rows.
+ val schema = makeSchema(
+ "region" -> AttributeType.STRING,
+ "revenue" -> AttributeType.INTEGER
+ )
+
+ val tuple1 = makeTuple(schema, "west", 100)
+ val tuple2 = makeTuple(schema, "east", null)
+ val tuple3 = makeTuple(schema, "west", 50)
+
+ val desc = new AggregateOpDesc()
+ desc.aggregations = List(makeAggregationOp(AggregationFunction.COUNT, "",
"row_count"))
+ desc.groupByKeys = List() // global aggregation
+
+ val exec = new AggregateOpExec(objectMapper.writeValueAsString(desc))
+ exec.open()
+ exec.processTuple(tuple1, 0)
+ exec.processTuple(tuple2, 0)
+ exec.processTuple(tuple3, 0)
+
+ val results = exec.onFinish(0).toList
+ assert(results.size == 1)
+ assert(results.head.getFields(0).asInstanceOf[Number].intValue() == 3)
+ }
}
diff --git a/docs/reference/operators/data-cleaning/aggregate/aggregate.md
b/docs/reference/operators/data-cleaning/aggregate/aggregate.md
index e2b3f0e997..1fcd883d7c 100644
--- a/docs/reference/operators/data-cleaning/aggregate/aggregate.md
+++ b/docs/reference/operators/data-cleaning/aggregate/aggregate.md
@@ -14,10 +14,12 @@ tags: [data-cleaning, aggregate]
|----------|-------------|------|---------|-------------|
| Aggregations | ✓ | List<Aggregation> | - | Multiple aggregation functions
(min: 1,<br>aggregations cannot be empty) |
| ↳ Aggregate Func | ✓ | sum, count, average, min, max, concat | - | Sum,
count, average, min, max, or concat |
-| ↳ Attribute | ✓ | String | - | Column to calculate average value |
-| ↳ Result Attribute | ✓ | String | - | Column name of average result |
+| ↳ Attribute | ✓ (optional for `count`) | String | - | Column to aggregate
on. Required for every function except `count`: leave it empty with `count` to
count all rows (`COUNT(*)`), or pick a column to count its non-null values |
+| ↳ Result Attribute | ✓ | String | - | Column name of the aggregation result |
| Group By Keys | | List | - | Group by columns |
+> **Counting rows**: with the `count` function, leave **Attribute** empty to
count every row (`COUNT(*)`, including rows with nulls), or choose a column to
count only that column's non-null values.
+
### Output Ports
| Port | Mode |
diff --git
a/frontend/src/app/workspace/component/property-editor/operator-property-edit-frame/operator-property-edit-frame.component.spec.ts
b/frontend/src/app/workspace/component/property-editor/operator-property-edit-frame/operator-property-edit-frame.component.spec.ts
index eb0088ba18..abe4aad391 100644
---
a/frontend/src/app/workspace/component/property-editor/operator-property-edit-frame/operator-property-edit-frame.component.spec.ts
+++
b/frontend/src/app/workspace/component/property-editor/operator-property-edit-frame/operator-property-edit-frame.component.spec.ts
@@ -19,7 +19,11 @@
import { ComponentFixture, discardPeriodicTasks, fakeAsync, TestBed, tick }
from "@angular/core/testing";
-import { OperatorPropertyEditFrameComponent } from
"./operator-property-edit-frame.component";
+import {
+ AGGREGATE_COUNT,
+ isAggregateAttributeRequired,
+ OperatorPropertyEditFrameComponent,
+} from "./operator-property-edit-frame.component";
import { WorkflowActionService } from
"../../../service/workflow-graph/model/workflow-action.service";
import { OperatorMetadataService } from
"../../../service/operator-metadata/operator-metadata.service";
import { StubOperatorMetadataService } from
"../../../service/operator-metadata/stub-operator-metadata.service";
@@ -52,6 +56,18 @@ import { MockComputingUnitStatusService } from
"../../../../common/service/compu
import { commonTestProviders } from "../../../../common/testing/test-utils";
const { marbles } = configure({ run: false });
+
+describe("Aggregate attribute requirement", () => {
+ it("makes the attribute optional for count and required for every other
function", () => {
+ // count -> optional (empty attribute means COUNT(*))
+ expect(isAggregateAttributeRequired(AGGREGATE_COUNT)).toBe(false);
+ // every other aggregate function -> attribute required
+ ["sum", "average", "min", "max", "concat"].forEach(fn => {
+ expect(isAggregateAttributeRequired(fn)).toBe(true);
+ });
+ });
+});
+
describe("OperatorPropertyEditFrameComponent", () => {
let component: OperatorPropertyEditFrameComponent;
let fixture: ComponentFixture<OperatorPropertyEditFrameComponent>;
diff --git
a/frontend/src/app/workspace/component/property-editor/operator-property-edit-frame/operator-property-edit-frame.component.ts
b/frontend/src/app/workspace/component/property-editor/operator-property-edit-frame/operator-property-edit-frame.component.ts
index 2512fecdac..b5f01423db 100644
---
a/frontend/src/app/workspace/component/property-editor/operator-property-edit-frame/operator-property-edit-frame.component.ts
+++
b/frontend/src/app/workspace/component/property-editor/operator-property-edit-frame/operator-property-edit-frame.component.ts
@@ -77,6 +77,17 @@ import { map, switchMap, take } from "rxjs/operators";
Quill.register("modules/cursors", QuillCursors);
+// The Aggregate "count" function. With an empty attribute it means COUNT(*)
(all rows);
+// with a column it counts that column's non-null values. It is the only
function whose
+// attribute is optional.
+export const AGGREGATE_COUNT = "count";
+
+// The Aggregate attribute is required for every function except `count` (an
empty
+// attribute on count means COUNT(*), which needs no column).
+export function isAggregateAttributeRequired(aggFunction: unknown): boolean {
+ return aggFunction !== AGGREGATE_COUNT;
+}
+
/**
* Property Editor uses JSON Schema to automatically generate the form from
the JSON Schema of an operator.
* For example, the JSON Schema of Sentiment Analysis could be:
@@ -545,6 +556,17 @@ export class OperatorPropertyEditFrameComponent implements
OnInit, OnChanges, On
mappedField.type = "datasetversionselector";
}
+ // Aggregate: the attribute is optional for `count` (an empty attribute
means COUNT(*),
+ // counting all rows) and required for every other function. Show the
required marker
+ // (red *) accordingly, based on the sibling aggFunction within the same
row.
+ if (this.currentOperatorSchema?.operatorType === "Aggregate" &&
mappedField.key === "attribute") {
+ mappedField.expressions = {
+ ...mappedField.expressions,
+ "props.required": (field: FormlyFieldConfig) =>
+ isAggregateAttributeRequired(field.parent?.model?.aggFunction),
+ };
+ }
+
if (this.currentOperatorSchema?.operatorType === "FileScanOp" &&
mappedField.key === "outputFileName") {
mappedField.expressions = {
...mappedField.expressions,