| name | data-pipeline |
| description | Use this skill when the user asks to process, clean, transform, merge, validate, deduplicate, aggregate, convert, or export data files such as CSV, JSON, JSONL, Excel, spreadsheets, logs, or datasets. Triggers on phrases like clean this CSV, transform data, merge spreadsheets, dedupe rows, convert JSONL, build a data pipeline, validate records, aggregate this dataset, export to CSV, and process a large file. Use it before writing data-processing code or designing repeatable ETL-style workflows. |
| emoji | 🧩 |
| version | 1.1.0 |
| triggers | clean this CSV, transform data, merge spreadsheets, dedupe rows, convert JSONL, build a data pipeline, validate records, aggregate dataset, export to CSV, process large file, parse JSON, normalize spreadsheet, ETL workflow, data cleanup, convert data file |
Data Pipeline
Read, transform, and write data correctly. Use this before writing any data processing code.
1. Reading Data
CSV
import csv
with open("data.csv", newline="", encoding="utf-8") as f:
reader = csv.DictReader(f)
rows = list(reader)
rows = []
with open("data.csv", newline="", encoding="utf-8") as f:
for row in csv.DictReader(f):
rows.append({
"id": int(row["id"]),
"value": float(row["value"]),
"name": row["name"].strip()
})
JSON
import json
with open("data.json", encoding="utf-8") as f:
data = json.load(f)
records = []
with open("data.jsonl", encoding="utf-8") as f:
for line in f:
if line.strip():
records.append(json.loads(line))
Excel (openpyxl)
import openpyxl
wb = openpyxl.load_workbook("data.xlsx")
ws = wb.active
headers = [cell.value for cell in ws[1]]
rows = []
for row in ws.iter_rows(min_row=2, values_only=True):
rows.append(dict(zip(headers, row)))
2. Transforming Data
Clean & Normalize
def clean_row(row):
return {
"id": str(row.get("id", "")).strip(),
"email": row.get("email", "").strip().lower(),
"name": row.get("name", "").strip().title(),
"amount": float(row.get("amount") or 0),
}
cleaned = [clean_row(r) for r in rows if r.get("id")]
Filter
active = [r for r in rows if r["status"] == "active"]
recent = [r for r in rows if r["date"] >= "2024-01-01"]
Deduplicate
seen = set()
unique = []
for row in rows:
key = row["email"]
if key not in seen:
seen.add(key)
unique.append(row)
Aggregate
from collections import defaultdict
by_category = defaultdict(list)
for row in rows:
by_category[row["category"]].append(row)
totals = {cat: sum(r["amount"] for r in items) for cat, items in by_category.items()}
Sort
sorted_rows = sorted(rows, key=lambda r: r["date"], reverse=True)
sorted_rows = sorted(rows, key=lambda r: (r["category"], -r["amount"]))
3. Writing Output
CSV
import csv
with open("output.csv", "w", newline="", encoding="utf-8") as f:
if rows:
writer = csv.DictWriter(f, fieldnames=rows[0].keys())
writer.writeheader()
writer.writerows(rows)
JSON
with open("output.json", "w", encoding="utf-8") as f:
json.dump(data, f, indent=2, ensure_ascii=False, default=str)
with open("output.jsonl", "w", encoding="utf-8") as f:
for record in records:
f.write(json.dumps(record, ensure_ascii=False) + "\n")
4. Large Files (Chunked Processing)
Never load huge files entirely into memory. Process in chunks:
def process_large_csv(input_path, output_path, chunk_size=10_000):
processed = 0
with open(input_path, newline="", encoding="utf-8") as infile, \
open(output_path, "w", newline="", encoding="utf-8") as outfile:
reader = csv.DictReader(infile)
writer = None
chunk = []
for row in reader:
chunk.append(transform(row))
if len(chunk) >= chunk_size:
if writer is None:
writer = csv.DictWriter(outfile, fieldnames=chunk[0].keys())
writer.writeheader()
writer.writerows(chunk)
processed += len(chunk)
print(f"Processed {processed} rows...")
chunk = []
if chunk:
if writer is None:
writer = csv.DictWriter(outfile, fieldnames=chunk[0].keys())
writer.writeheader()
writer.writerows(chunk)
processed += len(chunk)
print(f"Done. Total: {processed} rows")
5. Validation
Always validate before writing output:
def validate_row(row, row_num):
errors = []
if not row.get("id"):
errors.append(f"Row {row_num}: missing id")
if not row.get("email") or "@" not in row["email"]:
errors.append(f"Row {row_num}: invalid email")
if row.get("amount") is not None and float(row["amount"]) < 0:
errors.append(f"Row {row_num}: negative amount")
return errors
all_errors = []
for i, row in enumerate(rows, 1):
all_errors.extend(validate_row(row, i))
if all_errors:
print(f"Found {len(all_errors)} validation errors:")
for e in all_errors[:20]:
print(f" {e}")
6. Pipeline Pattern (composable)
def run_pipeline(input_path, output_path):
print(f"Reading: {input_path}")
rows = read_csv(input_path)
print(f" Read {len(rows)} rows")
rows = [clean_row(r) for r in rows]
rows = [r for r in rows if r["id"]]
rows = deduplicate(rows, key="email")
print(f" After cleaning: {len(rows)} rows")
errors = []
for i, row in enumerate(rows, 1):
errors.extend(validate_row(row, i))
if errors:
raise ValueError(f"{len(errors)} validation errors found")
write_csv(rows, output_path)
print(f"Output written: {output_path} ({len(rows)} rows)")
return len(rows)
7. Quick Checklist
Before processing any data file: