Dalaran chunk processing
The pipeline layer between raw data and an DLR: readers produce Chunks,
streams transform them, terminal calls execute. This skill is the generic
mechanics only. Decide the data model first (dalaran-data-model), then pick the
importer skill for each source:
| Source |
Reader |
Skill |
| MCAP file (ROS2, protobuf, Foxglove) |
McapReader(path).stream() |
dalaran-mcap |
| URDF robot model (+ joint states → FK) |
UrdfTree.from_file_path(...).stream() |
dalaran-urdf |
| Parquet table (trajectories, sensor logs) |
ParquetReader(path).stream() |
dalaran-parquet |
| LeRobot dataset directory |
built-in importer, then RrdReader |
dalaran-lerobot |
| Existing DLR |
RrdReader(path) |
here, below |
| Sidecar files (JSON calib, metadata) |
Chunk.from_columns + from_iter |
here, below |
The API is dalaran.experimental; when
behavior matters, check the installed surface:
python -c "from dalaran.experimental import LazyChunkStream; help(LazyChunkStream)".
Decision rule: where does each component come from?
Default: a reader produces the chunks; lenses shape them. Walk this before
writing any conversion code — most "build it by hand" instincts are wrong here:
- Source a reader supports? Use the reader's
.stream(); never hand-parse
and re-log. MCAP→McapReader, URDF→UrdfTree, parquet→ParquetReader,
DLR→RrdReader, LeRobot dir→log_file_from_path.
- A decoder already emits the archetype? Foxglove gives
Transform3D,
Pinhole, VideoStream (real sample bytes) ready-made — pass it through,
do not re-derive. Only custom-protobuf topics arrive as <Name>:message and
need a lens (see dalaran-mcap).
- Fix an existing component in place (swapped resolution, recolor, unit
convert)?
MutateLens, output_mode="forward_unmatched".
- Derive a new component/entity (FK→
/tf, scalars from a message)?
DeriveLens. To scatter one row into N (a joint batch → per-joint /tf),
use the two-lens pair: derive the batch with output_mode="forward_all"
(keeps the originals, e.g. the joint states), then a second
DeriveLens with scatter=True and output_mode="drop_unmatched" (emits
only the scattered rows). See the robot_data_preprocessing example.
- Genuine sidecar no reader or lens can produce (JSON calibration offsets,
hand-measured extrinsics, external metadata)?
Chunk.from_columns + from_iter.
- Finish with
LazyChunkStream.merge(...) →
.collect(optimize=OptimizationProfile.OBJECT_STORE) →
write_dlr(application_id, recording_id).
Why this order: the pipeline stays lazy, columnar, multithreaded, and
OBJECT_STORE-optimizable. A hand-built row loop or an out-of-lens pa.array
throws all of that away — that is the path we are deliberately avoiding.
Anti-patterns (use a reader + lens instead)
If you are writing the left, stop and use the right:
for-loop building rows/components → a lens with a Selector(...).pipe(...)
PyArrow-compute callback.
rr.init + rr.log per message for conversion → that is live logging;
for ingestion, read with a reader and write_dlr.
chunk.to_record_batch() + pc.filter then rebuilding via
Chunk.from_columns (row-thinning by hand) → stream.drop(content=...),
.split(...), or a MutateLens returning a filtered pa.array.
pa.array / pa.RecordBatch / np.frombuffer assembled OUTSIDE a lens →
move the transform inside a MutateLens/DeriveLens selector callback.
rr.send_columns hand-assembled from a custom parser → use the matching
reader; it produces chunks directly.
- Parsing MCAP/URDF with a non-Dalaran library then re-logging →
McapReader
/ UrdfTree.
Chunk.from_columns for data a reader already decodes (Pinhole
intrinsics, VideoStream, Transform3D from a transforms topic) → keep it in
the reader stream; fix with a MutateLens if needed.
A wall of pyarrow.compute "missing-attribute" type errors (pc.filter,
pc.list_element) usually means pc.* calls sit in module-level helpers instead
of inside Selector.pipe lens callbacks. Refactor into a lens before suppressing
the checker — the errors are a smell that the hand-building should not exist.
Porting an existing converter? Hand-built converters predate decoder
improvements and are not ground truth. Re-verify the decoder output (step 2) and
check every Chunk.from_columns / for-loop against this list before copying.
Core model
LazyChunkStream is a lazy pipeline DAG, not a collection. Building
filters, lenses, maps, splits, and merges reads no source data.
- Execution starts at terminal calls:
write_dlr(...), collect(),
to_chunks(), or iterating the stream.
- Execution is streaming, multithreaded, and mostly GIL-free. Prefer
stream/lens operations and PyArrow compute over Python row loops.
- Move semantics: builder calls (
filter, drop, lenses, map,
flat_map) consume the input stream; reusing a consumed stream raises.
Reassign after each step. Terminal calls do not consume, but each terminal
call re-executes the whole pipeline; collect() once if that is too costly.
ChunkStore is materialized in memory (stream.collect(),
ChunkStore.from_chunks). LazyStore is manifest-indexed, loads chunks on
demand (RrdReader(path).store(), catalog segment stores). Both have
schema(), summary(), stream(), and write_dlr(...).
Stream composition
from dalaran.experimental import Chunk, LazyChunkStream, OptimizationProfile
stream.filter(content=, has_timeline=, is_static=, components=) keeps the
matching portion of each chunk; stream.drop(...) is its complement, same
keyword filters. content takes an entity-path glob or a list of them.
stream.map(fn) applies Chunk -> Chunk; stream.flat_map(fn) applies
Chunk -> Iterable[Chunk]. Escape hatches for chunk-level Python logic;
prefer lenses for columnar work.
stream.split(content=, ...) returns (matching, non_matching); both
branches share the same upstream.
LazyChunkStream.merge(*streams) fans in any number of sources.
LazyChunkStream.from_iter(chunks) wraps hand-built chunks.
stream = source_stream() # any importer skill
stream = stream.drop(content="/video_raw/**")
stream = stream.lenses(fix_lens, content="/cam/**", output_mode="forward_unmatched")
merged = LazyChunkStream.merge(stream, sidecar_stream)
merged.write_dlr(out_path, application_id="my_app", recording_id=recording_id)
Hand-built chunks — sidecar only
Use Chunk.from_columns ONLY for data no reader or lens can emit — JSON/CSV
calibration, frame offsets, external metadata. If a reader
(McapReader/UrdfTree/ParquetReader) decodes the topic or a lens can derive
it, that is the idiomatic path; do not hand-assemble it here. In the
robot_data_preprocessing example the only hand-built chunk is the JSON
offsets sidecar; the camera fix, FK→/tf, meshes, and recolor are all
readers + lenses.
Chunk.from_columns(entity_path, indexes, columns) mirrors
rr.send_columns(...) and accepts the same archetype .columns(...) helpers.
Empty indexes means static.
chunk = Chunk.from_columns(
"/tf_static/robot_offsets",
indexes=[], # static
columns=rr.Transform3D.columns(
translation=translations,
quaternion=quaternions_xyzw,
parent_frame=parents,
child_frame=children,
),
)
sidecar_stream = LazyChunkStream.from_iter([chunk])
rr.AnyValues.columns(...) covers non-standard metadata fields. For
inspection, a Chunk exposes entity_path, num_rows, is_static,
timeline_names, to_record_batch(), and format() (human-readable table).
Lenses
Lenses reshape, fix, or derive components without iterating rows. Apply with
stream.lenses(lenses, output_mode=..., content=...).
MutateLens(component, selector, keep_row_ids=False) modifies an existing
component in place.
DeriveLens(component, output_entity=None, scatter=False) creates new
columns, optionally at another entity. Chain .to_component(descriptor, selector) per output; .to_timeline(name, "sequence" | "duration_ns" | "timestamp_ns", selector) extracts a time column from the data itself.
scatter=True explodes one input row into N output rows (one per list
element).
- Scope with
content= whenever the same component name exists under multiple
entities.
Output modes, and the default is drop_unmatched:
drop_unmatched (default): only lens outputs survive. Right for derive-only
intermediate streams; silently discards everything else if applied broadly.
forward_unmatched: lens outputs plus the original components no lens
consumed. Right for targeted fixes that preserve the rest of the stream.
forward_all: lens outputs plus all originals, including consumed ones. Can
duplicate data.
In-place fix (keep Arrow type and length intact):
stream = stream.lenses(
MutateLens(
"Pinhole:resolution",
Selector(".").pipe(
lambda res: pa.array(
[(h, w) for w, h in res.to_pylist()],
type=res.type,
)
),
),
content=["/external/cam_low", "/external/cam_high"],
output_mode="forward_unmatched",
)
Derive with unit conversion (PyArrow compute, no Python loop):
DeriveLens("schemas.proto.JointState:message", output_entity="/joints_deg/waist").to_component(
rr.Scalars.descriptor_scalars(),
Selector(".joint_positions").pipe(lambda arr: pc.multiply(pc.list_element(arr, 0), 180.0 / math.pi)),
)
Selector grammar
Selector("<query>") navigates nested Arrow data, jq-style:
. current value; .field struct field
[] iterate list elements; [N] index a list
? suppress errors / skip missing optionals; ! assert non-null
| pipe one expression into another
.pipe(fn) chains a Python/PyArrow transform (or another Selector).
.execute(array) runs it eagerly; .execute_per_row(array) guarantees the
output row count matches the input (use inside lens callbacks that must stay
row-aligned).
Writing RRDs
stream.write_dlr(path, application_id=..., recording_id=...) executes and
writes in one streaming pass.
stream.collect(optimize=OptimizationProfile.OBJECT_STORE).write_dlr(...)
materializes, optimizes chunk layout, then writes. Memory scales with the
materialized chunks.
- Profiles:
OBJECT_STORE (large chunks, for storage/query/catalog) and
LIVE (small chunks, low-latency viewer).
- Multiple physical RRDs form one logical recording when they share a
recording_id; use this to separate base data, model/URDF data, and layers.
Always use OptimizationProfile.OBJECT_STORE when the DLR is headed for a
Dalaran catalog or Hub, unless explicitly asked otherwise.
Chunk API vs logging API
- Logging (
rr.log, rr.send_columns, RecordingStream) is for live logging
from user code; chunk processing is for ingestion, conversion, and
postprocessing existing recordings.
- Logging → chunks: write an DLR, read it back with
RrdReader.
RrdReader(path) lists recordings() / blueprints() (each a StoreEntry
with kind, application_id, recording_id); .stream(store=entry) for
sequential passes, .store(store=entry) for indexed access.
- Chunks → logging:
dalaran.experimental.send_chunks(chunks, recording=...)
accepts a Chunk, LazyChunkStream, LazyStore, ChunkStore, or any
iterable of chunks. The source store's application_id/recording_id are
not preserved; the active recording's identity wins.
Common gotchas
- The default lens
output_mode is drop_unmatched; forgetting to set
forward_unmatched on a targeted fix silently drops the rest of the stream.
- Do not reuse a consumed
LazyChunkStream; reassign or split deliberately.
- Scope lenses with
content=; the same component name often exists under
many entities.
- Preserve Arrow array type and length in
MutateLens transforms.
- For catalog layers, the layer
recording_id must equal the segment id.
- This is
dalaran.experimental; pin-check signatures when upgrading.
References
- End-to-end example (MCAP + URDF + JSON sidecar, lenses, merge, optimize):
https://github.com/Flaminis/Dalaran/tree/main/examples/python/robot_data_preprocessing
- Docs:
https://dalaran.dev/docs/concepts/logging-and-ingestion/chunk-processing-api,
https://dalaran.dev/docs/concepts/query-and-transform/lenses
1---2name: dalaran-chunk-processing3description: Core mechanics of the Dalaran Chunk Processing API (dalaran.experimental) — LazyChunkStream pipelines, Chunk, lenses (MutateLens/DeriveLens/Selector), RrdReader, writing optimized RRDs. Read BEFORE writing any ingestion/conversion/preprocessing code (convert an MCAP, build a recording from a dataset, preprocess an .dlr, port an old converter): it mandates reader+lens pipelines and steers away from hand-built chunks — no Chunk.from_columns for data a reader/lens can produce, no per-message rr.log, no manual pa.array assembly. Source-specific knowledge lives in the importer skills (dalaran-mcap, dalaran-urdf, dalaran-parquet, dalaran-lerobot); read dalaran-data-model first to decide what the data should become.4---56# Dalaran chunk processing78The pipeline layer between raw data and an DLR: readers produce `Chunk`s,9streams transform them, terminal calls execute. This skill is the generic10mechanics only. Decide the data model first (`dalaran-data-model`), then pick the11importer skill for each source:1213| Source | Reader | Skill |14| ----------------------------------------- | --------------------------------------- | --------------- |15| MCAP file (ROS2, protobuf, Foxglove) | `McapReader(path).stream()` | `dalaran-mcap` |16| URDF robot model (+ joint states → FK) | `UrdfTree.from_file_path(...).stream()` | `dalaran-urdf` |17| Parquet table (trajectories, sensor logs) | `ParquetReader(path).stream()` | `dalaran-parquet` |18| LeRobot dataset directory | built-in importer, then `RrdReader` | `dalaran-lerobot` |19| Existing DLR | `RrdReader(path)` | here, below |20| Sidecar files (JSON calib, metadata) | `Chunk.from_columns` + `from_iter` | here, below |2122The API is `dalaran.experimental`; when23behavior matters, check the installed surface:24`python -c "from dalaran.experimental import LazyChunkStream; help(LazyChunkStream)"`.2526## Decision rule: where does each component come from?2728Default: **a reader produces the chunks; lenses shape them.** Walk this before29writing any conversion code — most "build it by hand" instincts are wrong here:30311. **Source a reader supports?** Use the reader's `.stream()`; never hand-parse32 and re-log. MCAP→`McapReader`, URDF→`UrdfTree`, parquet→`ParquetReader`,33 DLR→`RrdReader`, LeRobot dir→`log_file_from_path`.342. **A decoder already emits the archetype?** Foxglove gives `Transform3D`,35 `Pinhole`, `VideoStream` (real sample bytes) ready-made — **pass it through**,36 do not re-derive. Only custom-protobuf topics arrive as `<Name>:message` and37 need a lens (see `dalaran-mcap`).383. **Fix an existing component in place** (swapped resolution, recolor, unit39 convert)? `MutateLens`, `output_mode="forward_unmatched"`.404. **Derive a new component/entity** (FK→`/tf`, scalars from a message)?41 `DeriveLens`. To scatter one row into N (a joint batch → per-joint `/tf`),42 use the **two-lens pair**: derive the batch with `output_mode="forward_all"`43 (keeps the originals, e.g. the joint states), then a second44 `DeriveLens` with `scatter=True` and `output_mode="drop_unmatched"` (emits45 only the scattered rows). See the `robot_data_preprocessing` example.465. **Genuine sidecar** no reader or lens can produce (JSON calibration offsets,47 hand-measured extrinsics, external metadata)? `Chunk.from_columns` + `from_iter`.486. Finish with `LazyChunkStream.merge(...)` →49 `.collect(optimize=OptimizationProfile.OBJECT_STORE)` →50 `write_dlr(application_id, recording_id)`.5152Why this order: the pipeline stays lazy, columnar, multithreaded, and53`OBJECT_STORE`-optimizable. A hand-built row loop or an out-of-lens `pa.array`54throws all of that away — that is the path we are deliberately avoiding.5556## Anti-patterns (use a reader + lens instead)5758If you are writing the left, stop and use the right:5960- **`for`-loop building rows/components** → a lens with a `Selector(...).pipe(...)`61 PyArrow-compute callback.62- **`rr.init` + `rr.log` per message for conversion** → that is _live_ logging;63 for ingestion, read with a reader and `write_dlr`.64- **`chunk.to_record_batch()` + `pc.filter` then rebuilding via65 `Chunk.from_columns`** (row-thinning by hand) → `stream.drop(content=...)`,66 `.split(...)`, or a `MutateLens` returning a filtered `pa.array`.67- **`pa.array` / `pa.RecordBatch` / `np.frombuffer` assembled OUTSIDE a lens** →68 move the transform inside a `MutateLens`/`DeriveLens` selector callback.69- **`rr.send_columns` hand-assembled from a custom parser** → use the matching70 reader; it produces chunks directly.71- **Parsing MCAP/URDF with a non-Dalaran library then re-logging** → `McapReader`72 / `UrdfTree`.73- **`Chunk.from_columns` for data a reader already decodes** (`Pinhole`74 intrinsics, `VideoStream`, `Transform3D` from a transforms topic) → keep it in75 the reader stream; fix with a `MutateLens` if needed.7677A wall of `pyarrow.compute` "missing-attribute" type errors (`pc.filter`,78`pc.list_element`) usually means `pc.*` calls sit in module-level helpers instead79of inside `Selector.pipe` lens callbacks. Refactor into a lens before suppressing80the checker — the errors are a smell that the hand-building should not exist.8182**Porting an existing converter?** Hand-built converters predate decoder83improvements and are not ground truth. Re-verify the decoder output (step 2) and84check every `Chunk.from_columns` / for-loop against this list before copying.8586## Core model8788- `LazyChunkStream` is a lazy pipeline DAG, not a collection. Building89 filters, lenses, maps, splits, and merges reads no source data.90- Execution starts at terminal calls: `write_dlr(...)`, `collect()`,91 `to_chunks()`, or iterating the stream.92- Execution is streaming, multithreaded, and mostly GIL-free. Prefer93 stream/lens operations and PyArrow compute over Python row loops.94- **Move semantics**: builder calls (`filter`, `drop`, `lenses`, `map`,95 `flat_map`) consume the input stream; reusing a consumed stream raises.96 Reassign after each step. Terminal calls do not consume, but each terminal97 call re-executes the whole pipeline; `collect()` once if that is too costly.98- `ChunkStore` is materialized in memory (`stream.collect()`,99 `ChunkStore.from_chunks`). `LazyStore` is manifest-indexed, loads chunks on100 demand (`RrdReader(path).store()`, catalog segment stores). Both have101 `schema()`, `summary()`, `stream()`, and `write_dlr(...)`.102103## Stream composition104105```python106from dalaran.experimental import Chunk, LazyChunkStream, OptimizationProfile107```108109- `stream.filter(content=, has_timeline=, is_static=, components=)` keeps the110 matching portion of each chunk; `stream.drop(...)` is its complement, same111 keyword filters. `content` takes an entity-path glob or a list of them.112- `stream.map(fn)` applies `Chunk -> Chunk`; `stream.flat_map(fn)` applies113 `Chunk -> Iterable[Chunk]`. Escape hatches for chunk-level Python logic;114 prefer lenses for columnar work.115- `stream.split(content=, ...)` returns `(matching, non_matching)`; both116 branches share the same upstream.117- `LazyChunkStream.merge(*streams)` fans in any number of sources.118- `LazyChunkStream.from_iter(chunks)` wraps hand-built chunks.119120```python121stream = source_stream() # any importer skill122stream = stream.drop(content="/video_raw/**")123stream = stream.lenses(fix_lens, content="/cam/**", output_mode="forward_unmatched")124merged = LazyChunkStream.merge(stream, sidecar_stream)125merged.write_dlr(out_path, application_id="my_app", recording_id=recording_id)126```127128## Hand-built chunks — sidecar only129130Use `Chunk.from_columns` ONLY for data no reader or lens can emit — JSON/CSV131calibration, frame offsets, external metadata. If a reader132(`McapReader`/`UrdfTree`/`ParquetReader`) decodes the topic or a lens can derive133it, that is the idiomatic path; do not hand-assemble it here. In the134`robot_data_preprocessing` example the _only_ hand-built chunk is the JSON135offsets sidecar; the camera fix, FK→`/tf`, meshes, and recolor are all136readers + lenses.137138`Chunk.from_columns(entity_path, indexes, columns)` mirrors139`rr.send_columns(...)` and accepts the same archetype `.columns(...)` helpers.140Empty `indexes` means static.141142```python143chunk = Chunk.from_columns(144 "/tf_static/robot_offsets",145 indexes=[], # static146 columns=rr.Transform3D.columns(147 translation=translations,148 quaternion=quaternions_xyzw,149 parent_frame=parents,150 child_frame=children,151 ),152)153sidecar_stream = LazyChunkStream.from_iter([chunk])154```155156`rr.AnyValues.columns(...)` covers non-standard metadata fields. For157inspection, a `Chunk` exposes `entity_path`, `num_rows`, `is_static`,158`timeline_names`, `to_record_batch()`, and `format()` (human-readable table).159160## Lenses161162Lenses reshape, fix, or derive components without iterating rows. Apply with163`stream.lenses(lenses, output_mode=..., content=...)`.164165- `MutateLens(component, selector, keep_row_ids=False)` modifies an existing166 component in place.167- `DeriveLens(component, output_entity=None, scatter=False)` creates new168 columns, optionally at another entity. Chain `.to_component(descriptor,169selector)` per output; `.to_timeline(name, "sequence" | "duration_ns" |170"timestamp_ns", selector)` extracts a time column from the data itself.171 `scatter=True` explodes one input row into N output rows (one per list172 element).173- Scope with `content=` whenever the same component name exists under multiple174 entities.175176Output modes, and **the default is `drop_unmatched`**:177178- `drop_unmatched` (default): only lens outputs survive. Right for derive-only179 intermediate streams; silently discards everything else if applied broadly.180- `forward_unmatched`: lens outputs plus the original components no lens181 consumed. Right for targeted fixes that preserve the rest of the stream.182- `forward_all`: lens outputs plus all originals, including consumed ones. Can183 duplicate data.184185In-place fix (keep Arrow type and length intact):186187```python188stream = stream.lenses(189 MutateLens(190 "Pinhole:resolution",191 Selector(".").pipe(192 lambda res: pa.array(193 [(h, w) for w, h in res.to_pylist()],194 type=res.type,195 )196 ),197 ),198 content=["/external/cam_low", "/external/cam_high"],199 output_mode="forward_unmatched",200)201```202203Derive with unit conversion (PyArrow compute, no Python loop):204205```python206DeriveLens("schemas.proto.JointState:message", output_entity="/joints_deg/waist").to_component(207 rr.Scalars.descriptor_scalars(),208 Selector(".joint_positions").pipe(lambda arr: pc.multiply(pc.list_element(arr, 0), 180.0 / math.pi)),209)210```211212## Selector grammar213214`Selector("<query>")` navigates nested Arrow data, jq-style:215216- `.` current value; `.field` struct field217- `[]` iterate list elements; `[N]` index a list218- `?` suppress errors / skip missing optionals; `!` assert non-null219- `|` pipe one expression into another220221`.pipe(fn)` chains a Python/PyArrow transform (or another Selector).222`.execute(array)` runs it eagerly; `.execute_per_row(array)` guarantees the223output row count matches the input (use inside lens callbacks that must stay224row-aligned).225226## Writing RRDs227228- `stream.write_dlr(path, application_id=..., recording_id=...)` executes and229 writes in one streaming pass.230- `stream.collect(optimize=OptimizationProfile.OBJECT_STORE).write_dlr(...)`231 materializes, optimizes chunk layout, then writes. Memory scales with the232 materialized chunks.233- Profiles: `OBJECT_STORE` (large chunks, for storage/query/catalog) and234 `LIVE` (small chunks, low-latency viewer).235- Multiple physical RRDs form one logical recording when they share a236 `recording_id`; use this to separate base data, model/URDF data, and layers.237238**Always use `OptimizationProfile.OBJECT_STORE`** when the DLR is headed for a239Dalaran catalog or Hub, unless explicitly asked otherwise.240241## Chunk API vs logging API242243- Logging (`rr.log`, `rr.send_columns`, `RecordingStream`) is for live logging244 from user code; chunk processing is for ingestion, conversion, and245 postprocessing existing recordings.246- Logging → chunks: write an DLR, read it back with `RrdReader`.247 `RrdReader(path)` lists `recordings()` / `blueprints()` (each a `StoreEntry`248 with `kind`, `application_id`, `recording_id`); `.stream(store=entry)` for249 sequential passes, `.store(store=entry)` for indexed access.250- Chunks → logging: `dalaran.experimental.send_chunks(chunks, recording=...)`251 accepts a `Chunk`, `LazyChunkStream`, `LazyStore`, `ChunkStore`, or any252 iterable of chunks. The source store's `application_id`/`recording_id` are253 **not** preserved; the active recording's identity wins.254255## Common gotchas256257- The default lens `output_mode` is `drop_unmatched`; forgetting to set258 `forward_unmatched` on a targeted fix silently drops the rest of the stream.259- Do not reuse a consumed `LazyChunkStream`; reassign or `split` deliberately.260- Scope lenses with `content=`; the same component name often exists under261 many entities.262- Preserve Arrow array type and length in `MutateLens` transforms.263- For catalog layers, the layer `recording_id` must equal the segment id.264- This is `dalaran.experimental`; pin-check signatures when upgrading.265266## References267268- End-to-end example (MCAP + URDF + JSON sidecar, lenses, merge, optimize):269 `https://github.com/Flaminis/Dalaran/tree/main/examples/python/robot_data_preprocessing`270- Docs: `https://dalaran.dev/docs/concepts/logging-and-ingestion/chunk-processing-api`,271 `https://dalaran.dev/docs/concepts/query-and-transform/lenses`