[feat] add tq update and tq empty function and api - #182
Open
OutstanderWang wants to merge 5 commits into
Open
OutstanderWang wants to merge 5 commits into
OutstanderWang wants to merge 5 commits into
Conversation
Cover argument normalization, unit-side parser/empty application, cross-unit field_schema assembly, controller non-tensor marking, and KV backend rejection. Signed-off-by: OutstanderWang <wangweiyanster@gmail.com>
kv_update rewrites selected fields of an existing key by running a custom parser(old, new) on the SimpleStorage unit that holds the row, so the payload never travels to the client and back. kv_empty stores None instead, keeping the key and its production status. The unit computes every value before writing, so a failing parser leaves the row unchanged, and it reports only the shape it stored per sample; the manager assembles one batch-level field_schema ordered to match the notified indexes. FieldMeta now accepts a column turning non-tensor, so an emptied field stops being described as a tensor. KV-based backends reject the call. Signed-off-by: OutstanderWang <wangweiyanster@gmail.com>
kv_update needs to run parser(old, new) where the value lives, which only SimpleStorage can do. The refusal used to list backend names in a string that would go stale as soon as another KV backend was added, and it never told the caller which backend they were actually on. Report the concrete manager class instead, and cover Mooncake, Yuanrong and Ray individually rather than only the shared base class. Signed-off-by: OutstanderWang <wangweiyanster@gmail.com>
The unit tests stub out the transport, so they cannot show that a parser really travels to the storage unit or that a batch spanning several units reports its per-sample shapes back to the controller in the right order. This script exercises both against a live Ray cluster, alongside kv_empty, multi-field updates, non-tensor payloads, parser failure atomicity, argument validation, the async entry points, and updates issued from remote workers. It refuses to run when a controller already exists, so it can never attach to somebody else's deployment and clear their keys on the way out. Signed-off-by: OutstanderWang <wangweiyanster@gmail.com>
With SSD offload on, put_data wrote, committed and cleaned up one field at a time. If writing a later field's files failed (for example a full disk), the earlier fields were already replaced and their old files deleted, yet the caller got an error. A plain put converges on retry, but kv_update parsers such as concat are not idempotent and the manager does not retry updates, so retrying applied the earlier fields twice. Write every field's files first and only then swap all references, update the SSD accounting and delete superseded files; any failure removes just the files this call wrote. The base put_data now checks every field's length before writing any, so the single commit cannot fail halfway either. apply_update's "a raise leaves storage unchanged" now holds for SSD write errors as well as parser errors. Signed-off-by: OutstanderWang <wangweiyanster@gmail.com>
OutstanderWang
force-pushed
the
feat_tq_update
branch
from
October 3, 2026 10:52
f735236 to
caf3753
Compare
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
Adds
kv_update, which rewrites selected fields of an existing key by running a userparser(old, new)on the SimpleStorage unit that holds the row, so the stored payload never makes a round trip through the client.kv_emptyis the same call withempty=True: it storesNonebut keeps the key and its production status, so a large field can be released as soon as its last consumer is done (for example, pixel values after rollout) while the rest of the sample stays available for training. Each request is all-or-nothing on a unit: if a parser raises or an SSD write fails, storage is left unchanged. To get that guarantee, SSD-offloadedput_datais now two-phase, which also changes plain multi-fieldput(see "Behavior change for SSD offload").Usage
kv_update: append the response to the promptFor a single field, pass
fields="input_ids"andvalues=torch.tensor([20, 21]). The async variant isawait tq.async_kv_update(...).kv_empty: release pixel values after rolloutIn a multimodal RL step, the rollout engine (vLLM) needs the raw
pixel_values, while training (Megatron-Core) only needs the processed tensors. Once rollout is done, the pixel values are dead weight in storage.kv_clearcannot release them, because it removes the whole key, and training still has to read the sample.kv_emptydrops just that field:After
kv_empty:pixel_valuesas a non-tensor field, so metadata no longer advertises its old shape.None.To release many keys concurrently,
await asyncio.gather(*(tq.async_kv_empty(key=k, partition_id="train", fields="pixel_values") for k in keys)).What's here
1. Public API (
transfer_queue/interface.py,__init__.py)kv_update(key, partition_id, fields, *, values=None, parser=None, empty=False) -> KVBatchMetakv_empty(key, partition_id, fields) -> KVBatchMetaasync_kv_update/async_kv_emptywith the same signatures. All four are exported fromtransfer_queue.Argument rules:
ValueErrorand is not created.empty, both a callableparserandvaluesare required. There is no default parser.empty=True, passingvaluesorparseris rejected.fieldsis one field name or a list; for a list,valuesmust be a dict naming exactly those fields.KVBatchMetalists every field stored for the key.How the parser runs:
oldisNonewhen the field was never written for that key.2. Data path (data plane)
client.update→SimpleStorageManager.update_data→ new ZMQ requestUPDATE_DATA(withUPDATE_DATA_RESPONSE/UPDATE_ERROR) →SimpleStorageUnit._handle_update→StorageUnitData.apply_update.apply_updatereads old values through the polymorphicget_dataand computes every parser result before a singleput_datacall, so a raise leaves the row unchanged._build_update_field_schemaassembles one batch-levelfield_schemainmetadata.global_indexesorder, because the controller readsper_sample_shapesby position. Mixed lengths make the column nested; if any unit stored non-tensors, the whole column is non-tensor.notify_data_update, and the client merges the schema into itsBatchMetavia the newBatchMeta.apply_field_schema(add_fieldsnow uses the same helper).max_attempts=1:parser(old, new)reads stored state, so a retry after a lost reply would apply the new value twice. Empty updates are plain overwrites and keep the normal retry policy. This mirrors whatputalready does with adata_parser.UPDATE_DATAis added to the unit's per-operation metrics.3. Controller metadata (control plane)
FieldMeta.updategains a non-tensor branch. Before, an incoming non-tensor schema left an existing tensor field unchanged, so an emptied column was still described as a tensor with its old shape. Now it setsis_non_tensor=Trueandis_nested=False, clearsshapeandper_sample_shapes, and keepsdtype(a later tensor write must still match the field's original dtype). The tensor branch is unchanged apart from indentation. The controller only receives schema; payloads stay in the storage units.4. Backend support: SimpleStorage only
KVStorageManager.update_dataraisesNotImplementedError("kv_update is not supported by <ManagerClass>; it requires the SimpleStorage backend.")for Mooncake, Yuanrong, Ray and generic KV backends. A KV backend stores each sample-field under its own key and has no hook to runparserwhere the value lives, so read-modify-write cannot be made atomic. The message names the concrete manager class rather than a hard-coded list that would go stale.5. Behavior change for SSD offload (#162)
HybridStorageUnitData.put_dataused to write, commit and clean up one field at a time. If a later field failed (for example, a full disk), earlier fields were already replaced and their old files deleted, yet the caller got an error. Plainputconverges on retry, butkv_updateparsers such as concat are not idempotent and are not retried.put_datais now two-phase:ssd_active_values/ssd_active_bytes, and unlink superseded files.The base
StorageUnitData.put_datanow checks every field's length before writing any field, so the single commit cannot fail halfway either.This also affects plain multi-field
putwith SSD offload on: a failure on field 2 no longer leaves field 1 committed. Reviewers of the SSD offload code should look at this part.6. Docs, tutorial and recipe
(async_)kv_updateadded to the KV API list.tutorial/02_kv_interface.py: a new step that concatenates a field withkv_updateand thenkv_emptys it; later steps renumbered.recipe/kv_update_e2e.py: an end-to-end script for a real Ray cluster, SimpleStorage with 4 units so batches actually split across units. 10 checks: concat keeps the field name; the parser seesold=Nonefor a new field; several fields in one call; a non-tensor payload;kv_emptykeeps the key and updates the controller; a failed parser leaves the row unchanged; rejected arguments;async_kv_update/async_kv_empty; a multi-sample update across units (checks the controller'sper_sample_shapes); updates from remote Ray workers. It refuses to run if aTransferQueueControlleris already deployed, so it never attaches to someone else's deployment or clears their keys on exit.Testing
python -m compileall -q transfer_queue tutorial testsandruffare clean.python -m pytest -q: 692 passed, 12 skipped (the Yuanrong e2e, which needs a Yuanrong deployment, was excluded).New and updated tests:
tests/test_kv_update.py: argument normalization (test_normalize_*);apply_updateconcat, atomicity on parser error, empty, missing field passesold=None; an SSD-offloaded old value is decoded before the parser sees it; atomicity when a later field fails to reach SSD; cross-unit schema assembly (test_build_update_field_schema_*); controller non-tensor marking; Mooncake, Yuanrong, Ray and base KV managers reject the call with their class name.tests/test_simple_storage_unit.py::test_failed_ssd_write_on_second_field_leaves_every_field_unchanged: a multi-field put with an injectedENOSPCis all-or-nothing.tests/e2e/test_kv_interface_e2e.py::TestKVUpdateE2E(sync and async): concat then empty; concatenating prompt and response keeps the field name; several fields at once; missing key / empty with values; KV-backend rejection.On a Ray cluster (Ray 2.51, Python 3.13) at the branch head:
tests/e2e/test_ssd_offload_e2e.py: 2 passed.recipe/kv_update_e2e.py: 10/10 with SSD offload off (the committed config) and on. For the offload run, the recipe'sCONFIGwas overridden withssd_offload: {enabled: true, path: <local dir>, threshold_bytes: 16}so its small tensors were actually offloaded; SSD file counts matched each unit's reportedssd_active_values, and the SSD directory was empty afterclose().main, the full suite and the recipe also passed on two-node clusters.Notes for reviewers / out of scope
kv_updateupdates one key per call. For a batch, useclient.kv_retrieve_meta(keys=..., partition_id=..., create=False)followed byclient.update(metadata, field_names, values=..., parser=...); the recipe's cross-unit check uses this path.oldin place.put/getsemantics for KV backends, and no change to SimpleStorage behavior when SSD offload is off.