
Transform Aggregate
CertifiedAggregate records by group
Transform Aggregate
Aggregate records by group
Group records by one or more fields and compute typed summary values such as counts, sums, minimums, maximums, first, or last values.
type: io.kestra.plugin.transform.AggregateExamples
Aggregate totals
id: aggregate_totals_records
namespace: company.team
tasks:
- id: normalize
type: io.kestra.plugin.core.output.OutputValues
values:
records:
- customer_id: c1
country: FR
total_spent: 10
- customer_id: c1
country: FR
total_spent: 5
- id: aggregate
type: io.kestra.plugin.transform.Aggregate
from: "{{ outputs.normalize.values.records }}"
groupBy:
- customer_id
- country
aggregates:
order_count:
expr: count()
type: INT
total_spent:
expr: sum(total_spent)
type: DECIMAL
onError: FAIL
Aggregate with stored output
id: aggregate_totals
namespace: company.team
tasks:
- id: fetch
type: io.kestra.plugin.core.output.OutputValues
values:
records:
- customer_id: "c1"
country: "FR"
total_spent: 10
- customer_id: "c1"
country: "FR"
total_spent: 5
- id: aggregate
type: io.kestra.plugin.transform.Aggregate
from: "{{ outputs.fetch.values.records }}"
outputType: STORE
groupBy:
- customer_id
- country
aggregates:
order_count:
expr: count()
type: INT
total_spent:
expr: sum(total_spent)
type: DECIMAL
Properties
aggregates *object
Aggregate definitions
Output fields to compute for each group. Each value can be a shorthand expression string or an object with expr and optional type.
from *object
Input records
Ion list or struct to transform, or a storage URI pointing to an Ion file.
groupBy *array
Group by
Fields to group on.
assets
Assets this task consumes as inputs or produces as outputs, for lineage tracking and the asset graph (Enterprise Edition). A flow declaring this property on a task is rejected in the open-source edition.
io.kestra.core.models.assets.AssetsDeclaration
IGNOREFAILWARNAsset failure behavior
Behavior applied to the task state when a declared asset fails to render, emit, or be persisted (e.g. a lock conflict): FAIL escalates it to FAILED, WARN (default) warns it if it would otherwise succeed, IGNORE leaves the state untouched.
Whether to auto-register assets referenced dynamically at runtime that are not statically declared in inputs or outputs.
The assets consumed as inputs.
io.kestra.core.models.assets.AssetIdentifier
1The assets produced as outputs.
io.kestra.plugin.ee.assets.Dataset
1150{}1150io.kestra.plugin.ee.assets.File
1150{}1150io.kestra.plugin.ee.assets.Table
1150{}1150io.kestra.plugin.ee.assets.VM
1150{}1150io.kestra.core.models.assets.External
1150{}1150io.kestra.core.models.assets.Custom
11501Custom asset type
{}1150onError string
FAILFAILSKIPNULLOn error behavior
FAIL stops the task on aggregate errors, SKIP drops failing input records, and NULL sets failing aggregate outputs to null.
outputFormat string
TEXTTEXTBINARYOutput format
Experimental: TEXT or BINARY. Only transform tasks can read binary Ion. Use TEXT as the final step.
outputType string
AUTOAUTORECORDSSTOREOutput type
AUTO stores to internal storage when the input is a storage URI; otherwise it returns records.
Outputs
records array
Aggregated records
JSON-safe records when output mode is RECORDS or AUTO resolves to RECORDS.
uri string
Stored Ion file URI
URI to the stored Ion file when output mode is STORE or AUTO resolves to STORE.