Postgres to DataFrame, Ten Years Later

Back in 2016 I wrote a short post on connecting to Postgres with Pandas and Jupyter. At the time I was collecting web data into a Postgres database on AWS, and the post was put together to help a team who knew their way around CSVs but had never pulled data straight out of a database.

I had a look at it recently and the code doesn’t run anymore. The chunked loading example relies on df.append(), which was removed in pandas 2.0. Plenty of other things have moved on in the past decade too: the datasets I work with are bigger, I’d no longer keep a password in a JSON file next to a notebook, and the “if it doesn’t fit in memory, go and find another tool” advice at the end has aged the most of all.

With Polars 2.0 just released, this felt like a good excuse to update the post with what I’d use today.

Everything here was run with Python 3.13, Polars 2.0.0, ConnectorX 0.4.6, SQLAlchemy 2.1.3, psycopg 3.3.6 and PostgreSQL 17. Some of it is specific to this release, so expect details to change in later versions.


The data

To keep the spirit of the original, I wanted some web data sitting in Postgres. Bluesky publishes everything happening on the network through Jetstream, a public WebSocket feed of JSON events that needs no API key. I wrote a small collector that stores every post, like, repost and follow in a single table:

CREATE TABLE events (
    seq        bigint PRIMARY KEY,
    did        text NOT NULL,       -- the account
    time       timestamptz NOT NULL,
    operation  text NOT NULL,       -- create, update or delete
    collection text NOT NULL,       -- app.bsky.feed.post, app.bsky.feed.like, ...
    rkey       text NOT NULL,       -- the record's key within the account
    record     jsonb                -- the record itself
);

Jetstream can replay recent history, so rather than leaving the collector running all morning I replayed six hours (03:48 to 09:48 UTC on 7 October) over three connections in parallel, which took about 15 minutes. That came to a little over 6.2 million events, taking up 3 GB in Postgres. Jetstream delivers events at least once, so the same event can turn up twice (especially when replaying), but each one has a unique sequence number, and making seq the primary key means duplicates are dropped on the way in.

Everything here is public, but it’s still people’s posts, so I’m only showing totals and won’t be quoting or naming anyone.

Credentials

The 2016 post saved connection details in a config.json. These days I keep them in environment variables (usually loaded from a .env file that is in .gitignore), using the standard Postgres names so that other tools like psql pick them up too:

PGHOST=localhost
PGPORT=5432
PGDATABASE=bluesky
PGUSER=postgres
PGPASSWORD=not-telling

read_database_uri takes a connection URI rather than a psycopg2 connection object, so the first step is building one. Any special characters in the username or password need escaping, which urllib.parse.quote handles. By default it leaves / alone, so safe="" is needed to escape everything:

import os
from urllib.parse import quote

user = quote(os.environ["PGUSER"], safe="")
password = quote(os.environ["PGPASSWORD"], safe="")

uri = (
    f"postgresql://{user}:{password}"
    f"@{os.environ['PGHOST']}:{os.environ.get('PGPORT', '5432')}"
    f"/{os.environ['PGDATABASE']}"
)

Reading a table

With the URI in hand, reading a query into a DataFrame is one line:

import polars as pl

counts = pl.read_database_uri(
    "SELECT collection, count(*) AS n FROM events GROUP BY collection",
    uri,
)
┌───────────────────────┬─────────┐
│ collection            ┆ n       │
│ ---                   ┆ ---     │
│ str                   ┆ i64     │
╞═══════════════════════╪═════════╡
│ app.bsky.feed.like    ┆ 4127331 │
│ app.bsky.feed.post    ┆ 1068509 │
│ app.bsky.feed.repost  ┆ 638855  │
│ app.bsky.graph.follow ┆ 428712  │
└───────────────────────┴─────────┘

These are counts of events rather than totals: they include deletes as well as creates, and a like in this window can be for a post written at any time. Counting creates only, there were about four likes for every new post during the six hours.

There’s no cursor or connection to manage. By default this uses ConnectorX, which reads results straight into Arrow memory. It can also split a query across several connections and read the partitions in parallel. That isn’t automatically faster, since each partition is a separate query, but it can help with bigger tables:

likes = pl.read_database_uri(
    "SELECT seq, did, time FROM events WHERE collection = 'app.bsky.feed.like'",
    uri,
    partition_on="seq",
    partition_num=8,
)

Tables that don’t fit in one go

The old post handled big tables by reading 5,000 rows at a time and appending each chunk to a dataframe. That stops being useful pretty quickly, because you still end up with the whole table in memory at the end.

The Polars equivalent is iter_batches, which yields a DataFrame per batch. Rather than gluing the batches back together, I write each one out to Parquet. This is also a good place to let Postgres pull the useful fields out of the JSON, so the Parquet files have proper columns:

from pathlib import Path
from sqlalchemy import create_engine

query = """
    SELECT
        seq, did, time, operation, collection, rkey,
        record->>'createdAt'        AS created_at,
        record->'langs'->>0         AS lang,
        record->>'text'             AS text,
        record->'subject'->>'uri'   AS subject_uri
    FROM events
"""

engine = create_engine(uri.replace("postgresql://", "postgresql+psycopg://"))
out = Path("data/events")
out.mkdir(parents=True, exist_ok=True)
for old in out.glob("part-*.parquet"):
    old.unlink()

with engine.connect() as conn:
    batches = pl.read_database(
        query,
        connection=conn.execution_options(stream_results=True),
        iter_batches=True,
        batch_size=500_000,
        schema_overrides={
            "created_at": pl.String,
            "lang": pl.String,
            "text": pl.String,
            "subject_uri": pl.String,
        },
    )
    for i, batch in enumerate(batches):
        batch = batch.with_columns(
            pl.col("time").dt.convert_time_zone("UTC"),
            pl.col("created_at").str.to_datetime(strict=False, time_zone="UTC"),
        )
        batch.write_parquet(out / f"part-{i:04d}.parquet")

The stream_results=True option matters: without it, SQLAlchemy and psycopg may buffer much more of the result set in memory, which undoes most of the benefit of batching.

The other two additions both came from errors on the first run, and both are Polars being strict about types rather than guessing:

  • schema_overrides: Polars works out each column’s type from the first 100 rows. The first 100 rows here were all likes, which have no language or text, so those columns came out as the Null type and the batch failed as soon as it reached a post in Japanese. Declaring the types up front fixes it.
  • convert_time_zone("UTC"): Postgres hands timestamps over labelled Etc/UTC, while I parsed created_at as UTC. They’re the same thing, but Polars won’t compare datetimes with different time zone labels, so I normalise them as they’re written.

The export took just under a minute, and the 6.2 million rows came to 236 MB of Parquet, compared with 3 GB in Postgres (which is storing the full JSON for every record). Once the batches are on disk the database isn’t needed for the analysis, and the Parquet files can be scanned lazily as one table:

events = pl.scan_parquet("data/events/*.parquet")

Querying with SQL

This is the part that wasn’t possible in 2016. If the team I wrote the original post for were most comfortable in SQL, they’d have had to do all their filtering and joining in Postgres before pulling anything into Python. Polars can run SQL directly against DataFrames and LazyFrames, and 2.0 has expanded how much SQL it understands considerably.

Frames get registered in a SQLContext under a table name. Each event has two timestamps: time, when Jetstream saw the event, and created_at, the timestamp inside the record itself, which is set by whoever created it. My first question was how often those disagree:

ctx = pl.SQLContext(events=pl.scan_parquet("data/events/*.parquet"))

ctx.execute("""
    SELECT
        COUNT(*) AS posts,
        SUM(CASE WHEN created_at < time - INTERVAL '1 day' THEN 1 ELSE 0 END) AS backdated,
        ROUND(AVG(CASE WHEN created_at < time - INTERVAL '1 day' THEN 100.0 ELSE 0 END), 1) AS pct
    FROM events
    WHERE collection = 'app.bsky.feed.post' AND operation = 'create'
""").collect()
┌─────────┬───────────┬──────┐
│ posts   ┆ backdated ┆ pct  │
│ ---     ┆ ---       ┆ ---  │
│ i64     ┆ i32       ┆ f64  │
╞═════════╪═══════════╪══════╡
│ 1029042 ┆ 409218    ┆ 39.8 │
└─────────┴───────────┴──────┘

