Sarth Calhoun

Append-only pipelines · Part 3 of 5 · Machine

Current truth with window functions

Adapted from a piece I wrote in July 2024, about work I did starting in the fall of 2022.

The unique key

Because this is an append-only system, all logic about what data is current and valid is performed on the query side. Which raises the question: how do you know data is an update rather than a new row?

For this, each file format has a unique key, which is a combination of columns, for example station_id, date, and time. If an entry matches existing data on those columns, it should replace the old data.

Consider the following example. Let's say all you have in a certain table is this:

station_id date time readings
North Bridge 2022-01-01 19:00 2000
Harbour West 2022-01-03 12:00 1000

Suppose you receive an updated file with the following data:

station_id date time readings
Harbour West 2022-01-03 12:00 1100
Old Mill 2022-01-05 12:00 1900

Since the Harbour West entry in the new file matches on the unique key (station_id, date, time), the new file should be considered the source of truth, and your current data is:

station_id date time readings
North Bridge 2022-01-01 19:00 2000
Harbour West 2022-01-03 12:00 1100
Old Mill 2022-01-05 12:00 1900

If you aren't using a CRUD model that replaces that entry with updated data, but instead keep every version of everything in the database, how do you surface only the most current information?

Window functions

For this you use window functions. (If window functions are new to you, this page has the mental model.) They take the form of:

ROW_NUMBER() OVER (PARTITION BY <...unique keys...> ORDER BY file_date DESC) = 1

Here is an example of that for the external table from part 2:

SELECT * FROM analytics.ext_acme_readings
QUALIFY ROW_NUMBER() OVER (
  PARTITION BY
    site_name,
    station_id,
    report_date,
    report_time
  ORDER BY file_date DESC
) = 1

QUALIFY filters on the result of a window function, the way WHERE filters on columns.

PARTITION BY creates a subset of rows the window function acts upon. Each PARTITION in this context is treated as its own thing.

ROW_NUMBER() will assign a unique number starting from 1 to each row in the PARTITION.

The rows are numbered by file_date DESC, in other words, the row with the highest value for file_date gets a row number of 1, the next 2, and so forth.

And then = 1 means only return the entry with a row number of 1.

This only works if the column you order by can't tie within a partition. If two rows can have the same file_date, ROW_NUMBER() will pick one of them arbitrarily, so either make sure the value is unique or add more columns to the ORDER BY to break the tie.

In other words, take all the data and create subsets based on the unique key, then return the most recent from each subset.

One thing worth noting here is the overloading of the term PARTITION in this pipeline. When using window functions, PARTITION is the subset of data that will be compared (because in this case, these share a unique key and are different versions of the same entry). For these purposes you return one, most recent, result from each PARTITION.

When setting up external tables, PARTITION is a subset of the data that Snowflake uses for optimization, and corresponds to a file path, which in this case is generated based on the file date.

So in the above query, you select all the data in the external table for a certain supplier's data, you PARTITION it by the unique key, and then you ORDER BY file_date... which just so happens to be the external table PARTITION.

You order the contents of the PARTITION by PARTITION.

Views

The window function logic, plus any other cleanup, goes in views, such as acme_readings_current or weekly_by_supplier. In the view definitions you use joins to handle exceptions, such as overriding bad identifiers, and window functions to surface the most recent version of data. Downstream code consumes the views.

When faster access is required, the view can be materialized or the data copied into regular tables.

It helps to pick a naming convention and stick to it. For example: