Implementing efficient batch feature computation is a complex engineering problem. Rather than ask users to think about how to engineer pipelines for incremental computation and maximal data reuse, Zipline automates these optimizations into its batch computation engine.

The result of this is faster, cheaper, and more stable data pipelines that power:

  • Daily batch feature computation
  • Training data generation
  • Experimentation

Batch upload job

The role of the batch upload job in Zipline is to compute batch features as of some current snapshot date/time, and upload those to the online store for serving. These jobs need to run on a scheduled cadence to keep values fresh.

Note: here we’re specifically referring to batch uploads for batch features. Streaming features also have a batch upload job to seed streaming state and apply batch corrections, but that’s covered in a different post [link].

For example: `average transaction amount by zip-code in the last 5 days.`

Naive approach

The way teams generally implement this is with a SQL pipeline such as:

SELECT
	zip_code,
	avg(amount)
FROM
	data.transactions
WHERE
	ds BETWEEN '{{ ds }}' AND DATE_SUB('{{ ds }}', 5)
GROUP BY zip_code

And run it on a schedule.

This will scan and aggregate 5 days of data everyday, with 4 of those 5 days being duplicated work on each successive run.

Zipline approach

In Zipline, the user would write something like:

GroupBy(
  source=transactions,
  Aggregation(
    input_column="transaction_amt",
    operation=Operation.AVERAGE,
    windows=["5d"]
  ),
  key="zip_code"
)

Zipline optimizes daily batch computation by implementing the following functions within its engine for all supported operations:

  • lift: turn a raw event into an intermediate representation (IR) which can be aggregated and combined with other IRs. In the case of average, the IR is a (count, sum) pair.
  • combine (⊕ operator below): Combines two IRs into a single IR. For average this looks like (c1, s1) ⊕ (c2, s2) = (c1 +c2, s1 + s2)
  • finalize: convert IRs back into a final value. For average this is Finalize(c, s) = s / c

Note: Zipline enforces commutativity and associativity of all supported aggregations. This starts to veer into abstract algebra which is out of scope of this blog. For the purpose of this post, it’s sufficient to understand that by breaking an average into (count, sum) pairs, we can now produce aggregations across various time windows, then roll those windows up further to produce larger windowed aggregations with a compressed state.

Furthermore, Zipline is aware of which operations can be reversed (Abelian Groups in abstract algebra), such that one IR can be subtracted from another to produce a correct result.

Important: Zipline is handling the complexity of this for the user, who only needs to select the aggregation that they want within the API (as shown above). Everything below happens below the surface as engine optimization.

Zipline leverages the reversibility of the average operation to produce the following:

Optimized windowed pipeline for reversible operations
Optimized windowed pipeline for reversible operations

Every batch (daily in this example):

  • New raw events are aggregated up to a Daily IR (count, sum) pair per key (zip code in this case). For example, this might be (10, 2000) for zip code 10012 for 2026-05-01.
  • A new running total IR gets computed by taking yesterday’s running total and merging in today’s Daily IR. If yesterday’s running total for 10012 was (150, 45000), we would now have (160, 47000) as the total IR for this key.
  • Produce the snapshot for today by subtracting the running total for (today - window) from today. I.e. For a 5 day window, you’d take the total from 5 days ago, say it was (120, 41000), and subtract it from today, getting: (160 - 120, 47000 - 41000) = (40, 6000). Next zipline will invoke the finalize function to get the output value: 6000/40 = 150.

Note that not all operations can be reversed (for example min and max). And if we cannot subtract to get a window, then Zipline simply combines Daily IRs to get the final windowed value:

Optimized windowed pipeline for non-reversible operations
Optimized windowed pipeline for non-reversible operations

Every batch (daily in this example):

  • Computes the new Daily aggregate IR by reading only the new raw data since yesterday.
  • Combine and finalize the last W days of IRs into a single value per key, where W is window size, and upload.

While this approach needs to do slightly more work (scan and agg W days of IRs rather than 2 days), it’s still substantially better than rescanning and aggregating the full window range of raw data over the window.

Point-in-time-correct (PITC) backfills for model training

The above system works well for producing daily uploads, however training data generation with point-in-time accuracy is more complex because feature values need to be accurate as of timestamps provided by the user in a “spine” or “observation” source.

For example, if we’re backfilling real-time features that we intend to use to train a model that predicts fraud at login time, then we need to compute those features as of various login (user_id, timestamp) pairs that are provided by the spine.

Partial Aggregate Re-Use

Zipline leverages partial aggregates in training data generation (Backfill Jobs) for PITC data by combining raw events for the head and tail of the window, with compressed IRs for the middle.

For example, if a user requests 10 days of training data (today -> today-10) with a feature that has a 60 day window, Zipline will pull raw data from the source for (today -> today-10) to get accuracy as-of head timestamps, and (today-60 -> today-70) to get tail accuracy on windows for observation timestamps. Everything in the middle, i.e. (today-59 -> today-11) does not need raw events, and can leveraged cached IRs.

Semantic Hashing

Sometimes for a newly requested backfill for a given spine, Zipline doesn’t need to recompute feature values at all if those exact values have already been produced by another job. Zipline achieves this level of compute re-use by using a concept of semantic hashing (a hash that encodes the exact semantics of a feature computation, incorporating both the feature semantics and the relevant keys and timestamps from the spine).

Computed values are saved in cache tables with column-level semantic hashing, so any new job leverages existing computation as much as possible.

Further details on how this works within the engine will be explored more deeply in a future post.

Conclusion

Above we explored all the optimizations that Zipline leverages as a result of being able to convert operations into IRs. In the case of AVERAGE, the conversion to IR, as well as the implementation of merge and finalize functions is relatively simple. However, authoring pipelines that correctly produce, cache, and reuse this incremental state introduces excessive overhead to a user’s workflow.

Multiply this complexity across all operations, and it’s effectively prohibitive for users to implement these optimizations themselves, which is why teams often end up with the naive approach described above.

Zipline’s philosophy is that the compute engine should abstract away these complexities, allowing users to focus on building models and not optimizing data pipelines.

Get in touch
Contact us with any questions.
[f] Talk to founders