Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
112 changes: 89 additions & 23 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -16,6 +16,8 @@
- **Audit trails**: Generate complete audit protocols with deterministic DataFrame hashing
- **Multi-backend support**: Polars (default) and DuckDB backends with a pluggable Backend ABC
- **Serializable pipelines**: Save and load transformation plans as JSON
- **Structured plans**: Group long pipelines into named sections and compose them from reusable blocks
- **Documented decisions**: Record *why* a step exists, not just what it did

## Quick Example

Expand All @@ -25,20 +27,19 @@ from transformplan import TransformPlan, Col
# Build readable pipelines with 89 chainable operations
plan = (
TransformPlan()
# Standardize column names
.section("Standardize")
.col_rename(column="PatientID", new_name="patient_id")
.col_rename(column="DOB", new_name="date_of_birth")
.str_strip(column="patient_id")
.col_drop(["tmp_a", "tmp_b"]).because("scratch columns from the export")

# Calculate derived values
.dt_age_years(column="date_of_birth", new_column="age")
.math_clamp(column="age", min_value=0, max_value=120)
.section("Derive")
.dt_age_years(birth_column="date_of_birth", new_column="age")
.col_cast("raw_weight", "Float64")
.map_discretize(column="age", bins=[18, 40, 65], labels=["minor", "young", "adult", "senior"], new_column="age_group")

# Categorize patients age
.map_discretize(column="age", bins=[18, 40, 65], labels=["young", "adult", "senior"], new_column="age_group")

# Filter and clean
.rows_filter(Col("age") >= 18)
.section("Filter")
.rows_filter(Col("age") >= 18).because("cohort is adults only")
.rows_drop_nulls(columns=["patient_id", "age"])
.col_drop(column="date_of_birth")
)
Expand All @@ -64,27 +65,87 @@ protocol.print(show_params=False)
======================================================================
TRANSFORM PROTOCOL
======================================================================
Input: 1000 rows × 5 cols [a4f8b2c1]
Output: 847 rows × 5 cols [e7d3f9a2]
Total time: 0.0247s
Input: 1000 rows x 5 cols [a8bfc98263e8aa4c]
Output: 841 rows x 4 cols [6f62d5e6d677cb81]
Total time: 0.0011s
----------------------------------------------------------------------

SECTIONS

Standardize 1000 → 1000 rows (5 steps, -2 cols)
Derive 1000 → 1000 rows (3 steps, +2 cols)
Filter 1000 → 841 rows (3 steps, -1 cols)
----------------------------------------------------------------------

# Operation Rows Cols Time Hash
----------------------------------------------------------------------
0 input 1000 5 - a4f8b2c1
1 col_rename 1000 5 0.0012s b2e4a7f3
2 col_rename 1000 5 0.0008s c9d1e5b8
3 str_strip 1000 5 0.0013s c9d1e5b8 ○
4 dt_age_years 1000 6 (+1) 0.0041s d4f2c8a1
5 math_clamp 1000 6 0.0015s e1b7d3f9
6 map_discretize 1000 7 (+1) 0.0028s f8a4c2e6
7 rows_filter 858 (-142) 7 0.0037s a2e9f4b7
8 rows_drop_nulls 847 (-11) 7 0.0019s b5c1d8e3
9 col_drop 847 6 (-1) 0.0006s e7d3f9a2
0 input 1000 5 - a8bfc98263e8aa4c
[Standardize]
1 col_rename 1000 5 0.0000s 4a3d3ba5644cf847
2 col_rename 1000 5 0.0000s e535f4ddd1e78927
3 str_strip 1000 5 0.0001s 14cf5b02000d13a4
4 col_drop 1000 4 (-1) 0.0001s ccaddd0bc65974d7
↳ scratch columns from the export
5 col_drop 1000 3 (-1) 0.0000s 08a006cf6351723c
↳ scratch columns from the export
[Derive]
6 dt_age_years 1000 4 (+1) 0.0001s 2e3c2e80b77f7e10
7 col_cast 1000 4 0.0001s e7e31d60ee332aa6
8 map_discretize 1000 5 (+1) 0.0002s f3c8c0998d8ee820
[Filter]
9 rows_filter 851 (-149) 5 0.0002s 97f54c043b457d0c
↳ cohort is adults only
10 rows_drop_nulls 841 (-10) 5 0.0002s 7bb8de08021dfe20
11 col_drop 841 4 (-1) 0.0001s 6f62d5e6d677cb81
======================================================================
○ = no effect (steps 3 did not change data)
```

Sections give a plan with a hundred steps a readable balance instead of one flat
list, and `because()` records the domain knowledge that otherwise lives only in a
code comment — the question an audit actually asks.

### Composing Long Pipelines

Repeated blocks can live in their own function or their own plan, without breaking
the chain:

```python
# A function that takes and returns a plan
def decimal_hour(plan, source, target):
return (
plan.dt_format(source, "%H", target)
.col_cast(target, "Float64")
)

# A reusable block, defined as a plan of its own
CLEAN_IDS = TransformPlan().str_strip("patient_id").str_upper("patient_id")

plan = (
TransformPlan()
.pipe(decimal_hour, "admitted_at", "admit_hour")
.pipe(decimal_hour, "discharged_at", "discharge_hour")
.extend(CLEAN_IDS) # or: plan_a + plan_b
)
```

`pipe()` registers no step of its own, so the protocol is unchanged. `extend()`
deep-copies the block, so the same block can be reused across plans without them
sharing state — and unlike `pipe()`, its steps are serialized with the plan.

Operations that work in place also take a sequence of columns, so repetition
collapses into one call:

```python
plan = (
TransformPlan()
.col_drop(["fiscal_year", "import_batch", "row_checksum"])
.str_slice(["primary_code", "secondary_code"], 0, 3)
.col_cast(["weight", "height"], "Float64")
)
```

Each column still becomes its own protocol step.

### DuckDB Backend

Run the same pipelines on DuckDB for SQL-based execution and native large-file handling:
Expand Down Expand Up @@ -119,6 +180,11 @@ result, protocol = plan.process(rel, backend=DuckDBBackend(con))
| **dt\_** | Datetime operations | `dt_year`, `dt_month`, `dt_parse`, `dt_age_years`, `dt_diff_days` |
| **map\_** | Value mapping & encoding | `map_values`, `map_discretize`, `map_onehot`, `map_ordinal` |

24 in-place operations accept either a single column or a sequence of columns.
Operations that name an output column (`new_column`) take one column at a time.

See the [roadmap](docs/roadmap.md) for planned additions.

## Installation

```bash
Expand Down
91 changes: 91 additions & 0 deletions docs/api/plan.md
Original file line number Diff line number Diff line change
Expand Up @@ -50,12 +50,103 @@ See [Backends](backends.md) for details on each backend.
- validate
- validate_chunked
- dry_run
- pipe
- extend
- section
- because
- to_dict
- from_dict
- to_json
- from_json
- to_python

## Structuring Long Plans

Plans with dozens of steps become hard to read as a flat chain. Four methods
give them structure without changing what they do.

### pipe

Apply a function that takes and returns a plan, so repeated blocks can live in
their own function without breaking the chain.

```python
def decimal_hour(plan, source, target):
return (
plan.dt_format(source, "%H", target)
.col_cast(target, "Float64")
)

plan = (
TransformPlan()
.dt_diff_days("admitted", "discharged", "stay_days")
.pipe(decimal_hour, "admitted_at", "admit_hour")
.pipe(decimal_hour, "discharged_at", "discharge_hour")
)
```

`pipe` registers no step of its own — the function registers ordinary steps, so
the protocol is unchanged.

### extend

Append another plan's steps, so a reusable block can be defined once as its own
plan. Steps are deep-copied, so a block can be reused across plans without them
sharing state.

```python
CLEAN_NAMES = TransformPlan().str_strip("name").str_upper("name")

plan = TransformPlan().col_drop("temp").extend(CLEAN_NAMES)
combined = plan_a + plan_b # same thing, as a new plan
```

Unlike `pipe`, the block's steps become part of the plan and are serialized with it.

### section

Group the following steps under a label. The protocol then reports a balance per
section instead of one flat list.

```python
plan = (
TransformPlan()
.section("Row filters")
.rows_drop_nulls("shipped_at")
.rows_drop_nulls("order_id")
.section("Split product code")
.col_duplicate("code", "code_group")
.str_slice("code_group", 0, 3)
)
```

```
SECTIONS

Row filters 1500 → 49 rows (2 steps)
Split product code 49 → 49 rows (2 steps, +1 cols)
```

Pass `None` to end a section. Plans without sections render exactly as before.

### because

Attach a reason to the step just registered. The protocol records that rows were
removed; the reason records why — the question an audit actually asks.

```python
plan = (
TransformPlan()
.rows_drop(Col("provider") != "MDK02")
.because("cases where the reviewer awards the insurer nothing")
.rows_drop_nulls("drg_code")
.because("PEPP cases carry no DRG")
)
```

When the preceding call covered several columns, the reason is attached to every
step it produced.

## Execution Methods

### process
Expand Down
106 changes: 106 additions & 0 deletions docs/roadmap.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,106 @@
# Roadmap

Planned additions that are deliberately **not** in the current release, with the
reasoning and the groundwork already in place.

Both items below add new `@abstractmethod` entries to the `Backend` ABC
(`transformplan/backends/base.py`). That breaks any backend implemented outside
this repository, so they belong in a minor release (0.3.0), not a patch.

---

## Aggregation — `group_by` with aggregates

**Status:** planned, largest known gap.

There is `rows_pivot` and `rows_melt`, but no way to group with aggregation. Any
pipeline whose job is to condense onto an entity — one row per case, customer, or
transaction — has to leave the plan at exactly that point. The consequence is that
the most complex part of such a pipeline is the one part with **no protocol**,
while the simpler row and column operations are documented in full.

### Proposed API

```python
.group_by("case_id", {
"n_positions": ("*", "count"),
"total_amount": ("amount", "sum"),
"n_categories": ("field", "n_unique"),
})
```

### Groundwork already present

Two pieces exist and should be reused rather than rebuilt:

- **`AggFunction`** (`backends/base.py`) already defines the aggregate vocabulary
`Literal["first", "sum", "mean", "median", "min", "max", "count"]`, used today by
`rows_pivot` and `math_diff_from_agg`. `group_by` should extend this literal
(`n_unique`, `last`, `std`) rather than introduce a second vocabulary.
- **`ChunkMode.GROUP_DEPENDENT` with `group_param`** (`chunking.py`) already models
"needs all rows of a group together", used by `math_cumsum`, `math_rank` and
`rows_unique`. Registering `group_by` as
`OperationMeta(ChunkMode.GROUP_DEPENDENT, group_param="by")` makes it correct
under `process_chunked()` whenever `partition_key == by` — so aggregation fits
*inside* the chunking model instead of standing beside it.

### Where the real cost is

Not in the backends. The expensive part is `SchemaTracker` in `validation.py`:
every existing validator *advances* the schema, whereas `group_by` must **replace**
it — after the step, only the grouping columns and the aggregate outputs exist.
This is the one place where the current forward-propagation logic does not carry
over. Beyond that: DuckDB SQL generation and `dry_run` output.

### Constraint to hold

Keep the aggregate vocabulary **closed** — names only, never arbitrary callables.
A lambda cannot be serialized, and serializability is the library's core promise.

---

## Time components — `dt_hour`, `dt_minute`, `dt_second`

**Status:** planned, small.

The `dt_` family covers date components thoroughly but has no time components.
Turning `16:30` into `16.5` currently takes six steps and a detour through text
formatting and back:

```python
.dt_format("admitted_at", "%H", "admit_h")
.dt_format("admitted_at", "%M", "admit_m")
.col_cast("admit_h", "Float64")
.col_cast("admit_m", "Float64")
.math_divide("admit_m", 60)
.math_add_columns("admit_h", "admit_m", "admit_hour")
```

With `dt_hour` and `dt_minute` this becomes three steps and no string round trip.

### Notes for implementation

- DuckDB is trivial here: `EXTRACT(hour FROM …)`.
- Each new operation touches **nine files**: `ops/datetime.py`, `backends/base.py`,
`backends/polars.py`, `backends/duckdb.py`, `validation.py` (validator plus
registry), `chunking.py` (registry), `tests/test_datetime.py`,
`tests/test_duckdb.py`, and `docs/api/ops/datetime.md`. Budget for that, not
just for the operation itself.
- **`dt_decimal_hour` is deliberately excluded.** With `dt_hour`/`dt_minute` the
block is already down to three steps, and the library has no other convenience
combinations. Adding one opens a door to an unbounded set of special cases.

---

## Delivered in this release

For reference, the related suggestions that **are** implemented:

| Item | What shipped |
|---|---|
| `.pipe()` | Apply a function that takes and returns a plan |
| `.extend()` / `+` | Compose plans from reusable blocks |
| `.section()` | Named sections, grouped with a balance in the protocol |
| `.because()` | Reasons recorded per step in the protocol |
| Column sequences | 24 in-place operations accept `str \| Sequence[str]` |
| `col_cast` names | Canonical dtype names keep plans serializable |
1 change: 1 addition & 0 deletions mkdocs.yml
Original file line number Diff line number Diff line change
Expand Up @@ -74,6 +74,7 @@ markdown_extensions:

nav:
- Home: index.md
- Roadmap: roadmap.md
- Getting Started:
- Installation: getting-started/installation.md
- Quickstart: getting-started/quickstart.md
Expand Down
Loading
Loading