# How Zephyr Dataset Normalization Handles Pipeline Boundaries

> Discover how Zephyr dataset normalization handles pipeline boundaries by normalizing paths for deterministic, platform-independent matching between readers and writers.

- Repository: [The Marin Project/marin](https://github.com/marin-community/marin)
- Tags: internals
- Published: 2026-08-29

---

**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:

```python
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`](https://github.com/marin-community/marin/blob/main/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:

```python

# 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`:

```python

# 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:

```python

# lib/zephyr/src/zephyr/dataset.py (lines 17-28)

functools.partial(format_shard_path, output_pattern)

```

## Practical Code Examples

### Normalizing Input and Output Paths

```python
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

```python
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`](https://github.com/marin-community/marin/blob/main/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`](https://github.com/marin-community/marin/blob/main/lib/rigging/src/rigging/filesystem/storage_path.py) | `StoragePath` definition and `normalize` static method |

## Summary

- **Every boundary** — input globs and output shards — triggers `StoragePath.normalize` according 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 `rigging` filesystem 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.