Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
22 commits
Select commit Hold shift + click to select a range
e5c2984
feat: save selective TransferQueue checkpoints by key
OutstanderWang Sep 29, 2026
6d51c3e
fix: isolate tensor storage in selective dumps
OutstanderWang Sep 23, 2026
1090202
fix: preserve previous dump until publication succeeds
OutstanderWang Sep 23, 2026
16b341d
test: use portable selective dump directories
OutstanderWang Sep 23, 2026
3a07240
feat: restore selective dumps directly on owner storage units
OutstanderWang Sep 23, 2026
63b7cfd
fix: preserve original schemas in selective dump restores
OutstanderWang Sep 24, 2026
579ed35
fix: reserve restore indexes until remote operations settle
OutstanderWang Sep 24, 2026
06898da
fix: serialize dump publication and recovery across processes
OutstanderWang Sep 24, 2026
461109e
fix: use union syntax in selective dump test
OutstanderWang Sep 28, 2026
99ff73a
fix: narrow validated dump indexes to integers
OutstanderWang Sep 28, 2026
d51882d
fix: aggregate dump results after checking exceptions
OutstanderWang Sep 28, 2026
2e08ab8
fix: check backend recovery support before cancelling restores
OutstanderWang Sep 28, 2026
704c325
fix: scope strict schema validation to selective restores
OutstanderWang Sep 28, 2026
3a603fb
fix: preserve pending restores after receive timeouts
OutstanderWang Sep 28, 2026
87e5ac7
refactor: share storage routing with restore reservations
OutstanderWang Sep 28, 2026
eb6912a
fix: recover missing nested dump schemas on storage owners
OutstanderWang Sep 28, 2026
52efe23
fix: maintain field schemas across wrapped tensor chunks
OutstanderWang Sep 28, 2026
ccbf583
[test] Shorten the live socket pool to trigger the restore receive ti…
OutstanderWang Sep 29, 2026
41a562e
[fix] Drop methods duplicated when rebasing onto the handler-dispatch…
OutstanderWang Sep 29, 2026
b9824fa
fix: preserve per-unit restore report failures
OutstanderWang Sep 29, 2026
888da5e
fix: explain stalled restores and settle failed claims
OutstanderWang Sep 29, 2026
2e66a78
docs: clarify selective dump format compatibility
OutstanderWang Sep 29, 2026
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
11 changes: 10 additions & 1 deletion docs/checkpoint.md
Original file line number Diff line number Diff line change
Expand Up @@ -167,4 +167,13 @@ client.load_controller_checkpoint(...) # (2) controller restored second

If step (1) partially succeeds and step (2) fails, the system is left in a mixed state: some storage units hold checkpoint data while the controller still reflects its pre-restore state. There is no rollback path.

**Workaround**: If `load_checkpoint` raises, call `tq.init()` again to reset the system to a clean state before retrying.
**Workaround**: If `load_checkpoint` raises, call `tq.init()` again to reset the system to a clean state before retrying.

## Exporting selected keys

Use [selective data dumps](data_dump.md) when restoring selected keys into an
existing system or a different number of storage units. SimpleStorage restores for
dump formats v2 and later read assigned row ranges directly on the current storage
owners and merge values instead of replacing entire unit and controller state.
New dumps use v3, which also preserves field schemas. Version-1 dumps use the
caller-side KV put fallback.
187 changes: 187 additions & 0 deletions docs/data_dump.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,187 @@
# Selective data dump

`dump_data_by_key` persists selected keys, their produced fields and tags.
`load_data_by_key` merges them into a running TransferQueue. It preserves existing
key indexes and unrelated rows, fields and tag entries; new keys receive indexes
from the current controller. It does not restore sampler or consumption state.

```python
import transfer_queue as tq

tq.init()
tq.dump_data_by_key("/shared/dumps/selected", ["sample-1", "sample-2"], "train")
index = tq.read_row_index("/shared/dumps/selected")
tq.load_data_by_key("/shared/dumps/selected")
```

Pause writes and clears for these keys during both operations. A dump is not an
atomic snapshot of concurrent writers. Calls on the same dump path are serialized
by an exclusive lock in a stable sibling `.lock` file. Different dump paths remain
independent, and storage units within a load still read in parallel.

## Distributed I/O

On export, the storage manager groups source indexes by their current storage
owner and concurrently asks those units to write their records. Only units holding
selected rows participate. Tensor storage is compacted during serialization, so a
row view cannot include the rest of its original batch, including inside tags.

On version-3 SimpleStorage restore:

