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

> Learn to implement data cleaning and transformation pipelines with our guide. Discover a three-stage architecture using pandas or Apache Flink for effective data processing and feature engineering.

- Repository: [DataExpert.io/data-engineer-handbook](https://github.com/DataExpert-io/data-engineer-handbook)
- Tags: how-to-guide
- Published: 2026-08-09

---

**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`](https://github.com/DataExpert-io/data-engineer-handbook/blob/main/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`](https://github.com/DataExpert-io/data-engineer-handbook/blob/main/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`](https://github.com/DataExpert-io/data-engineer-handbook/blob/main/data_cleaning.md) file provides a minimal end-to-end example that sequentially executes each pipeline stage:

```python
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:

```python
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")

```

### Stream Processing with Apache Flink

For real-time pipelines, the Flink tutorial in [`intermediate-bootcamp/materials/4-apache-flink-training/README.md`](https://github.com/DataExpert-io/data-engineer-handbook/blob/main/intermediate-bootcamp/materials/4-apache-flink-training/README.md) shows how to apply cleaning logic within a streaming context:

```java
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:

```python
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`](https://github.com/DataExpert-io/data-engineer-handbook/blob/main/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`](https://github.com/DataExpert-io/data-engineer-handbook/blob/main/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`](https://github.com/DataExpert-io/data-engineer-handbook/blob/main/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.