Skip to content

[feat] add tq update and tq empty function and api - #182

Open
OutstanderWang wants to merge 5 commits into
Ascend:mainfrom
OutstanderWang:feat_tq_update
Open

OutstanderWang wants to merge 5 commits into
Ascend:mainfrom
OutstanderWang:feat_tq_update

Conversation

@OutstanderWang

@OutstanderWang OutstanderWang commented Oct 3, 2026 •

Copy link
Copy Markdown
Contributor

Summary

Adds kv_update, which rewrites selected fields of an existing key by running a user parser(old, new) on the SimpleStorage unit that holds the row, so the stored payload never makes a round trip through the client. kv_empty is the same call with empty=True: it stores None but 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-offloaded put_data is now two-phase, which also changes plain multi-field put (see "Behavior change for SSD offload").

Usage

kv_update: append the response to the prompt

import torch
import transfer_queue as tq

def concat(old, new):
    # old is None if this field was never written for the key
    return new if old is None else torch.cat([old, new])

tq.kv_put(key="s0", partition_id="train", fields={
    "input_ids": torch.tensor([10, 11, 12]),
    "attention_mask": torch.tensor([1, 1, 1]),
})

# Append the response in place; the field names stay the same.
tq.kv_update(key="s0", partition_id="train", fields=["input_ids", "attention_mask"],
             values={"input_ids": torch.tensor([20, 21]), "attention_mask": torch.tensor([1, 1])},
             parser=concat)
# input_ids      -> tensor([10, 11, 12, 20, 21])
# attention_mask -> tensor([1, 1, 1, 1, 1])

For a single field, pass fields="input_ids" and values=torch.tensor([20, 21]). The async variant is await tq.async_kv_update(...).

kv_empty: release pixel values after rollout

In 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_clear cannot release them, because it removes the whole key, and training still has to read the sample. kv_empty drops just that field:

for key in keys:
    tq.kv_put(key=key, partition_id="train", fields={
        "input_ids": input_ids[key],
        "pixel_values": pixel_values[key],  # large; only vLLM reads it
        "mm_tensors": mm_tensors[key],      # what Megatron-Core trains on
    })

# ... vLLM rollout reads pixel_values and writes the responses back ...

# Free the pixels; the key, its other fields and its production status are kept.
for key in keys:
    tq.kv_empty(key=key, partition_id="train", fields="pixel_values")

# Training reads input_ids / responses / mm_tensors exactly as before.

After kv_empty:

  • The storage unit no longer holds the tensor. With SSD offload on, its SSD file is also deleted.
  • The controller marks pixel_values as a non-tensor field, so metadata no longer advertises its old shape.
  • Reading the field returns 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) -> KVBatchMeta
  • kv_empty(key, partition_id, fields) -> KVBatchMeta
  • async_kv_update / async_kv_empty with the same signatures. All four are exported from transfer_queue.

Argument rules:

  • The key must already exist; a missing key raises ValueError and is not created.
  • Without empty, both a callable parser and values are required. There is no default parser.
  • With empty=True, passing values or parser is rejected.
  • fields is one field name or a list; for a list, values must be a dict naming exactly those fields.
  • Nested tensors are rejected as values for a single-key update.
  • The returned KVBatchMeta lists every field stored for the key.

How the parser runs:

  • Inside the storage unit process, once per sample per field.
  • It is shipped with cloudpickle, so lambdas and closures work, but they run on a copy of their captured state; mutating driver-side variables has no visible effect.
  • old is None when the field was never written for that key.
  • SSD-offloaded values reach the parser decoded, not as file references.

2. Data path (data plane)

client.update → SimpleStorageManager.update_data → new ZMQ request UPDATE_DATA (with UPDATE_DATA_RESPONSE / UPDATE_ERROR) → SimpleStorageUnit._handle_update → StorageUnitData.apply_update.

  • Compute first, then write once. apply_update reads old values through the polymorphic get_data and computes every parser result before a single put_data call, so a raise leaves the row unchanged.
  • Shapes reported by units, assembled by the manager. Each unit describes only the per-sample shapes it stored. _build_update_field_schema assembles one batch-level field_schema in metadata.global_indexes order, because the controller reads per_sample_shapes by position. Mixed lengths make the column nested; if any unit stored non-tensors, the whole column is non-tensor.
  • The manager notifies the controller with the existing notify_data_update, and the client merges the schema into its BatchMeta via the new BatchMeta.apply_field_schema (add_fields now uses the same helper).
  • No replays for parser updates. These requests use 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 what put already does with a data_parser.
  • UPDATE_DATA is added to the unit's per-operation metrics.

3. Controller metadata (control plane)

FieldMeta.update gains 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 sets is_non_tensor=True and is_nested=False, clears shape and per_sample_shapes, and keeps dtype (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_data raises NotImplementedError("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 run parser where 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_data used 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. Plain put converges on retry, but kv_update parsers such as concat are not idempotent and are not retried. put_data is now two-phase:

  1. Write every field's SSD files. On any failure, unlink only the files this call wrote and re-raise; storage is unchanged.
  2. Swap all references, update ssd_active_values / ssd_active_bytes, and unlink superseded files.

The base StorageUnitData.put_data now checks every field's length before writing any field, so the single commit cannot fail halfway either.

This also affects plain multi-field put with 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

  • README: (async_)kv_update added to the KV API list.
  • tutorial/02_kv_interface.py: a new step that concatenates a field with kv_update and then kv_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 sees old=None for a new field; several fields in one call; a non-tensor payload; kv_empty keeps 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's per_sample_shapes); updates from remote Ray workers. It refuses to run if a TransferQueueController is 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 tests and ruff are clean.
  • Full local suite, 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_update concat, atomicity on parser error, empty, missing field passes old=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 injected ENOSPC is 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:

  • kv_update unit tests plus the KV interface e2e: 83 passed, 2 skipped (the skips are the non-SimpleStorage rejection e2e, skipped by design under SimpleStorage).
  • 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's CONFIG was overridden with ssd_offload: {enabled: true, path: <local dir>, threshold_bytes: 16} so its small tensors were actually offloaded; SSD file counts matched each unit's reported ssd_active_values, and the SSD directory was empty after close().
  • Before rebasing onto the current main, the full suite and the recipe also passed on two-node clusters.

Notes for reviewers / out of scope

  • kv_update updates one key per call. For a batch, use client.kv_retrieve_meta(keys=..., partition_id=..., create=False) followed by client.update(metadata, field_names, values=..., parser=...); the recipe's cross-unit check uses this path.
  • Parsers should return a new value rather than modify old in place.
  • Atomicity holds per storage-unit request; a batch spanning several units is not a cross-unit transaction.
  • No change to put / get semantics for KV backends, and no change to SimpleStorage behavior when SSD offload is off.

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 OutstanderWang changed the title [WIP] feat: add tq update function and api feat: add tq update and tq empty function and api Oct 3, 2026
@OutstanderWang OutstanderWang changed the title feat: add tq update and tq empty function and api [feat] add tq update and tq empty function and api Oct 3, 2026
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant