Ray's Knowledge Base

Aggregation tier state that spans builds and processes

FactVerified 27 Sep 2026Holds project: dftracer-utils
Fact. A statement and the evidence for it.

Statement#

Besides its rows, the aggregation tier (dftracer.agg) keeps state that must stay consistent across builds, ranks and processes. Keep these rules when you change it:

  • Time bounds (__time_bounds__): an incremental build that aggregates only new files must widen the stored bounds, never overwrite them. The old persist_time_bounds put only this build's min and max, so adding one file shrank the tier's bounds. Stage 10a uses agg::tier::widen_time_bounds.
  • Intern dictionary watermark: intern_for_index(path) shares one intern table per index path in the process, with flushed_entries saying which entries are on disk. When code clears the tier on disk, it must reset that watermark to 0 (agg::tier::clear does), or later rows use ids whose dictionary entries are gone.
  • Config across shards: merge_shard_set must refuse shards aggregated with different params. Before stage 10a it kept the first shard's config and merged every shard's rows, which mixes intervals silently.
  • One converter: the distributed workers and the coordinator must build the native AggregationConfig from the Python object with the same function (agg_config_from_py), so their params hashes agree.
  • Rows and entries in one write: the fold puts each file's dftracer.agg manifest entry in the same IndexWrite as its merge operands, so an entry never exists without its rows, also through distributed SSTs.

See also the aggregation tier holds no file id.

Evidence#

  • Found while making the tier the dftracer.agg extension (stage 10a), by reading event_aggregator.cpp (persist_time_bounds), aggregation_intern.cpp (intern_for_index), sharded_view.cpp (merge_shard_set) and sst_distribution.cpp.
  • The stage 10a tests (index/test_agg_extension, the aggregators, sharded view and Python distributed tests) pass with these rules. The time-bounds and watermark bugs were not reproduced by a test before the fix.