Data pipelines & processing
ETL-style flows, batch transforms, validation, and moving data between formats.
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.
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": …
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.
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)…
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.
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__":
…
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.
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…
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.
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)
…
Parallel Extract Multiple Sources with Threads in Python
Extract data from multiple sources in parallel using ThreadPoolExecutor and verify results match sequential processing.
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…
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.
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
…
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.
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(…
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.
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…
Union Multiple DataFrames with Aligned Columns in Python
Concatenate DataFrames with different columns, aligning them and filling missing values with NaN using pandas concat.
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({
…
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.
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")
…
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.