Almost 40% of the posts that arrived during those six hours have a created_at more than a day before Jetstream saw them. Digging into the biggest accounts explained it: people moving their history onto Bluesky. One account brought across three years of conversation in a few hours, and a local news site imported two years of articles, over 200,000 posts on its own. These look like imports rather than live activity, and since old posts rarely pick up new likes they’d drag down any measure of engagement.

That shaped the main query, which looks at likes per post for each language. It only keeps posts whose created_at is no more than an hour before Jetstream saw them, which removes most of the backfilled history (an import that happened within an hour of the original post would still get through). It also only includes languages with more than 1,000 accounts in this filtered sample, because one automated account (posting every second and a half to no likes at all) was enough to put a language in the results on its own. The language names come from a small CSV of two-letter codes that never went near the database, and SPLIT_PART trims region tags like en-GB and pt-BR down to the base language so they still match:

ctx = pl.SQLContext(
    events=pl.scan_parquet("data/events/*.parquet"),
    languages=pl.scan_csv("data/languages.csv"),
)

query = ctx.execute("""
    WITH posts AS (
        SELECT 'at://' || did || '/app.bsky.feed.post/' || rkey AS uri, did, lang
        FROM events
        WHERE collection = 'app.bsky.feed.post' AND operation = 'create'
          AND created_at > time - INTERVAL '1 hour'
    ),
    likes AS (
        SELECT subject_uri AS uri, COUNT(*) AS likes
        FROM events
        WHERE collection = 'app.bsky.feed.like' AND operation = 'create'
        GROUP BY subject_uri
    )
    SELECT
        l.language,
        COUNT(DISTINCT p.did) AS accounts,
        COUNT(*) AS posts,
        SUM(COALESCE(k.likes, 0)) / COUNT(*) AS likes_per_post
    FROM posts p
    JOIN languages l ON SPLIT_PART(p.lang, '-', 1) = l.code
    LEFT JOIN likes k ON p.uri = k.uri
    GROUP BY l.language
    HAVING COUNT(DISTINCT p.did) > 1000
    ORDER BY likes_per_post DESC
""")

execute() returns a LazyFrame, so nothing has run yet. The SQL is translated into the same query plan that the Python expression API would produce, so it goes through the same optimiser. You can see what it plans to do with explain(), and the useful part is at the bottom of each branch:

print(query.explain())
Parquet SCAN [data/events/part-0000.parquet, ... 12 other sources]
PROJECT 7/10 COLUMNS
SELECTION: ((col("collection") == "app.bsky.feed.post") & (col("created_at") > col("time").dt.offset_by(["-3600s"]))) & (col("operation") == "create")
...
Parquet SCAN [data/events/part-0000.parquet, ... 12 other sources]
PROJECT 3/10 COLUMNS
SELECTION: (col("collection") == "app.bsky.feed.like") & (col("operation") == "create")

Both CTEs have had their WHERE clauses pushed down into the Parquet scan, and each only reads the columns it needs: three of the ten for the likes, so the post text is never loaded at all.

Before running anything expensive, collect_schema() checks the column names and output types without reading any data:

query.collect_schema()
Schema({'language': String, 'accounts': Int64, 'posts': Int64, 'likes_per_post': Float64})

Then collect() runs it, which took under a second:

