How Zephyr Dataset Normalization Handles Pipeline Boundaries
Zephyr normalizes dataset paths at pipeline boundaries by calling StoragePath.normalize on every input glob pattern and output shard path, collapsing duplicate separators and stripping trailing slashes to ensure deterministic, platform-independent path matching between writers and readers.
The Zephyr dataset API in marin-community/marin treats every filesystem boundary as a critical normalization point. When pipelines ingest input files or produce output files, path strings are canonicalized before any storage operation occurs. This prevents subtle bugs such as double-slash key mismatches that can cause readers to fail to locate files written by upstream stages.
The Normalization Mechanism
Path normalization in Zephyr centers on a single utility: StoragePath.normalize from the rigging filesystem layer. This static method parses a raw path string into a StoragePath object, then re-serializes it to guarantee a canonical representation.
The normalization guarantees three invariants:
- Single separator — internal
//sequences collapse to/(/a//b→/a/b) - No trailing slash — terminal
/characters are stripped (gs://bucket/dir/→gs://bucket/dir) - Scheme-aware handling — local paths, object-store URLs (
gs://,s3://), and empty-authority schemes (mirror://) all normalize consistently
The implementation is intentionally minimal:
class StoragePath:
@staticmethod
def normalize(value: str) -> str:
"""``value`` in canonical single‑separator form (``str(StoragePath(value))``)."""
return str(StoragePath(value))
Source: lib/rigging/src/rigging/filesystem/storage_path.py (lines 26-30)
Where Normalization Applies in Pipelines
Zephyr applies StoragePath.normalize at three critical pipeline stages:
Glob Expansion for Input Files
When resolving user-provided glob patterns into FileEntry lists, resolve_glob normalizes before any filesystem traversal:
# lib/zephyr/src/zephyr/dataset.py (lines 60-63)
pattern = StoragePath.normalize(source.pattern)
This ensures that patterns like gs://bucket//data/**/*.parquet match actual object keys regardless of how many slashes the user included.
Shard-Aware Output Path Formatting
Per-shard output filenames are normalized in format_shard_path:
# lib/zephyr/src/zephyr/dataset.py (lines 11-15)
return StoragePath.normalize(formatted)
This guarantees that formatted paths containing template variables (shard-{shard:05d}) resolve to canonical strings before storage operations.
Callable Output Pattern Wrapping
User-supplied callable patterns are automatically wrapped via _normalize_output_pattern, which applies functools.partial to route through format_shard_path and its normalization:
# lib/zephyr/src/zephyr/dataset.py (lines 17-28)
functools.partial(format_shard_path, output_pattern)
Practical Code Examples
Normalizing Input and Output Paths
from zephyr.dataset import Dataset
# Input glob with duplicate slashes gets normalized
ds = (Dataset
.from_files("gs://my-bucket//data/**/file-*.parquet")
.load_parquet()
.write_parquet("gs://my-bucket//out/shard-{shard:05d}.parquet"))
# Internally:
# resolve_glob → StoragePath.normalize("gs://my-bucket//data/**/file-*.parquet")
# format_shard_path → StoragePath.normalize("gs://my-bucket//out/shard-00001.parquet")
Callable Patterns with Automatic Normalization
def my_pattern(shard_idx, total):
return f"gs://my-bucket/out//shard-{shard_idx:05d}.parquet"
# _normalize_output_pattern wraps the callable, ensuring
# format_shard_path → StoragePath.normalize runs on every invocation
ds = ds.write_parquet(my_pattern)
Even though my_pattern returns a double-slash path, the wrapper normalizes before the path reaches storage.
Key Source Files
| File | Role |
|---|---|
lib/zephyr/src/zephyr/dataset.py |
Core dataset API implementing resolve_glob, format_shard_path, and _normalize_output_pattern |
lib/rigging/src/rigging/filesystem/storage_path.py |
StoragePath definition and normalize static method |
Summary
- Every boundary — input globs and output shards — triggers
StoragePath.normalizeaccording to the marin-community/marin source code - Three pipeline stages apply normalization: glob resolution, shard formatting, and callable pattern wrapping
- Two source files implement the system: the Zephyr dataset layer and the underlying
riggingfilesystem abstraction - Guaranteed invariants: single separators, no trailing slashes, scheme-aware behavior across local and object storage
Frequently Asked Questions
What problems does path normalization solve?
Unequal string representations of the same logical path cause lookup failures. A writer storing to gs://bucket//out/file.parquet and a reader requesting gs://bucket/out/file.parquet see different keys without normalization, breaking pipeline continuity.
Does normalization affect performance?
No. StoragePath.normalize performs a single construction and string conversion — negligible overhead compared to actual storage I/O. The operation is idempotent; normalizing an already-normalized path returns the same string.
Which storage backends benefit from this normalization?
All supported backends: local filesystems, Google Cloud Storage (gs://), S3 (s3://), and custom schemes like mirror://. The StoragePath parser handles each scheme's authority and path components correctly before normalization.
Can users disable or override normalization?
No public API exposes this. Normalization is mandatory at pipeline boundaries to maintain the invariant that identical logical paths have identical string representations. Users who need custom path formatting can supply callable patterns, which are still normalized after formatting.
Have a question about this repo?
These articles cover the highlights, but your codebase questions are specific. Give your agent direct access to the source. Share this with your agent to get started:
curl -s "https://instagit.com/install.md" Maintain an open-source project? Get it listed too →