Reference library

Data pipelines & processing

ETL-style flows, batch transforms, validation, and moving data between formats.

44 matches
Data pipelines & processing easy

How to Partition Output Files by Date Key in Python

Group output files into a dictionary partitioned by a YYYYMMDD date key extracted from the filename prefix.

file-partitioning date-key pathlib
Python
from pathlib import Path
from collections import defaultdict

def partition_files_by_date(directory: str) -> dict:
    """Partition output files by date key extracted from filename (YYYYMMDD prefix)."""
    path = Path(directory)
    partitions = defaultdict(list)
    
    for file in path.iterdir():
        if file.i…
14 0 Open
Data pipelines & processing easy

How to Reduce Aggregate Counts from Mapped Chunks in Python

Combine a list of mapped chunk dictionaries into a single aggregated count dictionary using functools.reduce.

reduce aggregation dictionary
Python
from functools import reduce
from collections import defaultdict

def aggregate_chunks(mapped_chunks):
    """Combine mapped chunk counts into a single aggregate dict."""
    return reduce(
        lambda acc, chunk: {
            **acc,
            **{k: acc.get(k, 0) + v for k, v in chunk.items()}
        },
       …
14 0 Open
Data pipelines & processing easy

How to Register a Dataset Schema as JSON in Python

Define a catalog of dataset schemas and serialize them to JSON with the standard library json module.

json schema catalog
Python
import json

catalog = {
    "name": "sample_catalog",
    "version": "1.0",
    "datasets": [
        {
            "id": "users",
            "type": "table",
            "fields": [
                {"name": "id", "type": "integer", "key": True},
                {"name": "email", "type": "string", "nullable": False}…
13 0 Open
Data pipelines & processing easy

How to Safely Coerce Strings to Numbers in Python

A safe conversion function that turns strings into integers or floats, returning a fallback value when conversion fails.

type-conversion robust-parsing data-cleaning
Python
import math

def to_number(value, fallback=None):
    """Safely coerce a string to int or float, returning fallback on failure."""
    if isinstance(value, (int, float)):
        return value
    try:
        # Try int first for clean whole numbers
        return int(value)
    except (ValueError, TypeError):
        …
12 0 Open
Data pipelines & processing easy

How to Sort a List of Dictionaries by Key in Python

A reusable helper function that sorts a list of dictionaries by a specified key, with optional descending order support.

sorting dictionaries data-pipelines
Python
from typing import List

def sort_records(records: List[dict], key: str, descending: bool = False) -> List[dict]:
    """Sort a list of dictionaries by a specified key."""
    return sorted(records, key=lambda record: record[key], reverse=descending)


def demonstrate_sorting() -> None:
    users = [
        {"name": …
12 0 Open
Data pipelines & processing easy

How to Track Checkpoint Offset After Batch Commit in Python

A batch processor that tracks the last successfully committed offset after processing records in batches, advancing the checkpoint only when each batch commits successfully.

batch-processing checkpoint offset
Python
import json
from typing import Any


class BatchProcessor:
    """Tracks checkpoint offset after committing batches."""

    def __init__(self, batch_size: int = 3):
        self.batch_size = batch_size
        self.offset = 0  # last successfully committed offset (exclusive)
        self.total_committed = 0

    def …
12 0 Open
Data pipelines & processing easy

How to Unpivot Wide to Long with pandas melt in Python

This code demonstrates how to use pandas.melt to unpivot a wide DataFrame into a tidy long format, converting subject columns into rows.

pandas melt reshape
Python
import pandas as pd

# Sample wide-format data
df_wide = pd.DataFrame({
    'id': [1, 2, 3],
    'name': ['Alice', 'Bob', 'Charlie'],
    'math': [90, 85, 95],
    'science': [80, 92, 88]
})

print("Original wide DataFrame:")
print(df_wide)

# Melt: unpivot subject columns into rows
df_long = pd.melt(
    df_wide,
   …
15 0 Open
Data pipelines & processing easy

How to Validate Data in a Python Pipeline

A helper module to validate common record types — email, positive integer, and non-empty string list — before processing data in a pipeline.

data-validation pipelines type-checking
Python
from typing import Any, Iterable


def is_valid_email(email: str) -> bool:
    """Basic email check: one '@', no spaces, dot after '@'."""
    if "@" not in email or " " in email:
        return False
    local, _, domain = email.partition("@")
    return bool(local) and "." in domain


def is_positive_int(value: Any)…
12 0 Open
Data pipelines & processing easy

How to create a dated snapshot path for a dataset in Python

Generate a versioned directory path combining a base directory, dataset name, and today's date, ready for creating snapshots in data pipelines.

date pathlib datasets
Python
import datetime
import os
from pathlib import Path


def snapshot_path(base_dir: str, dataset_name: str) -> Path:
    """Return a dated snapshot path for a dataset under a base directory."""
    today = datetime.date.today().isoformat()
    return Path(base_dir) / dataset_name / today


if __name__ == "__main__":
    …
15 0 Open
Data pipelines & processing easy

How to detect anomalies in a column using z-score in Python

Detect outliers in a list of numbers using z-score statistics, flagging values that deviate significantly from the mean.

anomaly-detection z-score statistics
Python
import random

def z_score_anomaly_detection(data, threshold=2.0):
    """
    Detect anomalies in a list of numbers using z-score.
    """
    mean = sum(data) / len(data)
    variance = sum((x - mean) ** 2 for x in data) / len(data)
    std_dev = variance ** 0.5
    
    if std_dev == 0:
        return []
    
    a…
14 0 Open
Data pipelines & processing easy

How to route late-arriving data to a side output in Python

Separate late-arriving events from a streaming data batch into a dead-letter side output list using a timestamp threshold.

data pipelines streaming dead-letter
Python
from collections import defaultdict

def late_arriving_side_output(events, late_threshold_ts):
    """
    Mock a streaming pipeline that separates late-arriving data events
    into a side output list (e.g., for dead-letter analysis).

    events: list of (timestamp, data) tuples, timestamps as ints.
    late_thresho…
12 0 Open
Data pipelines & processing easy

Idempotent Pipeline Dedupe by Record ID Set in Python

Filters records against a persistent set of seen IDs, returning only new ones and the updated set for idempotent pipeline processing.

deduplication idempotency pipelines
Python
def dedupe_records(records, seen_ids=None):
    """Return records whose id has not been seen before."""
    if seen_ids is None:
        seen_ids = set()
    unique = []
    for record in records:
        record_id = record.get("id")
        if record_id not in seen_ids:
            seen_ids.add(record_id)
           …
12 0 Open
Data pipelines & processing easy

Implement Exactly-Once Transaction Log in Python

A mock transaction log that deduplicates transaction IDs so each is recorded only once, with a dataclass for records and simple in-memory storage.

transactions deduplication dataclass
Python
from dataclasses import dataclass
from typing import Dict, Optional


@dataclass
class TxnRecord:
    txn_id: str
    status: str


class ExactlyOnceTxnLog:
    def __init__(self) -> None:
        self._log: Dict[str, TxnRecord] = {}
        self._processed_ids: set = set()

    def record(self, txn_id: str, status: s…
14 0 Open
Data pipelines & processing easy

Parallel Extract Multiple Sources with Threads in Python

Extract data from multiple sources in parallel using ThreadPoolExecutor and verify results match sequential processing.

threads threadpoolexecutor concurrency
Python
import threading
from concurrent.futures import ThreadPoolExecutor

def extract_from_source(source):
    """Simulate extracting data from a source."""
    return f"Data from {source}"

def main():
    sources = ["source_a", "source_b", "source_c", "source_d"]
    
    # Sequential extraction for comparison
    sequent…
14 0 Open
Data pipelines & processing easy

Pipeline stage compose functions left to right in Python

Compose multiple functions into a left-to-right pipeline so each stage receives the output of the previous one.

composition pipeline functional
Python
def compose(*funcs):
    """Compose functions left to right: compose(f, g, h)(x) == h(g(f(x)))"""
    def composed(arg):
        result = arg
        for func in funcs:
            result = func(result)
        return result
    return composed

if __name__ == "__main__":
    def add_one(x):
        return x + 1

    …
16 0 Open
Data pipelines & processing easy

Rollback dataset to previous snapshot pointer in Python

A SnapshotManager class stores timestamped data snapshots and rolls back to the most recent snapshot at or before a target time.

snapshots rollback datetime
Python
from datetime import datetime, timedelta


class SnapshotManager:
    def __init__(self):
        self.snapshots = {}  # timestamp -> data
        self.current_pointer = None

    def create_snapshot(self, data):
        timestamp = datetime.now()
        self.snapshots[timestamp] = data
        self.current_pointer =…
13 0 Open
Data pipelines & processing easy

Test a Python Pipeline with Fixture Sample Rows

Test pipeline functions with sample rows provided by a pytest fixture, verifying required keys and value constraints.

pytest fixtures data-pipelines
Python
import pytest


def get_value(data: dict, key: str):
    return data.get(key)


def sample_rows():
    return [
        {"name": "Alice", "age": 30, "city": "London"},
        {"name": "Bob", "age": 25, "city": "Paris"},
        {"name": "Charlie", "age": 35, "city": "Berlin"},
    ]


@pytest.fixture
def sample_data(…
16 0 Open
Data pipelines & processing easy

Trigger a Pipeline When a New File Appears in a Directory

Poll a directory every 0.5 seconds and return the name of the first new file that appears, or None after a timeout.

polling filesystem file-watcher
Python
import time
from pathlib import Path


def watch_for_file(directory: str, interval: float = 0.5, timeout: float = 10.0) -> str | None:
    """Poll a directory and trigger when a new file appears."""
    watch_dir = Path(directory)
    watch_dir.mkdir(exist_ok=True)
    
    known_files = set(watch_dir.iterdir())
    s…
13 0 Open
Data pipelines & processing easy

Union Multiple DataFrames with Aligned Columns in Python

Concatenate DataFrames with different columns, aligning them and filling missing values with NaN using pandas concat.

pandas dataframes concat
Python
import pandas as pd
from io import StringIO

# Sample dataframes with different columns
df1 = pd.DataFrame({
    'id': [1, 2, 3],
    'name': ['Alice', 'Bob', 'Charlie'],
    'age': [25, 30, 35]
})

df2 = pd.DataFrame({
    'id': [4, 5],
    'name': ['Diana', 'Eve'],
    'city': ['NYC', 'LA']
})

df3 = pd.DataFrame({
…
14 0 Open
Data pipelines & processing easy

Validate dict schema at pipeline boundary in Python

This code validates a dictionary against a TypedDict schema at a pipeline boundary, enforcing required fields and types with custom error messages.

validation dict typeddict
Python
from typing import Any, TypedDict


class Person(TypedDict):
    name: str
    age: int
    email: str


def validate_person(data: dict[str, Any]) -> Person:
    errors: list[str] = []

    if not isinstance(data.get("name"), str) or not data["name"].strip():
        errors.append("name must be a non-empty string")
  …
13 0 Open

Browse by section

Each section groups closely related Python snippets.

Data pipelines & processing — Python code examples

What you will find here

This page collects data pipelines & processing snippets — short, copy-ready Python you can paste into our free online IDE and run without installing anything. Each sample includes a plain-English explanation and the full source code.

Samples vs tutorials and challenges

Samples are quick reference — one concept per page. For step-by-step teaching, use our Python tutorials. To test yourself, try quizzes or coding challenges. Clean up style with the Python formatter.