Skip to content

Improve async LeRobot persistence throughput and memory bounds - #730

Merged
yuecideng merged 9 commits into
mainfrom
codex/offline-save-optimization
Sep 30, 2026
Merged

yuecideng merged 9 commits into
mainfrom
codex/offline-save-optimization

Conversation

@yuecideng

@yuecideng yuecideng commented Sep 30, 2026 •

Copy link
Copy Markdown
Contributor

Description

This PR improves offline LeRobot persistence for long parallel-environment collections.

It adds an optional bounded AsyncLeRobotRecorder payload queue with queue depth and producer backpressure metrics, and an optional lossless PNG image_compress_level from 0 to 9. Existing defaults remain unchanged: the queue is unbounded when async_queue_maxsize=0, and LeRobot image compression remains level 6 unless explicitly configured.

The data-pipeline context documents the new queue behavior and memory tradeoff.

Validation

  • python -m pytest -q tests/gym/envs/managers/test_async_dataset_functors.py tests/gym/envs/managers/test_dataset_functors.py — 87 passed.
  • python docs/scripts/check_api_docs.py — 2387/2387 exports documented.
  • python -m black . — 1171 files unchanged.
  • Real LeRobot smoke: 1 episode / 10 frames with async_queue_maxsize=1 and PNG level 1; queue drained to zero.

100-trajectory comparison

Same task and hardware in separate worktrees: StayStillSave-v1, 4 environments, 100 trajectories × 100 frames, 320×240 RGB, 8 image-writer threads, Python 3.11, DexSim 0.5.0, LeRobot 0.4.4.

Configuration Total time Episodes Frames Dataset size Queue peak
origin/main, default image compression 99.06 s 100 10,000 665 MB —
PR, PNG level 1, queue max 8 69.78 s 100 10,000 776 MB 8
PR, PNG level 1, unbounded queue 68.81 s 100 10,000 776 MB 74

The bounded configuration reduced end-to-end time by 29.28 s (29.6%) versus main while keeping the queue peak at 8. The unbounded configuration was only 0.96 s faster but reached queue depth 74, so the bounded setting is the safer production choice.

Validation read every Parquet row, checked 100 episode sidecar entries of length 100, decoded all 10,000 RGB images, and verified 10,000 rows of finite action/state data. Numeric hashes matched across runs; rendered image hashes differed because each isolated simulator run produced different camera pixels.

Type of change

  • Bug fix
  • Enhancement
  • New feature
  • Breaking change
  • Documentation update

Dependencies

No new dependencies.

The three-camera demo task is included as a benchmark configuration and registered Gym task.

Checklist

  • I have run the black . command to format the code base.
  • I reviewed and updated affected data-pipeline context.
  • Public API docs are aligned; no public exports changed.
  • I have added tests that prove the queue and compression behavior.
  • Dependencies have not changed.

Standard three-camera demo

The PR now includes StayStillSave3Cam-v1: three 640×480 RGB cameras, one 300-step stationary segment, and a catalog/config entry with the candidate optimized defaults (AsyncLeRobotRecorder, 8 image-writer threads, PNG level 1, queue max 8).

Single-episode measurements (300 frames and 900 decoded RGB images):

Configuration Total time Dataset size
level 6 + 4 threads 36.31 s 222 MB
level 3 + 8 threads 16.19 s 247 MB
level 1 + 8 threads 13.70 s 266 MB

All three runs produced one 300-frame episode and passed image shape, Parquet row, and sidecar-length validation.

16-environment / 100-trajectory validation

Using StayStillSave3Cam-v1 with 16 parallel environments, 100 trajectories, 300 steps per trajectory, three 640×480 cameras, and the candidate defaults (level=1, 8 writer threads, queue max 8):

  • End-to-end: 677.92 s (11.30 min)
  • Startup / collection / finalize: 6.43 / 612.53 / 58.95 s
  • Dataset: 100 episodes, 30,000 frames, approximately 28.47 GB
  • Queue peak 8, 73 backpressure events, 448.32 s cumulative producer wait
  • Metadata and Parquet validation: 30,000 rows, 100 sidecar episodes of length 300 with completed=true, truncated=false, and accepted segments; 90,000 image entries were non-empty and 900 sampled RGB images decoded with shape 480×640×3 across all three cameras

@yuecideng yuecideng added enhancement New feature or request dataset data Related to data_pipeline module gym robot learning env and its related features labels Sep 30, 2026
@greptile-apps

greptile-apps Bot commented Sep 30, 2026 •

Copy link
Copy Markdown

RetriggerConfidence Score: 5/5

[Medium risk] Adds queue bounds and compression tuning to async data recording.

The PR appears safe to merge based on the changes reviewed.

Summary

The PR adds optional queue backpressure and PNG compression controls to LeRobot persistence, plus a three-camera benchmark task. The latest changes narrow task-package discovery to namespace shadowing and remove nested-module normalization.

  • Existing queue and compression defaults remain unchanged.
  • No accepted new findings arose from the changes since the previous review.
