Track timely master: un-shred arrange with multi-capability stamps - #838
Open
frankmcsherry wants to merge 3 commits into
Open
Track timely master: un-shred arrange with multi-capability stamps#838frankmcsherry wants to merge 3 commits into
frankmcsherry wants to merge 3 commits into
Conversation
Points the timely dependency at timely master, which stamps each message with a multiset of timestamps rather than exactly one. The adaptations are mechanical: Distributor implementations receive the stamp in place of a time and reproduce it on each produced sub-message, and capture events carry a stamp. All stamps remain singletons; behavior is unchanged. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01BrUdeCb6dsunVdk4acCPmh
The arrange operator retired capabilities one at a time, carving the batcher into per-capability tiles because each message could carry only one capability (its comment: 'Until timely dataflow supports multiple capabilities on messages, at least'). It now seals one batch per frontier advance and ships it under a CapabilitySet of the retiring capabilities; for totally ordered times the capability antichain has at most one element and behavior is unchanged. TraceReplayInstruction's capability hint becomes a Stamp (empty exactly for empty batches), and trace import replays batches under capability sets minted with delayed_stamp. Consumers accept multi-stamp batches: join retains the stamp and lower-bounds its unit's consolidation meet by the lattice meet of the stamp's elements; reduce, count, threshold, and arrange's own input retain each stamp element. The test drives two incomparable capabilities through arrange (one batch, stamped with both times), a concurrent trace import (replaying the same stamp), and join and reduce over the fused batch (correct outputs). Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01BrUdeCb6dsunVdk4acCPmh
All rounds are introduced before stepping, so nested iterative scopes hold many incomparable capabilities: the workload that distinguishes fused multi-capability batches from per-capability tiles. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01BrUdeCb6dsunVdk4acCPmh
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
Tracks timely master, whose messages are stamped by multisets of timestamps (TimelyDataflow/timely-dataflow#813), and removes the shredding of arranged batches into per-capability tiles.
Three commits:
master-next) plus mechanical adaptations:Distributorimplementations receive the stamp in place of a time and reproduce it on each produced sub-message, and capture events carry a stamp. All stamps remain singletons; behavior is unchanged.batcher.seal(frontier)shipped under aCapabilitySetof the retiring capabilities. For totally ordered times the capability antichain has at most one element and behavior is identical.TraceReplayInstruction's capability hint becomes aStamp(empty exactly for empty batches); trace import replays underdelayed_stamp. Consumers accept multi-stamp batches: join lower-bounds each unit's consolidation meet by the lattice meet of the stamp's elements (it must be the meet — any single element would be unsound for the history advance); reduce, count, threshold, and arrange's own input retain each stamp element. An end-to-end test drives two incomparable capabilities through arrange, a live trace import, join, and reduce.Measured on the pre-
master-nextport of this work (single worker, spreads under 1%): arrange fusion alone is worth 2–5% on scc_bench; fusing reduce's output formation as well reached 2.7–4.8x. That second step is deliberately not in this PR: reduce's output path is now organized around theReduceTacticcontract, whoseretirecurrently promises per-time tiles ("output must be ordered and tile [lower, upper)"), and fusing it properly means revisiting that contract across its implementations — better done as its own PR. Also still singleton-stamped: upsert, the CDC capture path, and dogsdogsdogs' half-joins.🤖 Generated with Claude Code
https://claude.ai/code/session_01BrUdeCb6dsunVdk4acCPmh