Python & Data Science

Data Flow Decomposition Why Every Pandas Sklearn P

You’ve been there. It’s Monday morning, and you’re running last week’s data-cleaning script. It worked perfectly on Friday. Today, it crashes with a cryptic KeyError because a column name changed from 'customer_id' to 'CustomerID'. You fix that, run it again, and now groupby complains about a float column that used to be int64. You spend an hour debugging a pipeline you thought was finished.

Why does this feel so fragile? What’s the hidden structure underneath that makes one pipeline break and another one hum along?

Here’s the secret: every pandas method chain and every sklearn Pipeline is already a data-flow diagram. You just haven’t drawn the arrows yet. Once you see the diagram, the fragility becomes obvious — and fixable.

By the end of this article, you’ll see your own code as a network of transformations. You’ll know how to design new pipelines from scratch without trial-and-error. And you’ll have a five-step framework that turns messy code into a clean, testable, visual design.

Let’s start with a pipeline you’ve probably written yourself.

Opening: The Pipeline You Already Write

Here’s a realistic pandas pipeline. It loads customer transaction data, cleans it, groups it, and merges it with customer info. You’ve written something like this before:

import pandas as pd
import numpy as np

# Load raw data
df = pd.read_csv('transactions.csv')

# A typical chain of transformations
result = (
    df
    .dropna(subset=['customer_id', 'amount'])   # Input: DataFrame (10000, 5) -> Output: DataFrame (9500, 5)
    .groupby('customer_id')['amount']            # Input: DataFrame (9500, 5) -> Output: GroupBy object
    .agg(['sum', 'count'])                       # Input: GroupBy object -> Output: DataFrame (500, 2)
    .reset_index()                               # Input: DataFrame (500, 2) -> Output: DataFrame (500, 3)
    .merge(
        pd.read_csv('customers.csv'),            # Input: DataFrame (500, 3) + DataFrame (1000, 4)
        on='customer_id'
    )                                            # Output: DataFrame (500, 7)
)

Look at the comments. Each line has an input type and an output type. That’s your first data-flow diagram — right there in the code.

But here’s the problem: you didn’t write those comments. You wrote the code, ran it, and hoped it worked. When it broke, you had to trace through each step manually, checking types in your head.

What if you made those arrows explicit from the start? What if you designed the diagram before writing a single line of code?

That’s what data-flow decomposition gives you. Let’s learn the five steps.

Step 1: Name Every Input (The Data That Flows In)

Before any transformation, you must inventory everything that enters your system. This is harder than it sounds — most data scientists skip this and pay for it later.

INPUTS are every piece of data that enters your pipeline: raw CSVs, API responses, configuration parameters, random seeds. For each input, specify its type (e.g., pd.DataFrame, np.ndarray, dict) and its shape (rows, columns, dtypes).

Plain-English rule: If you can’t name what’s coming in, you can’t trust what’s going out.

Let’s use a concrete example. We’ll load the classic UCI Wine Quality dataset, which has mixed dtypes:

import pandas as pd
import numpy as np

# Load the dataset
url = 'https://archive.ics.uci.edu/ml/machine-learning-databases/wine-quality/winequality-red.csv'
df = pd.read_csv(url, sep=';')

# Inspect the inputs
print("Shape:", df.shape)
print("\nInfo:")
df.info()

Output:

Shape: (1599, 12)

Info:
<class 'pandas.core.frame.DataFrame'>
RangeIndex: 1599 entries, 0 to 1598
Data columns (total 12 columns):
 #   Column                Non-Null Count  Dtype
---  ------                --------------  -----
 0   fixed acidity         1599 non-null   float64
 1   volatile acidity      1599 non-null   float64
 2   citric acid           1599 non-null   float64
 3   residual sugar        1599 non-null   float64
 4   chlorides             1599 non-null   float64
 5   free sulfur dioxide   1599 non-null   float64
 6   total sulfur dioxide  1599 non-null   float64
 7   density               1599 non-null   float64
 8   pH                    1599 non-null   float64
 9   sulphates             1599 non-null   float64
 10  alcohol               1599 non-null   float64
 11  quality               1599 non-null   int64

Now list every input explicitly:

# INPUT INVENTORY
# Input 1: winequality-red.csv
#   Type: pd.DataFrame
#   Shape: (1599, 12)
#   Columns: fixed acidity (float64), volatile acidity (float64), ..., quality (int64)
#
# Input 2: random_seed (optional, for reproducibility)
#   Type: int
#   Default: 42

This inventory is your contract with the data. If a column changes from float64 to int64, you’ll catch it here before it breaks downstream.

Step 2: Name Every Output (What You Actually Need)

Most data scientists start coding without clearly defining what they want at the end. This is like driving without a destination — you’ll move, but you won’t know if you’ve arrived.

OUTPUTS are the final data product: a cleaned DataFrame, a trained model, a prediction array, a plot. Specify type and shape for each output, just like inputs.

Here’s the hard part: distinguishing between intermediate outputs (you’ll discard) and final outputs (you’ll deliver). For example, a grouped DataFrame is intermediate; the final churn prediction is final.

Let’s define the output for our wine quality pipeline:

def predict_wine_quality(df: pd.DataFrame) -> pd.DataFrame:
    """
    Predict wine quality score.
    
    Parameters
    ----------
    df : pd.DataFrame
        Input DataFrame with wine features (11 columns).
    
    Returns
    -------
    pd.DataFrame
        DataFrame with columns ['predicted_quality', 'confidence']
        Shape: (n_samples, 2)
    """
    # We'll fill this in later
    pass

Notice the docstring specifies the output shape: (n_samples, 2). This is your destination. Every transformation in your pipeline must lead here.

Step 3: Draw the Diagram (Boxes = Transformations, Arrows = Data)

Now the magic happens. Draw boxes for every transformation and arrows for data flow. This is the ‘aha’ moment where code becomes visual.

The concept comes from structured analysis, first described by Tom DeMarco in 1978 (Structured Analysis and System Specification) and later refined by Yourdon & Constantine in 1979 (Structured Design). The core idea: every system can be represented as a network of transformations (boxes) connected by data flows (arrows).

Rule: one box per distinct transformation. Don’t combine .dropna() and .groupby() into one box — they’re separate transformations.

Arrows carry data from one box to the next. Label each arrow with the type/shape of the data at that point.

Let’s see how a pandas chain maps to a data-flow diagram. We’ll write each transformation as a separate function and chain them manually:

def remove_missing(df: pd.DataFrame) -> pd.DataFrame:
    """Remove rows with missing values."""
    return df.dropna(subset=['quality'])

def engineer_features(df: pd.DataFrame) -> pd.DataFrame:
    """Create new features."""
    df = df.copy()
    df['acid_ratio'] = df['fixed acidity'] / df['volatile acidity']
    df['sulfur_ratio'] = df['free sulfur dioxide'] / df['total sulfur dioxide']
    return df

def scale_features(df: pd.DataFrame) -> pd.DataFrame:
    """Scale numeric features to [0, 1]."""
    from sklearn.preprocessing import MinMaxScaler
    scaler = MinMaxScaler()
    numeric_cols = df.select_dtypes(include=[np.number]).columns
    df[numeric_cols] = scaler.fit_transform(df[numeric_cols])
    return df

# Chain them manually
print("Step 1: remove_missing")
df_clean = remove_missing(df)
print(f"  Shape: {df_clean.shape}, Type: {type(df_clean)}")

print("Step 2: engineer_features")
df_feat = engineer_features(df_clean)
print(f"  Shape: {df_feat.shape}, Type: {type(df_feat)}")

print("Step 3: scale_features")
df_scaled = scale_features(df_feat)
print(f"  Shape: {df_scaled.shape}, Type: {type(df_scaled)}")

Output:

Step 1: remove_missing
  Shape: (1599, 12), Type: <class 'pandas.core.frame.DataFrame'>
Step 2: engineer_features
  Shape: (1599, 14), Type: <class 'pandas.core.frame.DataFrame'>
Step 3: scale_features
  Shape: (1599, 14), Type: <class 'pandas.core.frame.DataFrame'>

Now here’s the data-flow diagram in ASCII art:

[raw CSV] ---> [remove_missing] ---> [engineer_features] ---> [scale_features] ---> [cleaned DataFrame]
   |                  |                       |                        |
   v                  v                       v                        v
DataFrame(1599,12)  DataFrame(1599,12)     DataFrame(1599,14)       DataFrame(1599,14)

Your code is already a diagram — you just haven’t drawn it yet. Once you draw it, the structure becomes obvious.

Step 4: Label Every Arrow with Its Type (The Type-Checking Rule)

This is the hardest part. Every arrow must have a type, and the output type of one box must exactly match the input type of the next. A mismatch means you’ve missed a transformation — add a box.

The type-mismatch rule: if the output of .groupby() is a GroupBy object, but the next step expects a DataFrame, you need a .apply() or .agg() box to bridge the gap.

This is where most pipelines break silently. A column changes from float to int, or a DataFrame becomes a Series. Let’s see a concrete mismatch:

# This will crash because groupby returns a GroupBy object, not a DataFrame
try:
    result = (
        df
        .groupby('quality')['alcohol']  # Returns GroupBy object
        .mean()                          # Returns Series
        .reset_index()                   # Returns DataFrame
    )
    print("Success:")
    print(result.head())
except Exception as e:
    print(f"Error: {e}")

Output:

Success:
   quality    alcohol
0        3  10.210000
1        4  10.228261
2        5   9.899706
3        6  10.445192
4        7  11.466250

Wait — that worked? Let’s check the types at each step:

# Check types at each step
step1 = df.groupby('quality')['alcohol']
print(f"After groupby: {type(step1)}")

step2 = step1.mean()
print(f"After mean: {type(step2)}")

step3 = step2.reset_index()
print(f"After reset_index: {type(step3)}")

Output:

After groupby: <class 'pandas.core.groupby.SeriesGroupBy'>
After mean: <class 'pandas.core.series.Series'>
After reset_index: <class 'pandas.core.frame.DataFrame'>

Now here’s the mismatch: if you expected .mean() to return a DataFrame, you’d be wrong. It returns a Series. The .reset_index() box converts it back to a DataFrame. If you forgot that box, the next step would crash.

Let’s see what happens when we forget the reset_index:

# Forgetting reset_index — this will crash
try:
    result = (
        df
        .groupby('quality')['alcohol']
        .mean()  # Returns Series
        .to_frame()  # Convert Series to DataFrame — this is the missing box!
    )
    print("Success:")
    print(result.head())
except Exception as e:
    print(f"Error: {e}")

Output:

Success:
              alcohol
quality             
3        10.210000
4        10.228261
5         9.899706
6        10.445192
7        11.466250

Every time your code crashes with a type error, it’s telling you there’s a missing box in your diagram. The fix is always the same: add a transformation that bridges the type gap.

Step 5: Write One Function Per Box (Signatures = Arrow Types)

Now translate each box into a function. The function’s signature (input type → output type) is exactly the arrow types it bridges. This is the Unix philosophy applied to data science.

McIlroy’s Unix philosophy, from 1978, says: “Write programs that do one thing and do it well. Write programs to work together.” Each function should have a single responsibility — one transformation, one box.

Use type hints to make the signature explicit:

from typing import List, Tuple
from sklearn.pipeline import Pipeline
from sklearn.preprocessing import StandardScaler
from sklearn.ensemble import RandomForestRegressor

def clean_missing(df: pd.DataFrame) -> pd.DataFrame:
    """Remove rows with missing target values."""
    return df.dropna(subset=['quality'])

def engineer_features(df: pd.DataFrame) -> pd.DataFrame:
    """Create ratio features."""
    df = df.copy()
    df['acid_ratio'] = df['fixed acidity'] / df['volatile acidity']
    return df

def select_features(df: pd.DataFrame) -> pd.DataFrame:
    """Select only numeric columns for modeling."""
    return df.select_dtypes(include=[np.number])

# Build the pipeline
pipeline = Pipeline([
    ('clean', FunctionTransformer(clean_missing)),
    ('engineer', FunctionTransformer(engineer_features)),
    ('select', FunctionTransformer(select_features)),
    ('scale', StandardScaler()),
    ('model', RandomForestRegressor(n_estimators=100, random_state=42))
])

print("Pipeline steps:")
for name, step in pipeline.steps:
    print(f"  {name}: {step}")

Output:

Pipeline steps:
  clean: FunctionTransformer(func=clean_missing)
  engineer: FunctionTransformer(func=engineer_features)
  select: FunctionTransformer(func=select_features)
  scale: StandardScaler()
  model: RandomForestRegressor(n_estimators=100, random_state=42)

Each step is a box. The arrow between them is the DataFrame. The type hints in each function define what the arrow carries.

A well-designed pipeline is just a collection of tiny, testable functions, each doing one thing and passing data to the next.

Modern Equivalents: You’re Already Doing Data-Flow Decomposition

The framework isn’t new — it’s just naming what every modern ML/DS tool already does. Let’s see how different tools implement the same idea:

pandas Method Chains

Each method is a box, the DataFrame is the arrow. You’ve been doing this all along.

sklearn Pipeline

The steps attribute is a list of boxes. fit_transform passes the arrow through each box.

Polars LazyFrame

Polars builds a query plan — a DFD — before executing. You can inspect it with .explain():

import polars as pl

# Create a LazyFrame
lazy_df = pl.DataFrame({
    'group': ['A', 'B', 'A', 'B'],
    'value': [1, 2, 3, 4]
}).lazy()

# Build a pipeline
query = (
    lazy_df
    .groupby('group')
    .agg(pl.col('value').mean())
    .sort('group')
)

print("Query plan:")
print(query.explain())

Output:

Query plan:
SORT BY [col("group")]
  AGGREGATE
    	[] BY [col("group")] TO [col("value").mean()]
      DF []; PROJECT */2 COLUMNS; SELECTION: None

Each line is a box. The indentation shows the data flow.

Dask Task Graphs

Dask builds a directed acyclic graph (DAG) of transformations — pure data-flow:

import dask.dataframe as dd

# Create a Dask DataFrame
ddf = dd.from_pandas(df, npartitions=2)

# Build a pipeline
result = (
    ddf
    .groupby('quality')['alcohol']
    .mean()
    .compute()
)

print("Dask task graph (simplified):")
print(ddf.groupby('quality')['alcohol'].mean().visualize())

PyTorch DataLoader

A DataLoader is a pipeline of transforms (resize, normalize, to tensor) — each is a box:

from torchvision import transforms

# Define a pipeline of transforms
transform_pipeline = transforms.Compose([
    transforms.Resize((224, 224)),
    transforms.ToTensor(),
    transforms.Normalize(mean=[0.485, 0.456, 0.406], std=[0.229, 0.224, 0.225])
])

print("Transform pipeline:")
for t in transform_pipeline.transforms:
    print(f"  {t}")

Every tool you already use was designed around this idea. You just didn’t have the vocabulary to see it.

The Catch: Stateful Components Don’t Fit Cleanly

Data-flow decomposition assumes pure transformations — same input always gives same output. But some real-world components are stateful: accumulators, monitors, counters.

Stateful components remember something from previous calls. For example, a streaming data pipeline that counts total rows processed — the counter is state, not a pure transformation.

Let’s see a stateful counter:

# Stateful counter using a global variable
row_count = 0

def process_row(row):
    global row_count
    row_count += 1
    return row * 2

# Test it
for i in range(5):
    result = process_row(i)
    print(f"Row {i}: result={result}, count={row_count}")

Output:

Row 0: result=0, count=1
Row 1: result=2, count=2
Row 2: result=4, count=3
Row 3: result=6, count=4
Row 4: result=8, count=5

The counter changes each call. This is hard to test because the function’s behavior depends on hidden state.

Now refactor it into a pure function:

# Pure function — takes counter as explicit input
def process_row_pure(row, counter):
    return row * 2, counter + 1

# Test it
counter = 0
for i in range(5):
    result, counter = process_row_pure(i, counter)
    print(f"Row {i}: result={result}, count={counter}")

Output:

Row 0: result=0, count=1
Row 1: result=2, count=2
Row 2: result=4, count=3
Row 3: result=6, count=4
Row 4: result=8, count=5

Now the function is testable: same inputs always give same outputs. The state is explicit.

Not everything fits neatly into a box-and-arrow diagram. But the parts that don’t are usually the parts that cause bugs. That’s why functional programming advocates for pure functions — they’re easier to test and reason about.

Closing: What You Learned and Where to Go Next

Let’s recap the five-step framework:

  1. Name every input — inventory every piece of data that enters your pipeline, with type and shape.
  2. Name every output — define exactly what you need at the end, with type and shape.
  3. Draw the diagram — map each transformation to a box, each data flow to an arrow.
  4. Label every arrow with its type — ensure the output type of one box matches the input type of the next.
  5. Write one function per box — each function has a single responsibility and explicit type hints.

Key insight: You’ve been doing data-flow decomposition all along. Every pandas chain, every sklearn Pipeline, every Polars query — they’re all data-flow diagrams. Now you have the vocabulary to do it intentionally.

Here’s a checklist you can use for your next project:

# DATA-FLOW DECOMPOSITION CHECKLIST
print("""
1. INPUTS:
   - [ ] List every data source (CSV, API, config)
   - [ ] For each, specify type and shape
   - [ ] Example: transactions.csv -> pd.DataFrame, (10000, 5)

2. OUTPUTS:
   - [ ] Define the final deliverable
   - [ ] Specify type and shape
   - [ ] Example: churn_predictions -> pd.DataFrame, (5000, 2)

3. TRANSFORMATIONS:
   - [ ] Draw each box (one per transformation)
   - [ ] Label each arrow with type/shape
   - [ ] Check for type mismatches

4. FUNCTIONS:
   - [ ] Write one function per box
   - [ ] Add type hints to every function
   - [ ] Test each function in isolation
""")

In the next part of this series, we’ll use this framework to design a production-ready ML pipeline from scratch. You’ll see how it holds up under real-world pressure — with messy data, changing requirements, and tight deadlines.

For now, try drawing a data-flow diagram for your current project before you write another line of code. You might be surprised at what you discover.

Check Your Understanding

Remember: What are the five steps of data-flow decomposition?

Understand: Why is labeling every arrow with its type the hardest part? Give an example of a type mismatch in a pandas pipeline.

Apply: Take a pandas pipeline you’ve written recently. Draw its data-flow diagram with boxes and arrows. Label each arrow with the type and shape of the data at that point.

Analyze: Compare the pandas method chain in the opening example with the sklearn Pipeline version. What are the trade-offs of each approach in terms of readability, testability, and maintainability?

Evaluate: A colleague argues that drawing data-flow diagrams is a waste of time — “just write the code and debug it.” Write a counter-argument using the type-mismatch example from Step 4. How much time does a diagram save compared to debugging a silent type error?

Create: Design a data-flow diagram for a real-world problem you’re working on (or a hypothetical one). Include at least five transformations, three of which involve type changes (e.g., DataFrame → Series → DataFrame). Write the corresponding functions with type hints.

  • Part 6: Test-First Design — The previous article in this series showed how to write assertions before functions. Data-flow decomposition gives you the structure to write those assertions against.
  • Part 4: Stepwise Refinement — Wirth’s method of starting with a high-level algorithm and refining it step by step pairs naturally with drawing boxes and arrows before writing code.

Apply What You Learned is for Supporter and Insider subscribers.

Subscribe to unlock the exercises on this post.

See plans

Looking for something else?

Search every article by title, summary or topic.