Skip to content

Improve async LeRobot persistence throughput and memory bounds - #730

Open
yuecideng wants to merge 7 commits into
mainfrom
codex/offline-save-optimization
Open

yuecideng wants to merge 7 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: 4/5

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

The PR does not appear safe to merge until task-module normalization preserves the registered class’s import identity.

Fix All in CodexFindings

  1. P1 Registered class points elsewhere ▶
Fix with agent prompt
### Issue 1
embodichain/lab/gym/utils/registration.py:481-484
If a task module is imported under its nested name before discovery, discovery can also load a distinct module under the flat name. `setdefault` keeps that flat module, but this code still changes the registered class’s `__module__` to point to it. The resulting Gym entry point identifies a different class, and serialization of the registered class can fail because the class at that path is not the same object. Only change the class path when the flat name resolves to that class.

---

For each issue above, determine whether it is valid and should be fixed. If so, fix it directly.

Summary

The PR adds bounded async LeRobot persistence and configurable PNG compression, documents their tradeoffs, and adds a three-camera benchmark task. The latest changes normalize task-module names after discovery, but can assign a registered class the name of a different module object.

Diagram
%%{init: {'theme': 'neutral'}}%%
flowchart LR
  A[Nested task import registers class] --> B[Discovery loads flat task module]
  B --> C[Duplicate registration is skipped]
  C --> D[Normalization retains existing flat module]
  D --> E[Registered class is assigned flat module name]
  E --> F[Entry point resolves a different class]
Loading

Reviews (7) · Last reviewed commit: "fix: normalize preloaded nested task mod..."

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 on lines +481 to +484
sys.modules.setdefault(flat_name, module)
module.__name__ = flat_name
module.__package__ = flat_name.rpartition(".")[0]
spec.cls.__module__ = flat_name

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P1 Registered class points elsewhere If a task module is imported under its nested name before discovery, discovery can also load a distinct module under the flat name. setdefault keeps that flat module, but this code still changes the registered class’s __module__ to point to it. The resulting Gym entry point identifies a different class, and serialization of the registered class can fail because the class at that path is not the same object. Only change the class path when the flat name resolves to that class.

Prompt To Fix With AI
This is a comment left during a code review.
Path: embodichain/lab/gym/utils/registration.py
Line: 481-484

Comment:
**Registered class points elsewhere** If a task module is imported under its nested name before discovery, discovery can also load a distinct module under the flat name. `setdefault` keeps that flat module, but this code still changes the registered class’s `__module__` to point to it. The resulting Gym entry point identifies a different class, and serialization of the registered class can fail because the class at that path is not the same object. Only change the class path when the flat name resolves to that class.

---

For each issue above, determine whether it is valid and should be fixed. If so, fix it directly.

Fix in Codex Fix in Claude Code

This branch has not been deployed

No deployments
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