How to Implement Data Cleaning and Transformation Pipelines: A Complete Guide

Implement data cleaning and transformation pipelines using a three-stage architecture—ingestion, cleaning, and transformation—leveraging pandas for batch processing or Apache Flink for streaming, with deterministic rules for deduplication, type enforcement, and feature engineering.

Data cleaning and transformation pipelines form the backbone of trustworthy analytics workflows. According to the DataExpert-io/data-engineer-handbook repository, a robust pipeline follows a layered architecture that ensures data integrity from raw ingestion to final publication. This guide demonstrates how to implement these pipelines using the exact patterns and code samples documented in the source files.

Three-Stage Pipeline Architecture

The handbook advocates for a logical three-stage design that separates concerns and facilitates testing. Each stage operates as a pure function with no side effects, enabling reproducible data workflows.

Ingestion Layer

The ingestion layer pulls raw records from files, APIs, or message queues. Key considerations include schema-on-read validation, back-pressure handling, and idempotency to prevent duplicate processing.

Typical tools include pandas.read_csv(), Spark's read API, or Flink's FlinkKafkaConsumer for streaming sources.

Cleaning Layer

The cleaning layer produces a deterministic, de-duplicated, type-safe dataset. As documented in data_cleaning.md, this stage applies deterministic rules to eliminate duplicates, standardize column names, resolve missing values, and enforce correct data types.

Critical operations include:

  • Deduplication – df.drop_duplicates() prevents data leakage (see line 4 of the handbook's cleaning script)
  • Column standardization – Convert to lowercase and underscore naming using df.columns = [c.lower().replace(" ", "_") for c in df.columns] (line 6)
  • Missing-value strategy – Fill numeric columns with median values: df[num_cols] = df[num_cols].fillna(df[num_cols].median()) (line 20)
  • Date parsing – Enforce consistent datetime types with pd.to_datetime() (lines 22-24)

Transformation Layer

The transformation layer converts the cleaned dataset into the shape required by downstream models or dashboards. This stage involves feature engineering, windowing, aggregations, and format conversions (e.g., JSON to Parquet, or batch to streaming).

The Apache Flink training material in intermediate-bootcamp/materials/4-apache-flink-training/README.md demonstrates how these transformation concepts apply to real-time data streams using ProcessFunction and windowing operations.

Loading and Publishing

The final stage persists data to warehouses, data lakes, or exposes it via APIs. Best practices include partitioning by time or key, using to_parquet() for columnar storage, and maintaining lineage metadata for traceability.

End-to-End Implementation Examples

Batch Processing with Pandas

The data_cleaning.md file provides a minimal end-to-end example that sequentially executes each pipeline stage:

import pandas as pd

# 1️⃣ Ingest

df = pd.read_csv("s3://my-bucket/raw/events.csv")

# 2️⃣ Clean

df = df.drop_duplicates()
df.columns = [c.lower().replace(" ", "_") for c in df.columns]

num_cols = df.select_dtypes(include="number").columns
df[num_cols] = df[num_cols].fillna(df[num_cols].median())

if "date" in df.columns:
    df["date"] = pd.to_datetime(df["date"])

# 3️⃣ Transform – simple feature engineering

df["year"] = df["date"].dt.year
df["is_weekend"] = df["date"].dt.weekday >= 5
df["revenue_per_user"] = df["revenue"] / df["active_users"]

# 4️⃣ Load

df.to_parquet("s3://my-bucket/cleaned/events.parquet", partition_cols=["year"])

Batch Processing with PySpark

For larger datasets, the handbook recommends PySpark for distributed cleaning:

spark = SparkSession.builder.getOrCreate()
raw = spark.read.csv("s3://my-bucket/raw/events.csv", header=True, inferSchema=True)
clean = raw.dropDuplicates()
clean = clean.withColumnRenamed("Old Name", "old_name")
clean = clean.na.fill({"numeric_col": raw.agg({"numeric_col":"median"}).first()[0]})
clean.write.partitionBy("year").parquet("s3://my-bucket/cleaned/events.parquet")

For real-time pipelines, the Flink tutorial in intermediate-bootcamp/materials/4-apache-flink-training/README.md shows how to apply cleaning logic within a streaming context:

DataStream<String> source = env.addSource(new FlinkKafkaConsumer<>(topic, new SimpleStringSchema(), props));
DataStream<Event> events = source.map(json -> parse(json))
    .keyBy(Event::getId)
    .process(new DeduplicationFunction())
    .map(event -> standardise(event))
    .filter(event -> event.isValid());

Orchestration with Apache Airflow

Production pipelines require orchestration. The handbook suggests wrapping these stages in an Airflow DAG:

from airflow import DAG
from airflow.operators.python import PythonOperator

def run_pipeline():
    # call the pandas script or spark job

    pass

with DAG('data_clean_transform', schedule='@daily') as dag:
    task = PythonOperator(task_id='run', python_callable=run_pipeline)

Production Maintenance and Monitoring

The intermediate-bootcamp/materials/6-data-pipeline-maintenance/README.md file emphasizes operational concerns for production environments. Any failure in cleaning or transformation should abort the job, preserving data integrity through atomic transactions.

Key operational practices include:

  • Versioning cleaning logic alongside code
  • Logging the number of rows removed or filled during cleaning
  • Maintaining run-book templates for pipeline ownership
  • Implementing lineage tracking to trace data from source to sink

Summary

  • Implement data cleaning and transformation pipelines using a three-stage architecture: ingestion, cleaning, and transformation.
  • Use pandas for batch workflows and Apache Flink for streaming applications, as demonstrated in data_cleaning.md and the Flink training materials.
  • Apply deterministic cleaning rules: deduplicate with drop_duplicates(), standardize columns to lowercase/underscore, fill missing numerics with median values, and parse dates with pd.to_datetime().
  • Ensure transformations are pure functions with no side effects to facilitate unit testing and reproducibility.
  • Orchestrate with Airflow and maintain pipeline lineage metadata for production reliability.

Frequently Asked Questions

What is the difference between data cleaning and data transformation?

Data cleaning focuses on fixing or removing incorrect, corrupted, or duplicate data—such as applying drop_duplicates() or fillna()—while data transformation converts the cleaned data into the required format for analysis, such as feature engineering, aggregations, or format conversions. According to the handbook's architecture, cleaning ensures type safety and deduplication, whereas transformation handles business logic and feature derivation.

How should missing values be handled in production pipelines?

The handbook recommends filling numeric columns with median values using df[num_cols].fillna(df[num_cols].median()) and explicitly logging the count of filled values for observability. Always keep the missing-value strategy versioned and deterministic to ensure reproducibility across pipeline runs.

When should I use batch processing versus streaming for data cleaning?

Use batch processing with pandas or PySpark for historical data loads, daily ETL jobs, or when processing large static datasets, as shown in data_cleaning.md. Use Apache Flink for streaming when you require real-time data cleaning and low-latency transformations, particularly when ingesting from Kafka or other message queues, as documented in the Flink training README.

How do I ensure data quality during the transformation stage?

Ensure transformations are implemented as pure functions with no side effects, enabling unit testing and reproducibility. The handbook recommends aborting the entire pipeline if any cleaning or transformation step fails, preventing partial or corrupted data from reaching downstream systems. Maintain lineage metadata and partition output data by time or key for traceability.

Have a question about this repo?

These articles cover the highlights, but your codebase questions are specific. Give your agent direct access to the source. Share this with your agent to get started:

Share the following with your agent to get started:
curl -s "https://instagit.com/install.md"

Works with
Claude Codex Cursor VS Code OpenClaw Any MCP Client

Maintain an open-source project? Get it listed too →