1. The caller reads the row index and shard manifest and validates all file ranges.
2. The controller resolves existing keys and allocates indexes for new keys.
3. The storage manager routes records by the **current** indexes, then sends one
load request to each participating unit concurrently.
4. Each target unit reads only its assigned byte ranges and merges those values
into local storage. Records are processed in batches of at most 128 rows per
shard; the caller never reads or forwards their payloads.
5. Each unit claims permission from the controller before writing, then reports
completion directly. The client commits the saved schemas and tags only after
every unit has completed. The controller reserves the destination partition
until commit or confirmed cancellation.

The number of source units can differ from the number of destination units.
Even a dump with one source shard can restore across several target units because
records are independently addressable. Empty rows are recreated from metadata.

The dump directory must be on a filesystem accessible to every participating
storage unit. Local temporary storage suffices for single-node deployments.

| Operation | State | Payload I/O | Unit count on restore |
| --- | --- | --- | --- |
| Checkpoint | Entire controller and storage state | Each unit reads/writes its whole file | Must match |
| Selective v3 dump | Selected fields and tags, merged by key | Each owner unit reads/writes its records | May differ |

`DUMP_ROWS` and `LOAD_ROWS` are included in storage operation metrics. Unit logs
record loaded rows and bytes; the manager logs total bytes and participating units.
These count application reads, not filesystem read-ahead or physical disk traffic.

## Format and compatibility

New dumps use `format_version: 3`:

```text
dump_info.json
row_index.pt
shards/
shard_info.json
shard_0_<source-unit-id>.pkl
...
```

Each shard is a sequence of independent pickle records containing a source global
index and a field/value mapping. `shard_info.json` records each source index's
`[offset, length]`. Source indexes only locate records; they are never reused as
current indexes without controller resolution. `row_index.pt` remains readable
with `read_row_index` without opening payload shards.

Version 3 also saves the original field schemas and only the selected nested row
shapes. Restore uses that schema regardless of target topology or batch boundaries;
non-tensor fields remain non-tensor even when a batch happens to contain only tensors.
Destination type conflicts are rejected before payload writes.

Legacy imports can leave nested fields with missing row shapes, for example when a
later checkpoint chunk wraps tensor values in `NonTensorStack`. Export asks the owner
units to inspect those values while writing their records. The caller merges only
shape/type metadata before publishing the dump; payloads still stay on the units.
Missing tensor shapes are filled from the actual values and their dtype is checked.
If a missing-shape row contains `None` or an object, the entire selected field is
saved as non-tensor, with a warning, so every unit restores the same field contract.
Fields already declared non-tensor remain non-tensor. This repairs the exported
schema without mutating live controller metadata or inventing shapes for objects.
New puts also carry tensor shape hints for homogeneous tensor values wrapped in
`NonTensorStack`. An existing tensor field uses those hints to keep its shape map
complete; real mixed values make the field non-tensor. A field originally declared
non-tensor stays non-tensor when later batches contain only tensors.

Version-1 dumps remain readable using the prior caller-side KV put path. Version-2
dumps retain direct reads, but lack original schemas and use the older inference
behavior; exact field-type preservation cannot be guaranteed for those files.
Restoring to a backend without direct selective loading uses KV puts. Version-1
and KV fallback restores do not provide distributed file reads. Old builds that only understand
versions 1 or 2 cannot read version-3 dumps. Export of nonempty dumps currently requires
SimpleStorage.

## Failure behavior

Publication writes and syncs `.tmp`, moves the old directory to `.old`, publishes
the new directory, and syncs its parent before deleting the backup. If publication
is interrupted while the main directory is absent, the next dump, load or row-index
read recovers `.old`. Readers perform recovery only while holding the same lock as
publishers, so a healthy rename window is never mistaken for a crashed writer.
The load keeps the lock through all remote reads; a pending load marker continues
to prevent replacement after a timeout or client exit. Do not delete the sibling
lock file: unlinking it can create two independent locks for the same dump.
The shared filesystem must provide cross-node advisory locking (not local-only
locks). A backup-cleanup error does not invalidate a published dump.

Restore is not transactional: payload writes before a failure remain. Every load
has a unique ID. The controller blocks clearing/reusing its destination indexes
and conflicting KV puts while an operation is unresolved. Units must claim that ID
before writing; cancellation rejects requests that have not yet claimed permission.
A receive timeout never releases a writer that has already claimed permission.
Timeouts leave the operation pending, even if every unit later succeeds. They do
not cancel it or require changing the timeout used by ordinary puts and gets.

`RestorePendingError` means remote work is still running or its outcome is unknown.
The dump also retains a sibling `.restore` marker so an interrupted client cannot
silently allow its files to be replaced. After an interruption, call:

```python
committed = tq.recover_data_load("/shared/dumps/selected")
```

Recovery asks units to resend terminal results and commits schemas and tags only
when all units succeeded. It returns `True` for committed loads (or no pending
load), and `False` for a failed or cancelled load after all claimed workers stopped.
Unresolved work raises `RestorePendingError`, which includes `reason`, `unit_states`
and `report_errors`. Reports that fail retain the unit ID and original error message.
The controller remembers terminal outcomes so a lost commit reply can be confirmed
safely by retrying recovery; a redundant report failure cannot reverse that outcome.

- `running` units have claimed permission but have not reported completion. Retry
recovery later; a reporting timeout does not prove that a worker has stopped.
- `pending` units have not claimed permission. A request may still be queued, so
recovery does not automatically cancel it. If the initiating client exited before
sending all requests, use `recover_data_load(dump_dir, cancel=True)` to abandon the
operation; simply repeating recovery cannot dispatch the missing requests.
- `reason="unknown_restore"` means the controller has no record of the ID in the
`.restore` marker. If the initiating client has stopped, use explicit cancellation.
After a whole-system restart, first ensure all old actors have stopped, then cancel
the stale operation. This marker blocks its dump path, not the new controller's
unrelated partitions or checkpoints.

A unit whose claim timed out caches a failure confirming that it did not write.
Recovery accepts this failure even when the claim never reached the controller,
cancels unclaimed work and waits for any other claimed workers to finish.

To abandon the load explicitly, use `recover_data_load(dump_dir, cancel=True)`.
This denies unclaimed work and retains the reservation until claimed workers stop.
Cancellation cannot undo a load that has already committed. A failed unit also
cancels remaining unclaimed work; partial payload writes are never rolled back.
After cancellation settles, retry the load or clear its keys. A lost unit requires
stopping the old TQ actors and restarting the whole TQ system; restarting only the
controller while old storage actors run is unsupported.
After restarting, cancel the old marker before loading the dump again:

```python
# Run only after all old TQ actors have stopped and the new system is initialized.
tq.recover_data_load(dump_dir, cancel=True)
tq.load_data_by_key(dump_dir)
```

Writers that already hold low-level metadata must remain paused throughout recovery.
The controller reservation covers only the destination partition. Each storage unit
still serves requests on one worker thread: other partitions using that unit can
wait behind a load. The 128-row batches bound memory, not request latency.

## Tests

Run the selective E2E suite with its default pytest-managed temporary directory:

```bash
python -m pytest -q tests/e2e/test_data_dump_e2e.py tests/e2e/test_data_dump_cross_topology_e2e.py
```

For a multi-node Ray cluster, set `TQ_DUMP_TEST_ROOT` to an existing shared directory.
Tests create and remove only their own child directories beneath that root.
32 changes: 32 additions & 0 deletions tests/e2e/conftest.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,32 @@
# Copyright 2025 Huawei Technologies Co., Ltd. All Rights Reserved.
# Copyright 2025 The TransferQueue Team
#
# Licensed under the Apache License, Version 2.0 (the "License");
# you may not use this file except in compliance with the License.
# You may obtain a copy of the License at
#
# http://www.apache.org/licenses/LICENSE-2.0
#
# Unless required by applicable law or agreed to in writing, software
# distributed under the License is distributed on an "AS IS" BASIS,
# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
# See the License for the specific language governing permissions and
# limitations under the License.

import os
import shutil
import tempfile
from pathlib import Path

import pytest


@pytest.fixture(scope="module")
def dump_test_root(tmp_path_factory):
shared_root = os.environ.get("TQ_DUMP_TEST_ROOT")
if shared_root:
root = Path(tempfile.mkdtemp(prefix="tq-dump-", dir=shared_root))
else:
root = tmp_path_factory.mktemp("tq-dump")
yield root
shutil.rmtree(root)
162 changes: 162 additions & 0 deletions tests/e2e/test_data_dump_cross_topology_e2e.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,162 @@
# Copyright 2025 Huawei Technologies Co., Ltd. All Rights Reserved.
# Copyright 2025 The TransferQueue Team
#
# Licensed under the Apache License, Version 2.0 (the "License");
# you may not use this file except in compliance with the License.
# You may obtain a copy of the License at
#
# http://www.apache.org/licenses/LICENSE-2.0
#
# Unless required by applicable law or agreed to in writing, software
# distributed under the License is distributed on an "AS IS" BASIS,
# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
# See the License for the specific language governing permissions and
# limitations under the License.

"""A dump taken with N storage units must restore into a system with M storage units.

This is the property that separates a data dump from a checkpoint. ``load_checkpoint``
sends each storage unit file back to the unit at the same position, so it requires the
same unit count; ``load_data_by_key`` writes rows back by key and lets TransferQueue
route them for the current topology.

Each test restarts TransferQueue with a different unit count, so this lives apart from
``test_data_dump_e2e.py``, whose fixtures hold one system for the whole module.

Run with:
pytest tests/e2e/test_data_dump_cross_topology_e2e.py -v
"""

