Processing huge tables with pandas and multiprocessing

Golden lines flowing across a dark grid

When a table has hundreds of millions of rows, a single pd.read_sql call is a good way to crash your machine. The pattern below keeps memory flat and uses every core.

The idea

  1. Read the source in chunks, never all at once.
  2. Hand each chunk to a worker process that transforms it.
  3. Write results out in batches, so one slow insert never blocks the rest.

A minimal sketch

import pandas as pd
from multiprocessing import Pool

def transform(chunk: pd.DataFrame) -> pd.DataFrame:
    chunk["email"] = chunk["email"].str.lower().str.strip()
    return chunk

def run(engine, query, workers=8, chunk_size=200_000):
    chunks = pd.read_sql(query, engine, chunksize=chunk_size)
    with Pool(workers) as pool:
        for out in pool.imap(transform, chunks):
            out.to_sql("target_table", engine, if_exists="append", index=False)

imap streams results back in order, so you never hold the whole dataset in memory at once.

Things that bite

  • Pickling cost. Everything passed to a worker is serialised, so send chunks, not connections.
  • Too many workers. More processes than cores just adds overhead.
  • Chunk size matters. Too small and overhead dominates, too large and memory spikes. Measure it.

I’ll share more patterns like this in upcoming posts.

Keep reading