Diagram
%%{init: {'theme': 'neutral'}}%%
flowchart LR
  A[Episode payload] --> B[Optional bounded FIFO queue]
  B --> C[Single persistence worker]
  C --> D[LeRobot dataset]
  D --> E[Image writer]
Loading

Reviews (9) · Last reviewed commit: "fix: keep task registration modules flat..."

Comment thread embodichain/lab/gym/envs/managers/async_datasets.py Outdated
Comment thread tests/gym/envs/managers/test_dataset_functors.py
@yuecideng

Copy link
Copy Markdown
Contributor Author

Addressed the review findings and the failing CI tests in commit ffcf9dc7:

  • Queue peak now records the configured max size when queue.Full is observed, avoiding an undercount caused by a worker dequeue between put_nowait() and qsize().
  • Added an actual AsyncImageWriter thread test for image_compress_level, including pixel equality after asynchronous PNG writing.
  • Restored the missing from copy import deepcopy import in configured Task Program integration decoding. The targeted Task Program/catalog/layout tests now pass.

Validation:

  • Dataset manager tests: 88 passed.
  • Task Program catalog/vertical-slice/official-layout tests: 59 passed.
  • API docs: 2388/2388 exports documented.
  • Black and diff checks pass.

1 similar comment
@yuecideng

Copy link
Copy Markdown
Contributor Author

Addressed the review findings and the failing CI tests in commit ffcf9dc7:

  • Queue peak now records the configured max size when queue.Full is observed, avoiding an undercount caused by a worker dequeue between put_nowait() and qsize().
  • Added an actual AsyncImageWriter thread test for image_compress_level, including pixel equality after asynchronous PNG writing.
  • Restored the missing from copy import deepcopy import in configured Task Program integration decoding. The targeted Task Program/catalog/layout tests now pass.

Validation:

  • Dataset manager tests: 88 passed.
  • Task Program catalog/vertical-slice/official-layout tests: 59 passed.
  • API docs: 2388/2388 exports documented.
  • Black and diff checks pass.

@yuecideng

Copy link
Copy Markdown
Contributor Author

Follow-up fixes are now pushed to codex/offline-save-optimization:

  • async_datasets.py: queue_peak records the configured bound when a full queue is observed, so a concurrent dequeue cannot undercount the known peak.
  • test_dataset_functors.py: added an end-to-end AsyncImageWriter thread test for configured PNG compression, with decoded pixel equality.
  • configured.py: restored the missing deepcopy import that caused CI Task Program catalog collection failures.
  • registration.py: hardened legacy outer-package shadow handling so discovered task modules keep their flat public names in the CI editable-install layout.

Local validation now passes:

  • manager tests: 88 passed
  • Task Program/catalog/layout/package-data tests: 99 passed
  • API docs: 2388/2388 exports documented
  • Black and diff checks pass

@yuecideng

Copy link
Copy Markdown
Contributor Author

The CI still exposed the editable-install variant where task modules were imported before entry-point discovery as embodichain_tasks.embodichain_tasks.*. The latest commit d87ca93a normalizes those already-registered task modules to the flat public names after discovery and aliases the loaded modules accordingly.

Local validation with the preloaded nested package scenario now passes the official task-layout test. The combined targeted suite passes 99 tests:

  • tests/test_task_program_package_data.py
  • tests/gym/envs/test_official_task_layout.py
  • tests/gym/envs/task_program/test_catalog.py
  • tests/gym/envs/task_program/test_task_vertical_slices.py
  • tests/gym/envs/managers/test_async_dataset_functors.py
  • tests/gym/envs/managers/test_dataset_functors.py

Comment thread embodichain/lab/gym/utils/registration.py Outdated
@yuecideng

Copy link
Copy Markdown
Contributor Author

CI’s previous test job stayed in progress after the nested-module normalization. I narrowed that fix: it now aliases the already-loaded module and updates registered class __module__, without mutating the module object's __name__ or __package__ metadata.

The local combined targeted suite passes 139 tests, including the preloaded nested-package official-layout scenario. The previous run was canceled to avoid leaving a stale long-running job; the push triggers a fresh CI run.

@yuecideng

Copy link
Copy Markdown
Contributor Author

Root cause of the remaining CI failure was the new demo test importing the repository container path (embodichain_tasks.embodichain_tasks...). That caused task registration during collection to retain nested module names in the full xdist suite. Fixed in c3dfc33b by:

  • importing the demo through the public flat package path;
  • adding the new StayStillSave3Cam-v1 module to the official registration expectations;
  • reverting the overly broad nested-module normalization from the previous attempt, since the correct test import removes the shadowing trigger.

Validation:

  • 962 gym/package tests passed (1 skipped) under xdist;
  • 98 targeted tests passed;
  • API docs 2388/2388 aligned;
  • Black and diff checks pass.

@yuecideng
yuecideng merged commit c38d9f3 into main Sep 30, 2026
9 checks passed
@yuecideng
yuecideng deleted the codex/offline-save-optimization branch September 30, 2026 12:18
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

data Related to data_pipeline module dataset enhancement New feature or request gym robot learning env and its related features

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant