Aggregate
- async AsyncCogniteClient.data_modeling.records.aggregate(
- aggregates: Mapping[str, Aggregate | dict[str, Any]],
- *,
- stream_id: str,
- last_updated_time: TimeRange | None = None,
- filter: Filter | dict[str, Any] | None = None,
- target_units: RecordTargetUnits | Sequence[RecordTargetUnit] | None = None,
- include_typing: bool = False,
Aggregate records from a stream.
- Parameters:
aggregates (Mapping[str, Aggregate | dict[str, Any]]) – Aggregate request tree keyed by client-defined aggregate IDs.
stream_id (str) – External ID of the stream to aggregate from.
last_updated_time (TimeRange | None) – Filter records by last-updated time. Required for immutable streams (must include a lower bound).
filter (Filter | dict[str, Any] | None) – Filter expression.
target_units (RecordTargetUnits | Sequence[RecordTargetUnit] | None) – Unit conversion specification.
include_typing (bool) – Include property type metadata in the response.
- Returns:
Aggregate results keyed by the requested aggregate IDs.
- Return type:
Examples
The examples below aggregate over a stream of padel game statistics records; each example builds on the previous one.
The property paths used below:
>>> game_time = ["paddle", "game_statistics", "game_time"] >>> player_name = ["paddle", "game_statistics", "player_name"] >>> points_scored = ["paddle", "game_statistics", "points_scored"]
Find the average points scored across all games, using a typed helper:
>>> from cognite.client import CogniteClient >>> from cognite.client.data_classes.data_modeling.aggregates import Average >>> client = CogniteClient() >>> res = client.data_modeling.records.aggregate( ... stream_id="my-stream", ... aggregates={"avg_points_scored": Average(points_scored)}, ... )
Count the total number of games, and how many of them have a recorded score, only considering games updated after a given time:
>>> from cognite.client.data_classes.data_modeling.aggregates import Count >>> from cognite.client.data_classes.data_modeling.records import TimeRange >>> res = client.data_modeling.records.aggregate( ... stream_id="my-stream", ... aggregates={ ... "total_games": Count(), ... "games_with_score": Count(points_scored), ... }, ... last_updated_time=TimeRange(gt=1759276800000), ... )
Group games by day, then by player, and for each player-day compute their total, highest, and average points scored, alongside the single highest score across all games:
>>> from cognite.client.data_classes.data_modeling.aggregates import ( ... Average, ... Max, ... Sum, ... TimeHistogram, ... UniqueValues, ... ) >>> res = client.data_modeling.records.aggregate( ... stream_id="my-stream", ... aggregates={ ... "my_groups_by_1d_range": TimeHistogram( ... property=game_time, ... calendar_interval="1d", ... aggregates={ ... "my_groups_by_player_name": UniqueValues( ... property=player_name, ... aggregates={ ... "my_player_daily_scores_sum": Sum(points_scored), ... "my_player_daily_scores_maximum": Max(points_scored), ... }, ... ), ... "my_daily_scores_average": Average(points_scored), ... }, ... ), ... "my_scores_maximum_across_all_games": Max(points_scored), ... }, ... )
Bucket games by day and smooth the daily count with a 7-day moving average, using the
MovingFunctionsenum so the pipeline function name cannot be mistyped:>>> from cognite.client.data_classes.data_modeling.aggregates import ( ... Count, ... MovingFunction, ... MovingFunctions, ... TimeHistogram, ... ) >>> res = client.data_modeling.records.aggregate( ... stream_id="my-stream", ... aggregates={ ... "games_per_day": TimeHistogram( ... property=game_time, ... calendar_interval="1d", ... aggregates={ ... "games": Count(), ... "games_7d_avg": MovingFunction( ... buckets_path="games", ... window=7, ... function=MovingFunctions.UNWEIGHTED_AVG, ... ), ... }, ... ), ... }, ... )