Repository navigation
[KVConnector][Tiering] Drive the tiering control plane from a dedicated thread - #58113
Closed
liranschour wants to merge 1 commit into
Closed
liranschour wants to merge 1 commit into
liranschour wants to merge 1 commit into
Conversation
…ed thread The secondary tiers' control plane only advanced from the per-step scheduler hooks: on_schedule_end polled every tier for finished jobs and then let each one serve inbound peer requests. The p2p tier's transports are non-blocking by design, so nothing moved between steps, and a rank busy with long prefills made every peer control-plane round trip wait for a step boundary — several round trips per transfer. TieringOffloadingManager now runs both from its own daemon thread, and serializes it against the scheduler with one lock held for a whole scheduler step rather than per call. Step granularity is required, not just tidier: the connector learns a chunk is a HIT in lookup() and only reads it in a later prepare_load() hook, and the chunk is evictable in between, so a promotion started while serving a peer lookup in that window could evict it out from under the pending load. The lock is released at the end of on_schedule_end, which puts the control plane's window over model execution — the time a rank mid-prefill has to spare. The read-mostly hooks (has_pending_work, get_stats, take_events) instead hold the lock only for their own duration. The engine ticks them every iteration, including iterations that schedule nothing and run no model, so folding them into the step-wide hold would leave an idle engine with no gap between steps at all — the state a producer waiting for a consumer to connect is in, and exactly when the control plane matters. Tiers re-enter the manager through _SecondaryTierFacingParent, whose caller already holds the lock, so the four ParentManager methods gained *_unlocked variants; exclude_tier_idx moved onto those, restoring the public signatures to match OffloadingManager. A sweep that raises is logged and re-raised on the next scheduler-side call, matching the engine crash these failures produced when they ran on the scheduler thread, instead of leaving the control plane silently dead. No interface changes, and no changes under tiering/p2p. The thread is on whenever a secondary tier is configured; control_plane_thread=False keeps everything on the calling thread, which is what the step-driven tests need. Signed-off-by: Liran Schour <lirans@il.ibm.com>
Contributor
|
This pull request has merge conflicts that must be resolved before it can be |
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.
Purpose
Addresses #55179. The P2P secondary tier only advances its control plane from the
per-step scheduler hooks:
TieringOffloadingManager.on_schedule_endpolls eachtier for finished jobs and then lets each one serve inbound peer requests. The
p2p tier's transports are non-blocking by design (
ZmqTransport.poll()returnsimmediately with no traffic), so nothing moves between steps and a peer's
control-plane round trip waits for a step boundary.
This moves that sweep onto a dedicated thread owned by the tiering manager.
Honest summary of what this does and does not buy
It does not, on its own, produce a measurable latency win. I benchmarked it
loaded (numbers below) and fetch RTT is parity-within-noise. I am submitting it
because the mechanism it removes is real and because the measurement localises
the actual bottleneck, not because it makes the benchmark faster.
The residual wait appears to be on the producer side: the fetch is not
waiting to be answered, it is waiting for the producer's scheduler thread to
reach
complete_store→tier.submit_storeand park the blocks. This changedeliberately leaves
_flush_pending_promotionsand the cascadesubmit_storecalls on the scheduler thread, so that wait is untouched. Getting the win
probably requires moving the producer-side block parking off the step boundary
as well.
Relationship to the existing PRs on this issue (per AGENTS.md)
This overlaps open work, and I want that on the record rather than buried:
the model"). Its
TieringOffloadingManager.poll_pending_workbody is the sametwo calls in the same order as this PR's
_sweep_control_plane. The fix is thesame; only the driver differs. [Core][KV Offload] Service the P2P control plane while the engine waits on the model #57165 adds
KVConnectorBase_V1.poll_pending_workand drives it fromEngineCoreduringthe window where the engine thread blocks on
future.result(), handing theconnector between threads under a lock. [Core][KV Offload] Service the P2P control plane while the engine waits on the model #57165 also bundles
add_remote_peer_async.same goal, also via
EngineCore+ the connector base API + executor changes.In #55179 @nilig asked for a read on "a thread touching the connector from
inside
EngineCore, even phase-exclusive, versus a control thread owned by thetier". This PR is a working implementation of the second option, offered as
input to that question:
EngineCorepoller threadTieringOffloadingManagerKVConnectorBase_V1.poll_pending_workIf the maintainers prefer either existing PR, close this one. It is up as a
concrete alternative and as a carrier for the measurement, not a land grab.
Design: why the lock is step-scoped
Per-call locking is not sufficient. The connector learns a chunk is a
HITfromlookup()(underget_num_new_matched_tokens) and only reads it in a laterprepare_load()hook (underupdate_state_after_alloc). In between the chunkstill has
ref_cnt == 0, i.e. it is evictable. A promotion started while servinga peer lookup in that window can evict it via
prepare_write→CPUOffloadingManager.prepare_store, andprepare_loadthentrips
assert chunk is not None.So the lock is taken on the first scheduler-side call of a step and released at
the end of
on_schedule_end. That places the control thread's window over modelexecution. The read-mostly hooks (
has_pending_work,get_stats,take_events) are the exception: they hold the lock only for their own duration,because the engine ticks them on every iteration including ones that schedule
nothing and run no model, and folding them into the step-wide hold leaves an idle
engine with no gap at all — the state a producer waiting for a consumer sits in.
Tiers re-enter through
_SecondaryTierFacingParent, whose caller already holdsthe lock, so the four
ParentManagermethods gained*_unlockedvariants.exclude_tier_idxmoved onto those, which restores the public signatures tomatch
OffloadingManager. A sweep that raises is logged and re-raised on thenext scheduler-side call, matching the engine crash these failures produced when
they ran on the scheduler thread rather than leaving the control plane silently
dead.
No interface changes, and no changes under
tiering/p2p/.Tests
Unit, on this branch (upstream main base):
Also run with the change cherry-picked onto the #57139 base for the live tests
below:
tests/v1/kv_offload/tiering/-> 387 passed, 11 skipped.Ten new tests in
tests/v1/kv_offload/tiering/test_tiering_offloading.pycover:sweeping with no further steps, sweeps running off the scheduler thread, the step
locking out the control thread,
on_schedule_endnot double-serving, a failedsweep re-raising on the scheduler thread,
shutdown()joining,reset_cache()with the thread live, idle-engine hooks not starving the thread, and the
disabled/no-tier cases.
Two of those were verified by mutation rather than assertion, since they are the
substantive claims: reverting
lookupto per-call locking failstest_step_locks_out_the_control_thread, and reverting_observation_lockto astep-wide hold fails
test_idle_engine_hooks_do_not_block_the_control_thread.Existing step-driven suites pass
control_plane_thread=Falseto staydeterministic; that flag is the only new knob and it is constructor-only.
Notes
not alter what is computed, only which thread services the tier control plane.
The correctness checks above are the serving-level evidence.
description. Per AGENTS.md it stays a draft until I have reviewed every changed
line and can defend it end to end.
🤖 Generated with Claude Code