Processing huge tables with pandas and multiprocessing

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
- Read the source in chunks, never all at once.
- Hand each chunk to a worker process that transforms it.
- 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.