import os

import pytest
import ray
import torch
from omegaconf import OmegaConf
from tensordict import NonTensorStack, TensorDict

import transfer_queue as tq

os.environ["RAY_DEDUP_LOGS"] = "0"


def _tq_config(num_storage_units: int) -> OmegaConf:
return OmegaConf.create(
{
"controller": {"polling_mode": True},
"backend": {
"storage_backend": "SimpleStorage",
"SimpleStorage": {
"total_storage_size": 400,
"num_data_storage_units": num_storage_units,
},
},
}
)


@pytest.fixture(scope="module")
def ray_init():
if not ray.is_initialized():
ray.init(namespace="TestDataDumpCrossTopology")
yield
if ray.is_initialized():
ray.shutdown()


@pytest.fixture
def dump_dir(dump_test_root, request):
return dump_test_root / request.node.name / "dump"


def _row_input_ids(row: int) -> torch.Tensor:
return torch.tensor([row * 10, row * 10 + 1, row * 10 + 2])


def _put_rows(partition_id: str, keys: list[str]) -> None:
tq.kv_batch_put(
keys=keys,
partition_id=partition_id,
fields=TensorDict(
{"input_ids": torch.stack([_row_input_ids(row) for row in range(len(keys))])},
batch_size=len(keys),
),
tags=[{"idx": row} for row in range(len(keys))],
)


def _assert_rows_equal(actual: torch.Tensor, expected_rows: list[torch.Tensor]) -> None:
actual_rows = list(actual.unbind()) if actual.is_nested else list(actual)
assert len(actual_rows) == len(expected_rows)
for actual_row, expected_row in zip(actual_rows, expected_rows, strict=True):
assert torch.equal(actual_row, expected_row)


@pytest.mark.parametrize(
("dump_units", "load_units"),
[(4, 2), (2, 4), (3, 3), (1, 4)],
)
def test_dump_restores_across_storage_unit_counts(ray_init, dump_dir, dump_units, load_units):
# Define test data
partition_id = "cross"
keys = [f"c{i}" for i in range(8)]

# Dump with one topology
tq.init(_tq_config(dump_units))
try:
_put_rows(partition_id, keys)
report = tq.dump_data_by_key(dump_dir, keys, partition_id)
assert report["rows_with_data"] == len(keys)
assert report["shards"] <= dump_units
finally:
tq.close()

# Restore into a different topology
tq.init(_tq_config(load_units))
try:
_put_rows("bystander", ["unrelated"])
_put_rows(partition_id, [keys[3]])
old_index = tq.get_client().kv_retrieve_meta([keys[3]], partition_id).global_indexes[0]
tq.load_data_by_key(dump_dir)
assert tq.get_client().kv_retrieve_meta([keys[3]], partition_id).global_indexes == [old_index]
assert tq.kv_batch_get(["unrelated"], "bystander", ["input_ids"]).batch_size[0] == 1

# Check restored state: every row readable, payload and tag intact
retrieved = tq.kv_batch_get(keys=keys, partition_id=partition_id, select_fields=["input_ids"])
_assert_rows_equal(retrieved["input_ids"], [_row_input_ids(row) for row in range(len(keys))])

controller = ray.get_actor("TransferQueueController", namespace="transfer_queue")
snapshot = ray.get(controller.get_partition_snapshot.remote(partition_id))
for row, key in enumerate(keys):
assert snapshot.custom_meta[snapshot.keys_mapping[key]]["idx"] == row
finally:
tq.close()


@pytest.mark.parametrize("row_count", [2, 127, 128, 129, 130])
def test_preserves_nontensor_schema_across_topology(ray_init, dump_dir, row_count):
keys = [f"k{i}" for i in range(row_count)]
values = [torch.tensor([i], dtype=torch.int64) for i in range(row_count - 1)] + [torch.tensor([1.5])]
if row_count > 2:
values[-2] = None
tq.init(_tq_config(1))
try:
tq.kv_batch_put(keys, "mixed", TensorDict({"x": NonTensorStack(*values)}, batch_size=row_count))
tq.dump_data_by_key(dump_dir, keys, "mixed")
finally:
tq.close()
tq.init(_tq_config(2))
try:
tq.load_data_by_key(dump_dir)
schema = tq.get_client().kv_retrieve_meta(keys, "mixed").field_schema["x"]
assert schema["is_non_tensor"]
assert schema["dtype"] is None
assert not schema["is_nested"]
for key, expected in zip(keys, values, strict=True):
value = tq.kv_batch_get([key], "mixed", ["x"])["x"][0]
if expected is None:
assert value is None
else:
torch.testing.assert_close(value, expected)
finally:
tq.close()
Loading
Loading