Repository navigation
ES|QL: Chunk time-series aggregation output - #151670
leontyevdv merged 8 commits into
Conversation
|
Pinging @elastic/es-storage-engine (Team:StorageEngine) |
|
Hi @leontyevdv, I've created a changelog YAML for you. |
🔍 Preview links for changed docs⏳ Building and deploying preview... View progress This comment will be updated with preview links when the build is complete. |
ℹ️ Important: Docs version tagging👋 Thanks for updating the docs! Just a friendly reminder that our docs are now cumulative. This means all 9.x versions are documented on the same page and published off of the main branch, instead of creating separate pages for each minor version. We use applies_to tags to mark version-specific features and changes. Expand for a quick overviewWhen to use applies_to tags:✅ At the page level to indicate which products/deployments the content applies to (mandatory) What NOT to do:❌ Don't remove or replace information that applies to an older version 🤔 Need help?
|
9d92715 to
6954a81
Compare
6954a81 to
498a649
Compare
|
Hi @leontyevdv, I've updated the changelog YAML for you. |
|
Thanks Dima! I have some high-level feedback:
|
94607b1 to
4244259
Compare
|
Thanks Nhat!
Done! So, I understood there are two different things - periodic emit and chunking. I disabled the periodic emit.
We now materialize DimensionValues's values block once and copy each page out of it.
I added a per-page tsid dictionary there.
Removed! |
| timestamps.appendLong(p, finalHash.getKey2(p)); | ||
| final int groupId = selected.getInt(p); | ||
| final long globalOrd = finalHash.getKey1(groupId); | ||
| if (localOrd < 0 || globalOrd != prevGlobalOrd) { |
There was a problem hiding this comment.
Here we rely on the ordering of TSIDs; however, this ordering is not always guaranteed. Can you keep a map of ordinals instead of solely relying on the previous ordinal?
There was a problem hiding this comment.
Done! Added a test too: testTimeSeriesBlockHashChunkedPageDeduplicatesUnsortedTsids
| final Block[] blocks; | ||
| if (OrdinalBytesRefBlock.isDense(positionCount, tsidHash.size())) { | ||
| blocks = buildOrdinalKeys(positionCount); | ||
| blocks = buildOrdinalKeys(selected); |
There was a problem hiding this comment.
Can we have a non-chunking version that reuses the full dictionary as we do today?
There was a problem hiding this comment.
Yes, I returned the existing (in main) logic for cases when there is only one page (non-chunked). Thanks for pointing out on that!
| @@ -241,26 +243,29 @@ public GroupingAggregatorFunction.PreparedForEvaluation prepareEvaluateIntermedi | |||
|
|
|||
| private void evaluate(Block[] blocks, int offset, IntVector selectedInPage) { | |||
| int positionCount = selectedInPage.getPositionCount(); | |||
There was a problem hiding this comment.
I think we lose the optimization for non-chunked cases. How about something like this:
@Override
public GroupingAggregatorFunction.PreparedForEvaluation prepareEvaluateIntermediate(
IntVector selected,
GroupingAggregatorEvaluationContext ctx
) {
values = builder.build();
return (blocks, offset, selectedInPage) -> {
if (selected == selectedInPage && values.getPositionCount() == selected.getPositionCount()) {
values.incRef();
blocks[offset] = values;
return;
}
final int positionCount = selectedInPage.getPositionCount();
final BytesRef scratch = new BytesRef();
try (var outputBuilder = driverContext.blockFactory().newBytesRefBlockBuilder(positionCount)) {
for (int p = 0; p < positionCount; p++) {
int groupId = selectedInPage.getInt(p);
if (groupId < values.getPositionCount()) {
outputBuilder.copyFrom(values, groupId, scratch);
} else {
outputBuilder.appendNull();
}
}
blocks[offset] = outputBuilder.build();
}
};
}
@Override
public GroupingAggregatorFunction.PreparedForEvaluation prepareEvaluateFinal(
IntVector selected,
GroupingAggregatorEvaluationContext ctx
) {
return prepareEvaluateIntermediate(selected, ctx);
}There was a problem hiding this comment.
Thanks! I use your method almost as is but instead of selected == selectedInPage I added an explicit check that the groups are not reordered or dropped by TimeSeriesAggregationOperator.selectedForValuesAggregator when window is not an exact multiple of the bucket. See this method. Does that make sense?
| final boolean outputFinal = aggregatorMode.isOutputPartial() == false; | ||
| final boolean outputPartial = aggregatorMode.isOutputPartial(); | ||
| final boolean outputFinal = outputPartial == false; | ||
| final int effectiveTargetChunkRows = outputPartial ? targetChunkRows : Integer.MAX_VALUE; |
There was a problem hiding this comment.
I think we can enable chunking in both partial and final, but let’s enable it for partial first as this PR does.
72f2104 to
d87704f
Compare
dnhatn
left a comment
There was a problem hiding this comment.
Thanks for all iterations, Dima!
| * each page sent to the coordinator. Specific to the time-series operator and independent of the regular | ||
| * aggregation emit settings. | ||
| */ | ||
| public static final Setting<Integer> TIME_SERIES_TARGET_CHUNK_SIZE = Setting.intSetting( |
There was a problem hiding this comment.
Can we name this setting with ROWS instead of SIZE? Right now we use "row" and "size" interchangeably, but at some point we'll need "size" to refer to actual byte size.
| final int groupId = selected.getInt(p); | ||
| final int globalOrd = (int) finalHash.getKey1(groupId); | ||
| final int localOrd; | ||
| if (globalToLocalOrd.containsKey(globalOrd)) { |
There was a problem hiding this comment.
Nit: Maybe use indexOf / indexGet / indexInsert for a single lookup/insert instead.
There was a problem hiding this comment.
Yes, it's globalToLocalOrd.indexExists(slot). Fixed this and the insert part. Thanks!
833e961 to
259c4d4
Compare
|
Here are benchmarking results. scenario: v1.0.2-hostmetrics_8clientsort-270m avg_avgot_memory_by_host_5m avg_rate_cpu_by_host_5m sum_rate_sys_cpu_time_large_clause_5m Q_MEM —
|
| Config | median (ms) | min | max | exch_pages | Chunking (exp_pages) |
|---|---|---|---|---|---|
| main · no-pragma | 77 | 60 | 367 | 12 | off (main) |
| branch · no-pragma | 44 | 40 | 51 | 12 | off |
| branch · chunk=2147483647 | 42 | 40 | 47 | 12 | not-chunked (1) |
| branch · chunk=100000 | 41 | 39 | 50 | 12 | not-chunked (1) |
| branch · chunk=50000 | 39 | 38 | 45 | 12 | not-chunked (1) |
| branch · chunk=10000 | 41 | 38 | 45 | 12 | CHUNKED (3) |
| branch · chunk=1000 | 41 | 38 | 58 | 37 | CHUNKED (29) |
Q_RATE — avg_rate_cpu_by_host_5m (G = 192,000)
| Config | median (ms) | min | max | exch_pages | Chunking (exp_pages) |
|---|---|---|---|---|---|
| main · no-pragma | 291 | 251 | 387 | 12 | off (main) |
| branch · no-pragma | 240 | 234 | 276 | 12 | off |
| branch · chunk=2147483647 | 239 | 233 | 265 | 12 | not-chunked (1) |
| branch · chunk=100000 | 235 | 227 | 266 | 12 | CHUNKED (2) |
| branch · chunk=50000 | 241 | 229 | 254 | 12 | CHUNKED (4) |
| branch · chunk=10000 | 238 | 231 | 261 | 25 | CHUNKED (20) |
| branch · chunk=1000 | 242 | 232 | 271 | 199 | CHUNKED (192) |
Q_WIDE — sum_rate_sys_cpu_time_large_clause_5m (G = 192,000)
| Config | median (ms) | min | max | exch_pages | Chunking (exp_pages) |
|---|---|---|---|---|---|
| main · no-pragma | 499 | 484 | 518 | 12 | off (main) |
| branch · no-pragma | 477 | 465 | 505 | 12 | off |
| branch · chunk=2147483647 | 476 | 466 | 501 | 12 | not-chunked (1) |
| branch · chunk=100000 | 473 | 462 | 507 | 12 | CHUNKED (2) |
| branch · chunk=50000 | 477 | 463 | 507 | 12 | CHUNKED (4) |
| branch · chunk=10000 | 473 | 465 | 506 | 24 | CHUNKED (20) |
| branch · chunk=1000 | 479 | 468 | 510 | 197 | CHUNKED (192) |
Summary — main vs branch (chunking off) and chunking overhead
| Query | main median | branch off (chunk=MAX) | branch default (100000) | branch heavy (1000) | median spread across branch sweep |
|---|---|---|---|---|---|
| Q_MEM | 77 | 42 | 41 | 41 | 39–44 (~5 ms) |
| Q_RATE | 291 | 239 | 235 | 242 | 235–242 (~7 ms) |
| Q_WIDE | 499 | 476 | 473 | 479 | 473–479 (~6 ms) |
The time-series aggregation operator previously emitted its entire intermediate output as a single page, so a data node could send one jumbo page to the coordinator and trip the circuit breaker while merging. Enable chunking of the partial/intermediate output: - TimeSeriesAggregationOperator emits partial results periodically and slices each emission into _tsid-aligned pages: it fills up to targetChunkRows, then cuts at the next _tsid change so a _tsid that starts in a page is sent whole. maxChunkRows (2x the target) is a hard ceiling that only splits a single oversized _tsid. Final mode is unchanged. - Add a nextPageSliceEnd hook on HashAggregationOperator so the slicing policy can be specialized by the time-series subclass. - Fix TimeSeriesBlockHash#getKeys to project only the selected group ids instead of always returning all keys, which the slicer needs. Decouple the time-series knobs from the regular aggregation settings: add esql.time_series.partial_agg_emit_keys_threshold, esql.time_series.partial_agg_emit_unique_threshold and esql.time_series.target_chunk_size cluster settings, read only by the time-series factory (with matching query pragmas). Defaults of 100,000 / 0.1 / 100,000 leave regular STATS aggregations untouched.
Chunk the time-series aggregation's intermediate output so data nodes send bounded pages to the coordinator instead of one jumbo page, while keeping the final (coordinator) path unchanged. Drop the periodic-emit and tsid-alignment and reduce the knobs to a single chunk size.
Address review feedback
259c4d4 to
592dbe6
Compare
The time-series aggregation operator previously emitted its entire intermediate output as a single page, so a data node could send one jumbo page to the coordinator and trip the circuit breaker while merging.
This PR enables chunking of the partial/intermediate output of the time-series aggregation.
TimeSeriesAggregationOperatornow slices its partial/intermediate output into pages of abouttargetChunkRowsrows, bounding the size of each page sent to the coordinator. A_tsidmay straddle a page boundary; the coordinator re-merges groups by key, so no_tsid-aligned slicing is needed. Final-mode output is not chunked.Chunking is controlled by a single dedicated, time-series-only knob, decoupled from the regular aggregation emit settings. Cluster setting
esql.time_series.target_chunk_rows(dynamic, default 100,000), with a matching per-query pragma of the same name.Closes #147286