Data Engineering Basics
9 examples to get you started with Data Engineering - 7 basic and 2 intermediate.
Search across all documentation pages
9 examples to get you started with Data Engineering - 7 basic and 2 intermediate.
uv pip install pandas pyarrow prefectLand raw data in columnar storage before transforms.
import pandas as pd
df = pd.read_csv("orders.csv", parse_dates=["ordered_at"], dtype={"region": "category"})
df.to_parquet("raw/orders.parquet", index=False, compression="zstd")dt=2025-01-15/ speed selective reads.Related: File Formats - Parquet and Arrow
Clean and shape data in a scriptable step.
import pandas as pd
df = pd.read_parquet("raw/orders.parquet")
clean = (
df.dropna(subset=["order_id"])
.assign(revenue=pd.to_numeric(df["revenue"], errors="coerce"))
.query("revenue >= 0")
)
clean.to_parquet("staging/orders_clean.parquet", index=False)Related: Cleaning & Transforming Data
Upsert or partition overwrite so retries do not duplicate.
import pandas as pd
from pathlib import Path
run_dt = "2025-01-15"
out = Path(f"mart/orders/dt={run_dt}/data.parquet")
out.parent.mkdir(parents=True, exist_ok=True)
df.to_parquet(out, index=False) # overwrite partition for this dtrun_dt replaces, not appends duplicates.Related: Workflow Reliability - idempotency decisions
Batch jobs run on a clock when latency tolerance is hours.
# crontab: 0 6 * * * /path/.venv/bin/python /path/jobs/daily_orders.py
import logging
logging.basicConfig(level=logging.INFO)
logging.info("daily_orders started")Related: Airflow - DAG scheduling
One function per step with typed inputs/outputs.
from pathlib import Path
import pandas as pd
def extract(path: Path) -> pd.DataFrame:
return pd.read_parquet(path)
def transform(df: pd.DataFrame) -> pd.DataFrame:
return df.groupby("region", observed=True)["revenue"].sum().reset_index()
def load(df: pd.DataFrame, path: Path) -> None:
path.parent.mkdir(parents=True, exist_ok=True)
df.to_parquet(path, index=False)Related: Prefect & Dagster
Reject bad data before it reaches marts.
import pandera.pandas as pa
schema = pa.DataFrameSchema({
"order_id": pa.Column(int, unique=True),
"revenue": pa.Column(float, pa.Check.ge(0)),
})
schema.validate(df)Related: Data Validation & Quality
Operational visibility starts with simple metrics.
import logging
logger = logging.getLogger(__name__)
def transform(df):
logger.info("input_rows=%s", len(df))
out = df.drop_duplicates("order_id")
logger.info("output_rows=%s", len(out))
return outRelated: Data Engineering Best Practices
Orchestrate Python tasks with retries and UI.
from prefect import flow, task
import pandas as pd
@task(retries=2, retry_delay_seconds=30)
def extract() -> pd.DataFrame:
return pd.read_parquet("raw/orders.parquet")
@flow(log_prints=True)
def daily_orders():
df = extract()
# transform + load ...
if __name__ == "__main__":
daily_orders()@task retries isolate transient S3/DB failures.Related: Prefect & Dagster
Scan only the dates you need.
import pandas as pd
df = pd.read_parquet("mart/orders", filters=[("dt", ">=", "2025-01-01")]scan_parquet offers the same predicate pushdown.Related: File Formats - partition layout
Stack versions: This page was written for Python 3.14.0 (stable 3.14, maintenance 3.13), FastAPI 0.115+, Django 5.2, Flask 3.1, Pydantic 2, PyTorch 2.6+, pandas 2.2+, Polars 1.x, ruff 0.9+, and uv 0.6+.
Reviewed by Chris St. John·Last updated Jul 16, 2026