result = query.collect()
┌────────────┬──────────┬────────┬────────────────┐
│ language   ┆ accounts ┆ posts  ┆ likes_per_post │
│ ---        ┆ ---      ┆ ---    ┆ ---            │
│ str        ┆ i64      ┆ i64    ┆ f64            │
╞════════════╪══════════╪════════╪════════════════╡
│ German     ┆ 9180     ┆ 31817  ┆ 3.436748       │
│ French     ┆ 5475     ┆ 18494  ┆ 2.332973       │
│ English    ┆ 91234    ┆ 274534 ┆ 2.184822       │
│ Italian    ┆ 1351     ┆ 4224   ┆ 1.946733       │
│ Spanish    ┆ 7531     ┆ 25314  ┆ 1.927234       │
│ Korean     ┆ 3500     ┆ 15467  ┆ 1.861318       │
│ Dutch      ┆ 2603     ┆ 8449   ┆ 1.5461         │
│ Japanese   ┆ 30021    ┆ 92194  ┆ 1.144825       │
│ Portuguese ┆ 2380     ┆ 6535   ┆ 1.062433       │
└────────────┴──────────┴────────┴────────────────┘

German posts had the highest rate in this sample, with over three likes per post, ahead of French and English. Japanese, the second most active language by some distance, was near the bottom, with only Portuguese lower.

That’s six hours of one network, rather than proof that German speakers are nicer to each other! The obvious problems:

  • It only counts likes that arrived within the same six hours, so posts from near the end of the window have had less time to collect them.
  • The time of day favours whichever countries were awake, and a handful of very popular accounts can move a language’s average.
  • It counts like events rather than final totals. A like that was later removed is still counted, although that only applied to 0.4% of likes in this window.

One thing to watch out for: in 2.0, joins and group-bys no longer guarantee the order of their output rows by default. If row order matters, ask for it explicitly: with ORDER BY in SQL (as above), or in the Python API with maintain_order=True on group_by and maintain_order="left" on join. The release notes suggest maintain_order=True for both, but join only accepts a string saying which side’s order to keep, and passing True raises a TypeError.

When it really doesn’t fit

The 2016 post finished by saying that if the data was too big to fit in memory, you’d need a different tool altogether. That’s the advice that has changed the most.

In Polars 2.0, collect() uses the streaming engine by default. It processes data in chunks rather than loading everything at once, so a query over a large Parquet dataset mostly needs memory for the parts it’s working on, not the whole table. For operations that do need to hold a lot at once, such as a full sort, Polars now spills to disk once memory use reaches around 80% of RAM. In the 2.0 release, spilling covers sorts and window functions, with joins and group-bys still to come.

Six hours of Bluesky fits comfortably in my laptop’s memory, so to see spilling happen I ran Polars in a container with its memory capped at 1.5 GB, and sorted the whole table by post text. Sorting needs every row in hand before it can return the first one, which makes it a good worst case:

(
    pl.scan_parquet("data/events/*.parquet")
    .sort("text", nulls_last=True)
    .sink_parquet("data/sorted.parquet")
)

sink_parquet writes the result straight to disk instead of collecting it.

With no memory limit, the sort took 19 seconds and peaked at around 2.6 GB. In the 1.5 GB container with default settings, it was killed for running out of memory. Polars did pick up the container’s limit (with POLARS_VERBOSE=1 it reports total memory: 1.500 GiB), but the default threshold left too little headroom for everything else in the process. The release notes do say the 80% figure “may need tuning”.

Setting the memory budget explicitly fixed it:

POLARS_OOC_MEMORY_BUDGET_MB=500 python sort.py

With a 500 MB budget the same sort finished in 37 seconds without going over the container’s limit, writing 1.5 GB of temporary files along the way. That’s about twice as slow as having enough memory, which seems a fair price for not falling over. POLARS_OOC_SPILL_DIR controls where those files go, and they’re cleaned up afterwards.

I couldn’t find these settings in the docs. I found them by searching the 2.0.0 library for POLARS_OOC_, then tested them. So don’t build anything around them - they’re implementation details and could disappear in the next release.

Still, the point stands: the 2016 answer to “it doesn’t fit in memory” was a different tool, and now it’s an environment variable.

Wrapping up

The steps are the same as they were ten years ago: keep the credentials somewhere sensible, get the data out of Postgres, and query it. The difference is how much of the work now happens in one place, and how much less I have to think about memory along the way.