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
167 changes: 167 additions & 0 deletions docs/plans/2026-03-08-data-pipeline-design.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,167 @@
# Data Pipeline Example — Design

**Date:** 2026-03-08
**Pattern:** `trigger → data → transform → data`
**Status:** Approved

## Goal

Fifth runnable example for workflow-patterns. A data processing pipeline that reads structured data, applies transformations (filtering, aggregation, enrichment), and writes results. Demonstrates the `data → transform → data` pattern — read, process, write.

## Architecture

```
select_pipeline() → load_data() → transform_data() → save_results()
trigger data transform data
```

## Module Structure

```
examples/data-pipeline/
├── run.py # CLI entry point
├── .env.example # No config needed
├── .gitignore # .env, __pycache__, output/
├── pyproject.toml # Zero runtime dependencies (stdlib only)
├── src/pipeline/
│ ├── __init__.py
│ ├── models.py # Dataclasses: Record, Dataset, Pipeline, TransformResult
│ ├── pipelines.py # 5 pipeline presets
│ ├── reader.py # Load data from CSV/JSON (data read layer)
│ ├── transforms.py # Transform functions (filter, aggregate, enrich)
│ ├── writer.py # Save results to CSV/JSON (data write layer)
│ └── display.py # Terminal formatting
├── tests/
│ ├── test_models.py
│ ├── test_pipelines.py
│ ├── test_reader.py
│ ├── test_transforms.py
│ ├── test_writer.py
│ ├── test_display.py
│ └── __init__.py
├── data/ # Sample input data (CSV)
│ └── sales_data.csv
└── output/ # Pipeline results
```

## Modules

### models.py

```python
@dataclass
class Record:
data: dict[str, str]

@dataclass
class Dataset:
headers: list[str]
records: list[Record]

@dataclass
class TransformStep:
name: str
description: str
operation: str # "filter", "aggregate", "sort", "add_column", "rename"
params: dict

@dataclass
class Pipeline:
name: str
description: str
steps: list[TransformStep]
```

### pipelines.py

5 curated pipeline presets:

| Pipeline | Focus |
|----------|-------|
| Sales Summary | Filter by region, aggregate revenue by product |
| Top Performers | Sort by metric, take top N |
| Data Cleanup | Remove nulls, normalize formats, deduplicate |
| Period Comparison | Filter by date range, compute period-over-period |
| Custom Report | Add computed columns, rename headers, export |

### reader.py

- `read_csv(path)` → Dataset from CSV file
- `read_json(path)` → Dataset from JSON file
- `from_dicts(headers, records)` → Dataset from in-memory data
- `generate_sample_data()` → Realistic sample sales dataset

### transforms.py

- `filter_records(dataset, column, operator, value)` → filtered Dataset
- `sort_records(dataset, column, ascending)` → sorted Dataset
- `aggregate(dataset, group_by, agg_column, operation)` → aggregated Dataset
- `add_column(dataset, name, expression)` → Dataset with computed column
- `rename_columns(dataset, mapping)` → Dataset with renamed headers
- `deduplicate(dataset, columns)` → deduplicated Dataset
- `apply_pipeline(dataset, steps)` → run all steps, return TransformResult

### writer.py

- `write_csv(dataset, path)` → save as CSV
- `write_json(dataset, path)` → save as JSON
- `to_table_string(dataset, max_rows)` → formatted ASCII table

### display.py

- `format_header(title)` → boxed ASCII header
- `format_pipeline_menu(pipelines)` → numbered pipeline list
- `format_stats(dataset)` → record count, column count

## UX Flow

```
$ uv run python run.py

╔══════════════════════════════════════════╗
║ Data Pipeline — Setup ║
╚══════════════════════════════════════════╝

Choose a pipeline:
1. Sales Summary Filter by region, aggregate revenue
2. Top Performers Sort by metric, take top N
3. Data Cleanup Remove nulls, normalize, deduplicate
4. Period Comparison Filter date range, period-over-period
5. Custom Report Computed columns, rename, export

Pipeline (1-5): 1

── Sales Summary ──

Loading data...
Loaded 50 records, 6 columns

Applying transforms:
✓ Filter: region = "EMEA" → 18 records
✓ Aggregate: sum revenue by product → 5 records

┌──────────┬─────────┐
│ product │ revenue │
├──────────┼─────────┤
│ Widget A │ 45,200 │
│ Widget B │ 32,100 │
│ ... │ ... │
└──────────┴─────────┘

Results saved to output/2026-03-08_sales-summary.csv
```

