[feat] selective saving and restoring of TQ samples by key instead of the whole - #181
Open
OutstanderWang wants to merge 22 commits into
Open
OutstanderWang wants to merge 22 commits into
OutstanderWang wants to merge 22 commits into
Conversation
Squashes the initial by-key checkpoint work with its follow-up refactors (helper extraction, restructured flow, fsync on publish). Signed-off-by: OutstanderWang <wangweiyanster@gmail.com>
Signed-off-by: OutstanderWang <wangweiyanster@gmail.com>
Signed-off-by: OutstanderWang <wangweiyanster@gmail.com>
Signed-off-by: OutstanderWang <wangweiyanster@gmail.com>
Signed-off-by: OutstanderWang <wangweiyanster@gmail.com>
Signed-off-by: OutstanderWang <wangweiyanster@gmail.com>
Signed-off-by: OutstanderWang <wangweiyanster@gmail.com>
Signed-off-by: OutstanderWang <wangweiyanster@gmail.com>
Signed-off-by: neowywang <neowywang@tencent.com>
Signed-off-by: neowywang <neowywang@tencent.com>
Signed-off-by: neowywang <neowywang@tencent.com>
Signed-off-by: neowywang <neowywang@tencent.com>
Signed-off-by: neowywang <neowywang@tencent.com>
Signed-off-by: neowywang <neowywang@tencent.com>
Signed-off-by: neowywang <neowywang@tencent.com>
Signed-off-by: neowywang <neowywang@tencent.com>
Signed-off-by: neowywang <neowywang@tencent.com>
…meout The socket pool applies its timeout when a socket connects, so patching TQ_SIMPLE_STORAGE_SEND_RECV_TIMEOUT after tq.init() no longer reached the sockets the load reuses and the request waited out the paused unit. Signed-off-by: OutstanderWang <wangweiyanster@gmail.com>
… refactor The controller and client now route requests through a dispatch table and a shared _request_controller helper, so the copies carried over from the old elif chain shadowed those and left dead code behind. Signed-off-by: OutstanderWang <wangweiyanster@gmail.com>
CLA Signature PassOutstanderWang, thanks for your pull request. All authors of the commits have signed the CLA. 👍 |
OutstanderWang
force-pushed
the
save-by-key-on-main
branch
from
September 29, 2026 07:36
c2fa8f6 to
41a562e
Compare
CLA Signature PassOutstanderWang, thanks for your pull request. All authors of the commits have signed the CLA. 👍 |
Signed-off-by: neowywang <neowywang@tencent.com>
Signed-off-by: neowywang <neowywang@tencent.com>
Signed-off-by: neowywang <neowywang@tencent.com>
CLA Signature PassOutstanderWang, thanks for your pull request. All authors of the commits have signed the CLA. 👍 |
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.
Summary
This PR adds a way to save a chosen set of keys from TransferQueue and later merge them back into a running deployment:
tq.dump_data_by_key(dump_dir, keys, partition_id)saves the selected keys, all fields produced for them, and their tags.tq.load_data_by_key(dump_dir)merges those rows back into the running system.The controller handles only metadata: which keys, indexes and storage units are involved, the field schemas, and restore reservations. All payload bytes are read and written by the storage units that own them, in parallel. Payloads never pass through the client.
In a real RL training job with partial (streaming) rollout, saving the rollout state at a training checkpoint went from 3.00 s to 0.90 s, about 3.3× faster. Loading was also faster even though the restored state was about 1.5× larger.
Motivation
The existing checkpoint saves and restores the whole system.
save_checkpoint/load_checkpointcapture the entire controller (every partition, the global index manager and the sampler) plus each storage unit's full contents. On load, the storage-unit count must match and the system must be in a clean state. That is right for restarting a whole job. It is the wrong tool for "put these N samples back into the TransferQueue I am already running".Partial rollout needs exactly that. With partial or streaming rollout, rollout samples that are unfinished or not yet consumed live only in TransferQueue when the trainer takes a checkpoint. If they are lost, resuming throws away generation work that was already paid for. So at every checkpoint the trainer must save exactly those keys next to its model checkpoint, and on resume it must put them back:
The save also runs on the training critical path at every checkpoint, so it has to be cheap.
The current workaround is slow because all payloads go through one process. The workaround (called V1 below):
kv_batch_getfor the keys on the client and serialize the result to disk.kv_batch_put.Every payload byte travels from the storage units to the client over ZMQ and is deserialized there. A single client process then serializes all of it to disk. Restore reverses this, so the payload crosses the client twice per direction. The cost grows with the buffer size and blocks the trainer while it runs.
This PR (V2) keeps the client on the metadata path only. Each owner storage unit writes its own rows during a dump and reads only its own rows during a restore, and all units work in parallel.
Usage
How
load_data_by_keymerges into the running system:dump_dirmust be on a filesystem that every participating storage unit can reach (any local directory works on a single node).What's here
Public API (
transfer_queue/data_dump.py, exported fromtransfer_queue)dump_data_by_key,load_data_by_key,read_row_index,recover_data_loadandRestorePendingError.<dir>.tmpand synced. The previous dump is moved to<dir>.oldand is deleted only after the new directory has been published and its parent directory synced. If a crash leaves the main directory missing, the next access recovers.old..lockfile, so this also works across processes. A sibling.restoremarker stops a dump from being replaced while a load may still be reading it.Controller (control plane only)
DESCRIBE_ROWS_BY_KEYreturns each key's global index, produced fields, tag and original field schema, with no payload.VALIDATE_DUMP_SCHEMArejects destination type conflicts before any payload is written.BEGIN_RESTORE,RESTORE_UNITclaim/complete,FINISH_RESTORE,LIST_RESTORES). Each load has a unique ID. A storage unit must claim permission before writing.RestorePendingErrorexplains why a load is still unresolved: per-unit states (pending/running), unit report failures, or an unknown ID after a restart. The controller remembers final outcomes, so retrying recovery after a lost commit reply is safe.Storage manager and SimpleStorage units (data plane)
utils/storage_routing.py.DUMP_ROWSandLOAD_ROWSare added to the storage operation metrics.Schema fidelity
format_version: 3, which saves the original field schemas. A restore reproduces those schemas regardless of target topology or batch boundaries; for example, a non-tensor field stays non-tensor.NonTensorStackand leave nested row shapes missing. For those fields, the owner units recover the shape and dtype from the actual values while writing. If the values are mixed, the field is saved as non-tensor, with a warning.Fallbacks
kv_batch_put.Docs
docs/data_dump.mdcovers the format, distributed I/O, failure handling and recovery.docs/checkpoint.mdgains a pointer explaining when to use a selective dump instead of a full checkpoint.Performance
Setup: a real VLM RL training job with partial (streaming) rollout on 64 H20 GPUs, resumed twice in a row:
V1 is the client-side path described in Motivation. V2 is
dump_data_by_key/load_data_by_key.Correctness was also checked at a larger scale (hundreds of GiB, about 15k keys, 32 shards). After restore, a sample of keys matched in both content and tags, and dumps with missing fields or a missing manifest were rejected.
Testing
python -m compileall -q transfer_queue tutorial tests: clean.tests/test_data_dump.py,tests/test_dump_lock.py,tests/test_restore_lifecycle.pytests/e2e/test_data_dump_e2e.py,tests/e2e/test_data_dump_cross_topology_e2e.py,tests/e2e/test_restore_timeout_e2e.pypython -m pytest -q tests --ignore=tests/test_yuanrong_storage_client_e2e.py: 806 passed, 10 skipped.What the new tests cover: