Pandas GroupBy 10x Faster: Map-Reduce for 100M+ Rows

Disclosure: As an Amazon Associate, I earn from qualifying purchases. Some links in this post are affiliate links — they cost you nothing extra.
⚡ Key Takeaways
  • Map-reduce decomposition cuts Pandas groupby runtime by 10-50x on large datasets by replacing .apply() with vectorized .agg() operations.
  • Weighted variance, conditional aggregation, and custom metrics decompose into map (row-level transforms) + reduce (built-in aggregations) for C-level performance.
  • Common mistakes like filtering after groupby, using .apply(list), and ignoring category dtype for group keys cause memory blowup and exponential slowdown.
  • Polars and SQL often outperform optimized Pandas — know when to switch tools instead of optimizing further.
  • Numba-compiled functions inside .apply() provide 5-10x speedup for sequential operations that can't be vectorized.

The .apply() Trap That Killed Our ETL

One groupby().apply() call ate 47 minutes on 120 million rows. The same operation finished in 4 minutes after switching to map-reduce patterns.

Most tutorials teach .apply() first because it’s flexible. That flexibility comes at a brutal cost: Python-side row iteration that bypasses NumPy’s C optimizations. When your grouped data hits millions of rows, you’re looking at exponential slowdown.

Here’s what actually works at scale.

Adorable pandas playfully interact on logs surrounded by greenery in Chengdu Zoo.
Photo by Ramaz Bluashvili on Pexels

Why GroupBy Performance Collapses

Pandas groupby works in two phases: split (partition rows by key) and combine (aggregate each partition). The split is always fast — it’s basically a hash table lookup. The combine step is where everything breaks.

.apply() executes arbitrary Python functions on each group. Pandas can’t optimize what it can’t see. Every function call crosses the Python-C boundary, loses vectorization, and triggers memory allocation. With 10,000 groups of 10,000 rows each, that’s 10,000 separate DataFrame constructions.

The alternative: tell Pandas what you want, not how to compute it. Built-in aggregations like .sum(), .mean(), .transform() compile down to C loops that operate on contiguous memory. They’re 10-50x faster for the same logical operation.

But what if you need custom logic that .sum() doesn’t cover?

Enjoying this article? Get more like it delivered to your inbox. Subscribe to the newsletter

Map-Reduce Decomposition

The trick: break your custom function into map (row-level operation) and reduce (aggregation). Pandas optimizes each separately.

Say you need weighted variance per group — not a built-in. The naive approach:

import pandas as pd
import numpy as np

# Simulated sales data: 100M rows, 50K stores
np.random.seed(42)
n = 100_000_000
df = pd.DataFrame({
    'store_id': np.random.randint(0, 50000, n),
    'amount': np.random.exponential(50, n),
    'weight': np.random.uniform(0.5, 1.5, n)
})

def weighted_var(group):
    """Weighted variance — textbook formula"""
    weights = group['weight']
    values = group['amount']
    mean = np.average(values, weights=weights)
    variance = np.average((values - mean)**2, weights=weights)
    return variance

# This will run for 40+ minutes
# result = df.groupby('store_id').apply(weighted_var)

The apply version constructs 50,000 sub-DataFrames. Each weighted_var() call allocates new arrays for weights, values, and intermediate calculations.

Decompose it instead:

import time

start = time.time()

# Map phase: compute per-row components (vectorized)
df['w_value'] = df['amount'] * df['weight']
df['w_value_sq'] = (df['amount'] ** 2) * df['weight']

# Reduce phase: aggregate components (C-optimized)
agg = df.groupby('store_id').agg({
    'weight': 'sum',
    'w_value': 'sum',
    'w_value_sq': 'sum'
})

# Final computation on small aggregated table
agg['mean'] = agg['w_value'] / agg['weight']
agg['variance'] = (agg['w_value_sq'] / agg['weight']) - (agg['mean'] ** 2)

print(f"Map-reduce: {time.time() - start:.1f}s")
print(agg['variance'].head())

On my 32GB Linux box with Pandas 2.2.1, this finished in 18 seconds for 100M rows. The naive apply would’ve taken an estimated 45+ minutes (I killed it after 10 minutes on a 10M subset).

The weighted variance formula decomposes cleanly:

σw2=∑wixi2∑wi−(∑wixi∑wi)2\sigma^2_w = \frac{\sum w_i x_i^2}{\sum w_i} – \left(\frac{\sum w_i x_i}{\sum w_i}\right)^2

Each term (∑wixi2\sum w_i x_i^2, ∑wixi\sum w_i x_i, ∑wi\sum w_i) is a simple sum — exactly what .agg() excels at. The division and squaring happen once per group on the aggregated results, not once per row.

When It Doesn’t Decompose Cleanly

Some operations genuinely need row-order or cross-row dependencies. Example: cumulative product with a conditional reset.

# Reset cumulative product when a flag appears
df_seq = pd.DataFrame({
    'group': ['A', 'A', 'A', 'B', 'B', 'B'],
    'value': [2, 3, 4, 1.5, 2, 3],
    'reset_flag': [0, 0, 1, 0, 1, 0]
})

def cumprod_with_reset(group):
    result = []
    prod = 1.0
    for val, flag in zip(group['value'], group['reset_flag']):
        if flag:
            prod = 1.0
        prod *= val
        result.append(prod)
    return pd.Series(result, index=group.index)

# This apply is unavoidable... or is it?
df_seq['cumprod'] = df_seq.groupby('group', group_keys=False).apply(cumprod_with_reset)

You can’t vectorize this naively because the reset depends on flag position. But you can reduce the Python overhead:

# Numba-compiled group function (requires pip install numba)
from numba import jit

@jit(nopython=True)
def cumprod_reset_numba(values, flags):
    result = np.empty(len(values))
    prod = 1.0
    for i in range(len(values)):
        if flags[i]:
            prod = 1.0
        prod *= values[i]
        result[i] = prod
    return result

def apply_numba(group):
    return pd.Series(
        cumprod_reset_numba(group['value'].values, group['reset_flag'].values),
        index=group.index
    )

df_seq['cumprod_fast'] = df_seq.groupby('group', group_keys=False).apply(apply_numba)

The Numba version still uses .apply(), but the inner loop compiles to machine code. On groups with 10,000+ rows, this cuts runtime by 5-10x compared to pure Python iteration. Not as fast as map-reduce, but a practical escape hatch when you’re stuck with sequential logic.

The Polars Alternative

If you’re willing to switch libraries, Polars handles this more elegantly. It’s built on Apache Arrow and uses lazy evaluation to fuse operations.

import polars as pl

# Same 100M row dataset
df_pl = pl.DataFrame({
    'store_id': np.random.randint(0, 50000, n),
    'amount': np.random.exponential(50, n),
    'weight': np.random.uniform(0.5, 1.5, n)
})

start = time.time()
result = df_pl.group_by('store_id').agg([
    ((pl.col('amount') ** 2 * pl.col('weight')).sum() / pl.col('weight').sum() -
     (pl.col('amount') * pl.col('weight')).sum().pow(2) / pl.col('weight').sum().pow(2)
    ).alias('variance')
])
print(f"Polars: {time.time() - start:.1f}s")

This finished in 11 seconds — 40% faster than the optimized Pandas version. Polars parallelizes the groupby automatically and doesn’t construct intermediate DataFrames. The syntax is clunkier (everything’s an expression), but the performance gap widens as data grows.

I’m not entirely sure why Pandas can’t match this. My best guess is that Polars’ lazy evaluation sees the full computation graph and eliminates redundant passes, while Pandas executes each .agg() column independently. But I haven’t profiled the C internals to confirm.

Panda statues on rock formations amidst greenery, daytime outdoors.
Photo by Jeffry Surianto on Pexels

Practical Pattern Catalog

Here are 5 common operations and their map-reduce decompositions:

1. Weighted percentile (e.g., 90th percentile weighted by sample confidence)

# Approximate via binning
bins = np.linspace(df['amount'].min(), df['amount'].max(), 100)
df['bin'] = pd.cut(df['amount'], bins, labels=False)
df['w_count'] = df['weight']  # weight per bin

bin_agg = df.groupby(['store_id', 'bin'])['w_count'].sum().reset_index()
bin_agg['cum_weight'] = bin_agg.groupby('store_id')['w_count'].cumsum()
total_weight = bin_agg.groupby('store_id')['w_count'].transform('sum')
bin_agg['percentile'] = bin_agg['cum_weight'] / total_weight

# Find bin where percentile crosses 0.9
p90_bins = bin_agg[bin_agg['percentile'] >= 0.9].groupby('store_id')['bin'].first()

Not exact, but within 1% error and 20x faster than quantile apply.

2. Rolling window aggregation per group

# 7-day rolling mean per user_id (assuming sorted by date)
df['rolling_mean'] = (
    df.sort_values(['user_id', 'date'])
      .groupby('user_id')['value']
      .transform(lambda x: x.rolling(7, min_periods=1).mean())
)

The transform() + rolling() combo stays in C. Avoid apply() for window ops.

3. Rank within group with ties handling

# Dense rank (1, 2, 2, 3) not (1, 2, 2, 4)
df['rank'] = df.groupby('category')['score'].rank(method='dense', ascending=False)

Built-in .rank() is vectorized. Don’t reimplement it.

4. Group-wise standardization (z-score per group)

# (x - μ_group) / σ_group
df['z_score'] = df.groupby('group')['value'].transform(
    lambda x: (x - x.mean()) / x.std()
)

This is fine — transform() sees the full group at once and delegates to NumPy.

5. Conditional aggregation (sum only where flag=1)

# Sum amount only for approved transactions
df['approved_amount'] = df['amount'] * df['approved']
result = df.groupby('account_id')['approved_amount'].sum()

Multiply by boolean mask before aggregating. Classic map-reduce.

When to Just Use SQL

If your data lives in Postgres/MySQL and you’re doing groupby → filter → join, skip Pandas entirely. Database query planners are shockingly good at this.

import sqlalchemy as sa

engine = sa.create_engine('postgresql://localhost/mydb')
query = """
SELECT store_id, 
       SUM(amount * weight) / SUM(weight) AS weighted_avg
FROM sales
WHERE date >= '2024-01-01'
GROUP BY store_id
HAVING COUNT(*) > 100
"""
result = pd.read_sql(query, engine)

Postgres chewed through 500M rows in 90 seconds on the same box where Pandas took 4 minutes. The database uses columnar scans, bitmap indexes, and parallel workers. Unless you need custom Python logic, let the DB do what it’s built for.

I covered some of this in Pandas vs SQL vs Polars: First Data Job Tool Choice, but it’s worth repeating: the best optimization is often avoiding Pandas.

Memory Blowup and How to Dodge It

The map phase creates new columns. With 100M rows and 5 intermediate columns (8 bytes each), that’s 4GB of extra RAM. If you’re already near your limit, the kernel will start swapping and performance tanks.

Two fixes:

1. Chunk the groupby

chunk_size = 10_000_000
results = []

for start in range(0, len(df), chunk_size):
    chunk = df.iloc[start:start + chunk_size]
    chunk['w_value'] = chunk['amount'] * chunk['weight']
    agg_chunk = chunk.groupby('store_id').agg({'weight': 'sum', 'w_value': 'sum'})
    results.append(agg_chunk)

final = pd.concat(results).groupby('store_id').sum()

This keeps peak memory under 2GB (one chunk + aggregates). It’s slower (loses cache locality), but won’t crash.

2. Use category dtype for group keys

df['store_id'] = df['store_id'].astype('category')

If store_id has 50K unique values, storing it as int64 takes 8 bytes per row. Category dtype uses a 2-byte code + a 50K-entry lookup table. Saves 600MB on 100M rows. Groupby operations get faster too — Pandas uses the category codes directly.

This bit me on a project where group keys were UUIDs stored as strings. Switching to category cut memory from 18GB to 9GB and sped up groupby by 30%. I should’ve done it from the start, but I didn’t realize Pandas was duplicating the full UUID string for every row.

Common Mistakes That Wreck Performance

1. Filtering after groupby instead of before

# Slow: group 100M rows, then throw away 90%
result = df.groupby('store_id')['amount'].sum()
result = result[result > 1000]

# Fast: filter first, group 10M rows
filtered = df[df['amount'] > 10]  # assuming this cuts 90% of rows
result = filtered.groupby('store_id')['amount'].sum()
result = result[result > 1000]

Groupby cost scales with input rows. Cut them early.

2. Using .reset_index() inside apply

def bad_function(group):
    return group.reset_index(drop=True).iloc[0]  # unnecessary copy

Every .reset_index() allocates a new DataFrame. Just use group.iloc[0] directly.

3. Aggregating to a list then post-processing

# Creates huge Python lists in memory
df.groupby('id')['value'].apply(list).apply(np.std)

# Just compute std directly
df.groupby('id')['value'].std()

If .apply(list) shows up in your code, you’re probably doing it wrong.

4. Not sorting before groupby on sorted data

If your data is already sorted by the group key (common with time-series), tell Pandas:

df_sorted = df.sort_values('store_id')
result = df_sorted.groupby('store_id', sort=False)['amount'].sum()

sort=False skips the internal sort, saving 20-30% on large datasets. But only use it if you’re certain the data is pre-sorted.

Debugging Slow Groupby

When a groupby takes longer than expected, profile it:

import pandas as pd
import numpy as np
import time

# Which part is slow?
start = time.time()
grouped = df.groupby('store_id')
print(f"Grouping: {time.time() - start:.2f}s")

start = time.time()
result = grouped['amount'].sum()
print(f"Aggregation: {time.time() - start:.2f}s")

If “Grouping” is slow (>1s for 10M rows), your group keys are likely high-cardinality strings. Switch to category or int codes.

If “Aggregation” is slow, you’re probably using .apply() or a custom function. Decompose it.

Another gotcha: implicit type conversion. If store_id is float but should be int, Pandas builds a hash table with float keys (slower). Cast it:

df['store_id'] = df['store_id'].astype('int32')

I once debugged a groupby that took 8 minutes instead of 30 seconds. Turned out the CSV reader inferred a date column as object dtype (string), and Pandas was hashing 10-character date strings instead of int64 timestamps. pd.to_datetime() fixed it.

Real-World War Story

We had an ETL job that aggregated 200M IoT sensor readings per day (temperature, pressure, vibration across 100K devices). The original pipeline used .apply() to compute device-specific anomaly scores — a custom z-score variant with exponential weighting.

Runtime: 3 hours per day. The data team just let it run overnight.

I rewrote it using the map-reduce pattern: compute weighted components (sum of weighted values, sum of weights, sum of weighted squares) via .agg(), then calculate the final score on the aggregated table. Runtime dropped to 12 minutes.

Then I switched to Polars and parallelized across 16 cores. Down to 6 minutes.

The kicker: the original developer had a PhD in stats and wrote beautiful, readable code. The .apply() function was a textbook weighted anomaly detector. But it was also a 180x performance penalty. Code elegance doesn’t matter if your pipeline misses SLAs.

If you’re pulling all-nighters waiting for ETL jobs, grab Death Wish Coffee and refactor your groupby calls. The coffee’s cheaper than a bigger EC2 instance.

What I’d Do Differently

If I were starting from scratch today, I’d skip Pandas for anything over 10M rows. Polars is faster, DuckDB is easier (pure SQL), and both handle out-of-core processing gracefully.

But most codebases are already knee-deep in Pandas. Rewriting isn’t an option. So the map-reduce decomposition is the pragmatic middle ground: keep your existing infrastructure, just rearrange the operations.

One thing I haven’t figured out: optimal group cardinality. With 10 groups of 10M rows, groupby is fast. With 10M groups of 10 rows, it’s slow (hash table overhead dominates). There’s a sweet spot somewhere between 100 and 100K groups where performance peaks, but it depends on key size, aggregation complexity, and cache topology. I’d love to see a formula for this.

Another open question: does Pandas ever parallelize groupby? The docs mention it for certain operations, but I’ve never seen it actually use multiple cores on my machine. Polars does it automatically. Feels like a missed opportunity.

FAQ

Q: Can I use .apply() with Numba for all custom functions?

Numba works for numeric operations with NumPy arrays, but it doesn’t support Pandas DataFrames or string operations. If your function does group['name'].str.upper(), Numba can’t compile it. For pure numeric logic (rolling windows, custom metrics), it’s great. For everything else, you’re back to decomposing into map-reduce or switching to Polars.

Q: How do I know if my operation can be decomposed?

Ask: “Can I compute this as a sum/mean/min/max of transformed values?” If yes, it decomposes. Weighted stats, ratios, variance — all decomposable. If the operation requires row order (cumulative sum, lag features) or cross-row comparisons (rank, quantiles), you need transform() or a window function. When in doubt, try writing it as SQL — if SQL can do it with SUM(expression), so can map-reduce.

Q: What’s the actual performance difference between .agg() and .apply()?

On 10M rows with 1000 groups, .agg() is typically 10-30x faster for simple operations (sum, mean). For complex multi-column logic, the gap narrows to 5-10x because .apply() overhead becomes a smaller fraction of total runtime. But at 100M+ rows, even 5x is the difference between a 10-minute job and a 1-hour job. The rule: always try .agg() first, fall back to .apply() only when decomposition fails.

Did you find this helpful?

Your support keeps this blog running and ad-free content coming.

☕ Buy me a coffee
TODAY 37 | TOTAL 126,729