## Dependencies

- Zero runtime dependencies (stdlib only)
- Uses `csv`, `json`, `pathlib` from stdlib
- `pytest` (dev)

## Testing Strategy

- CSV/JSON read/write round-trip tests
- Each transform function with edge cases
- Pipeline application with multiple steps
- Sample data generation
- Display formatting
- Target: ~25 tests
2 changes: 2 additions & 0 deletions examples/data-pipeline/.env.example
Original file line number Diff line number Diff line change
@@ -0,0 +1,2 @@
# Data Pipeline — no API keys required
# All functionality works with stdlib only
4 changes: 4 additions & 0 deletions examples/data-pipeline/.gitignore
Original file line number Diff line number Diff line change
@@ -0,0 +1,4 @@
.env
__pycache__/
output/
*.pyc
121 changes: 121 additions & 0 deletions examples/data-pipeline/data/sales_data.csv
Original file line number Diff line number Diff line change
@@ -0,0 +1,121 @@
date,region,product,units,revenue,category
2026-01-01,EMEA,Widget A,43,5160,hardware
2026-01-01,EMEA,Widget B,26,2210,hardware
2026-01-01,EMEA,Widget C,23,4600,hardware
2026-01-01,EMEA,Service X,24,8400,services
2026-01-01,EMEA,Service Y,17,2550,services
2026-01-01,APAC,Widget A,42,4536,hardware
2026-01-01,APAC,Widget B,13,994,hardware
2026-01-01,APAC,Widget C,41,7380,hardware
2026-01-01,APAC,Service X,5,1575,services
2026-01-01,APAC,Service Y,19,2565,services
2026-01-01,NA,Widget A,27,3564,hardware
2026-01-01,NA,Widget B,22,2057,hardware
2026-01-01,NA,Widget C,22,4840,hardware
2026-01-01,NA,Service X,26,10010,services
2026-01-01,NA,Service Y,28,4620,services
2026-01-01,LATAM,Widget A,19,1824,hardware
2026-01-01,LATAM,Widget B,32,2176,hardware
2026-01-01,LATAM,Widget C,36,5760,hardware
2026-01-01,LATAM,Service X,36,10080,services
2026-01-01,LATAM,Service Y,33,3960,services
2026-02-01,EMEA,Widget A,27,3240,hardware
2026-02-01,EMEA,Widget B,38,3230,hardware
2026-02-01,EMEA,Widget C,41,8200,hardware
2026-02-01,EMEA,Service X,23,8050,services
2026-02-01,EMEA,Service Y,44,6600,services
2026-02-01,APAC,Widget A,22,2376,hardware
2026-02-01,APAC,Widget B,42,3213,hardware
2026-02-01,APAC,Widget C,41,7380,hardware
2026-02-01,APAC,Service X,28,8820,services
2026-02-01,APAC,Service Y,30,4050,services
2026-02-01,NA,Widget A,39,5148,hardware
2026-02-01,NA,Widget B,8,748,hardware
2026-02-01,NA,Widget C,28,6160,hardware
2026-02-01,NA,Service X,33,12705,services
2026-02-01,NA,Service Y,26,4290,services
2026-02-01,LATAM,Widget A,44,4224,hardware
2026-02-01,LATAM,Widget B,41,2788,hardware
2026-02-01,LATAM,Widget C,10,1600,hardware
2026-02-01,LATAM,Service X,23,6440,services
2026-02-01,LATAM,Service Y,21,2520,services
2026-03-01,EMEA,Widget A,43,5160,hardware
2026-03-01,EMEA,Widget B,19,1615,hardware
2026-03-01,EMEA,Widget C,37,7400,hardware
2026-03-01,EMEA,Service X,36,12600,services
2026-03-01,EMEA,Service Y,30,4500,services
2026-03-01,APAC,Widget A,42,4536,hardware
2026-03-01,APAC,Widget B,32,2448,hardware
2026-03-01,APAC,Widget C,29,5220,hardware
2026-03-01,APAC,Service X,21,6615,services
2026-03-01,APAC,Service Y,27,3645,services
2026-03-01,NA,Widget A,17,2244,hardware
2026-03-01,NA,Widget B,19,1777,hardware
2026-03-01,NA,Widget C,35,7700,hardware
2026-03-01,NA,Service X,38,14630,services
2026-03-01,NA,Service Y,39,6435,services
2026-03-01,LATAM,Widget A,14,1344,hardware
2026-03-01,LATAM,Widget B,38,2584,hardware
2026-03-01,LATAM,Widget C,43,6880,hardware
2026-03-01,LATAM,Service X,34,9520,services
2026-03-01,LATAM,Service Y,25,3000,services
2026-04-01,EMEA,Widget A,35,4200,hardware
2026-04-01,EMEA,Widget B,13,1105,hardware
2026-04-01,EMEA,Widget C,21,4200,hardware
2026-04-01,EMEA,Service X,42,14700,services
2026-04-01,EMEA,Service Y,16,2400,services
2026-04-01,APAC,Widget A,17,1836,hardware
2026-04-01,APAC,Widget B,12,918,hardware
2026-04-01,APAC,Widget C,38,6840,hardware
2026-04-01,APAC,Service X,16,5040,services
2026-04-01,APAC,Service Y,27,3645,services
2026-04-01,NA,Widget A,20,2640,hardware
2026-04-01,NA,Widget B,42,3927,hardware
2026-04-01,NA,Widget C,39,8580,hardware
2026-04-01,NA,Service X,9,3465,services
2026-04-01,NA,Service Y,43,7095,services
2026-04-01,LATAM,Widget A,28,2688,hardware
2026-04-01,LATAM,Widget B,15,1020,hardware
2026-04-01,LATAM,Widget C,8,1280,hardware
2026-04-01,LATAM,Service X,32,8960,services
2026-04-01,LATAM,Service Y,11,1320,services
2026-05-01,EMEA,Widget A,41,4920,hardware
2026-05-01,EMEA,Widget B,22,1870,hardware
2026-05-01,EMEA,Widget C,32,6400,hardware
2026-05-01,EMEA,Service X,36,12600,services
2026-05-01,EMEA,Service Y,33,4950,services
2026-05-01,APAC,Widget A,37,3996,hardware
2026-05-01,APAC,Widget B,6,459,hardware
2026-05-01,APAC,Widget C,28,5040,hardware
2026-05-01,APAC,Service X,32,10080,services
2026-05-01,APAC,Service Y,23,3105,services
2026-05-01,NA,Widget A,9,1188,hardware
2026-05-01,NA,Widget B,34,3179,hardware
2026-05-01,NA,Widget C,41,9020,hardware
2026-05-01,NA,Service X,13,5005,services
2026-05-01,NA,Service Y,28,4620,services
2026-05-01,LATAM,Widget A,6,576,hardware
2026-05-01,LATAM,Widget B,21,1428,hardware
2026-05-01,LATAM,Widget C,30,4800,hardware
2026-05-01,LATAM,Service X,29,8120,services
2026-05-01,LATAM,Service Y,44,5280,services
2026-06-01,EMEA,Widget A,13,1560,hardware
2026-06-01,EMEA,Widget B,44,3740,hardware
2026-06-01,EMEA,Widget C,41,8200,hardware
2026-06-01,EMEA,Service X,35,12250,services
2026-06-01,EMEA,Service Y,34,5100,services
2026-06-01,APAC,Widget A,15,1620,hardware
2026-06-01,APAC,Widget B,5,382,hardware
2026-06-01,APAC,Widget C,14,2520,hardware
2026-06-01,APAC,Service X,12,3780,services
2026-06-01,APAC,Service Y,28,3780,services
2026-06-01,NA,Widget A,9,1188,hardware
2026-06-01,NA,Widget B,44,4114,hardware
2026-06-01,NA,Widget C,29,6380,hardware
2026-06-01,NA,Service X,42,16170,services
2026-06-01,NA,Service Y,44,7260,services
2026-06-01,LATAM,Widget A,41,3936,hardware
2026-06-01,LATAM,Widget B,19,1292,hardware
2026-06-01,LATAM,Widget C,37,5920,hardware
2026-06-01,LATAM,Service X,9,2520,services
2026-06-01,LATAM,Service Y,20,2400,services
19 changes: 19 additions & 0 deletions examples/data-pipeline/pyproject.toml
Original file line number Diff line number Diff line change
@@ -0,0 +1,19 @@
[project]
name = "data-pipeline"
version = "0.1.0"
description = "Data Pipeline workflow: trigger -> data -> transform -> data"
requires-python = ">=3.12"
dependencies = []

[dependency-groups]
dev = ["pytest"]

[build-system]
requires = ["hatchling"]
build-backend = "hatchling.build"

[tool.hatch.build.targets.wheel]
packages = ["src/pipeline"]

[tool.pytest.ini_options]
testpaths = ["tests"]
Loading