Connection Pooling Is a Distributed Systems Problem

One of the first things we usually do when building a Java application that talks to a relational database is configure a connection pool. The reason is straightforward: opening a database connection is expensive, and an application runs a huge number of database operations over its lifetime. Paying that price every single time is not something we want to do. With PostgreSQL the cost is even more tangible than with other databases, because every connection is served by its own backend process on the server. Each one carries its own memory, its own share of the scheduler, and its own contribution to lock and snapshot bookkeeping, which is why max_connections starts to hurt long before the machine looks busy. So, instead of creating a new connection for every operation, we keep a collection of connections around and reuse them. In the Java ecosystem, HikariCP has become the default implementation of this pattern, and frameworks such as Spring Boot will configure it for us without asking us to think about the details.

The architecture is quite simple. The application borrows a connection from HikariCP, uses it to execute one or more statements, and returns it to the pool. HikariCP takes care of creating connections when necessary, keeping the pool within its configured limits, retiring old connections, and coordinating the threads that are waiting for one. From the application’s perspective, it is a very convenient abstraction: the database connection becomes just another managed resource.

The problem appears when we stop running one application instance and start running many. Suppose our HikariCP pool is configured with a maximum of 20 connections. With a single instance, the application can hold at most 20 database connections. If we deploy ten instances, there are now ten independent pools, each allowed to open up to 20 connections, so we are potentially at 200. If we are in an environment with auto-scaling, and the number of instances grows to, let’s say, 100, the same configuration can end up asking PostgreSQL for 2,000 connections, and nobody changed a single line of configuration to get there.

Nothing strange has happened inside HikariCP. Every instance is behaving exactly as we asked it to behave. The problem is that HikariCP only has local knowledge. Instance A knows how many connections it is using, instance B knows how many it is using, and so on, but there is no shared resource manager coordinating those decisions. The database, on the other hand, is a shared resource. It has a finite amount of CPU, memory, and I/O, and a finite practical limit on how much concurrent work it can do.

This creates a mismatch between two scaling models. Application infrastructure is usually designed to scale horizontally: if traffic increases, we add more instances. Database infrastructure does not scale that way, or at least not as cheaply. Adding another application instance can indirectly increase the pressure on the database for no other reason than that the new instance brings another connection pool with it.

The obvious response is to make the HikariCP pool smaller. Instead of allowing every instance to hold 20 connections, perhaps we configure five. Many teams go one step further and do the arithmetic explicitly: take the connection budget PostgreSQL can afford, divide it by the maximum number of replicas the autoscaler is allowed to create, and use the result as the pool size. It works, but it is fragile. The number has to be recalculated every time someone raises the autoscaler limit, adds a new service against the same database, or runs a deployment that briefly doubles the number of pods during a rolling update. More importantly, it doesn’t change the fundamental relationship between application scaling and database connections. With 100 instances and a pool of five, we are still potentially at 500 connections, and most of them will be sitting idle on instances that happen to be quiet while the busy instances queue for their five. The database connection capacity is still coupled to the number of application instances.

This is where an external connection pooler such as PgBouncer changes the architecture. Instead of every application instance opening its physical connections directly against PostgreSQL, we introduce another layer between the applications and the database.

There are now two different kinds of connections. The connections between the applications and PgBouncer are client connections, while the connections between PgBouncer and PostgreSQL are server connections, the real backend processes. PgBouncer keeps a pool of the latter and multiplexes activity from a much larger number of clients onto a smaller number of actual PostgreSQL connections.

This distinction matters because the number of application-side connections no longer has to match the number of backend processes PostgreSQL has to run. We might have hundreds of application instances and thousands of client connections while deliberately capping the connections PgBouncer opens to PostgreSQL at, let’s say, 50. The pooler becomes the point at which we impose a global limit on the database, rather than relying only on independent limits inside each application process.

It also explains why, contrary to what many people think, HikariCP and PgBouncer are not really competitors. They operate at different layers. HikariCP is a library running inside each JVM, and manages the connections available to that application instance. PgBouncer is an external process, and manages the connections between itself and the database. It is perfectly valid to have both, and it is a very natural setup when we want local connection management and a central limit at the same time: HikariCP stops an individual instance from opening an uncontrolled number of connections, while PgBouncer stops the fleet as a whole from consuming an uncontrolled number of backend processes.

There is a detail I skipped over in the previous paragraphs, and it is the one that decides whether any of this works: PgBouncer’s pool mode. PgBouncer can operate in three modes, and they give very different results.

In session mode, a client connection is assigned a server connection when it connects and keeps it until it disconnects. That is the most compatible mode, but if we put HikariCP in front of it, which opens its connections at startup and keeps them open for as long as it can, we have gained almost nothing. Every Hikari connection pins a PostgreSQL backend, and we are back to 100 instances times 20 connections, now with an extra hop in the middle.

In transaction mode, a server connection is assigned to a client only for the duration of a transaction and goes back to PgBouncer’s pool when the transaction commits or rolls back. This is the mode where the multiplexing really happens, and the reason it works so well is that most application connections are idle most of the time. A Hikari connection spends the bulk of its life waiting between transactions, or waiting on the application to do something with the results. In transaction mode, that idle time no longer costs a backend process.

There is also a statement mode, which releases the server connection after every single statement and therefore forbids multi-statement transactions. It has its uses, but for the typical Java service it is not the mode you want.

Transaction mode comes with a catch, though, and it is one that Java applications run into quite frequently. Because consecutive transactions from the same client connection can land on different server connections, anything that relies on session state stops being reliable. A SET executed outside a transaction, session-level advisory locks, LISTEN/NOTIFY, temporary tables that outlive a transaction, all of them can silently end up on a different backend than the one we expect. The classic trap is prepared statements. The PostgreSQL JDBC driver switches to named server-side prepared statements after a statement has been executed a few times (controlled by prepareThreshold, five by default), and historically that produced errors such as prepared statement "S_1" does not exist when the next execution landed on a different backend. For years the answer was to add prepareThreshold=0 to the JDBC URL and give up server-side prepared statements. Since PgBouncer 1.21, setting max_prepared_statements makes PgBouncer track protocol-level prepared statements itself and re-prepare them transparently on whatever server connection the client lands on, which removes most of the pain.

As a reference, a minimal setup could look something like this:

; pgbouncer.ini
[databases]
orders = host=postgres.internal port=5432 dbname=orders
[pgbouncer]
listen_port = 6432
pool_mode = transaction
max_client_conn = 5000 ; client connections PgBouncer will accept
default_pool_size = 40 ; server connections per database/user pair
reserve_pool_size = 5
query_wait_timeout = 30 ; seconds a client may wait for a server connection
server_lifetime = 3600
server_idle_timeout = 600
max_prepared_statements = 200 ; PgBouncer 1.21+

And on the application side:

spring:
datasource:
# With PgBouncer < 1.21 add ?prepareThreshold=0 to the URL
url: jdbc:postgresql://pgbouncer.internal:6432/orders
hikari:
maximum-pool-size: 20
connection-timeout: 5000
max-lifetime: 1800000
idle-timeout: 600000

Notice that default_pool_size is per database and user pair, not per PgBouncer instance as a whole, so every service that connects with its own credentials gets its own pool of server connections. The real ceiling on PostgreSQL is the sum of those pools.

There is a price for introducing the second layer. We now have two systems managing database connectivity, and their configuration and behaviour have to be considered together. Connection limits, timeouts, connection lifetimes, and failure behaviour exist at both layers, and they interact in ways that are not obvious at first. Take waiting as an example. A request thread in our service first waits for HikariCP to hand it a connection, bounded by connection-timeout. Once it has one, the first statement of the transaction may wait again inside PgBouncer for a server connection, bounded by query_wait_timeout. That is two queues, with two timeouts, surfacing two different errors, and whoever is on call at three in the morning needs to know which one fired. The same happens with lifetimes: Hikari’s max-lifetime now governs the cheap connection to PgBouncer, while server_lifetime governs the expensive one to PostgreSQL, and tuning one does nothing for the other.

That queue inside PgBouncer is also worth pausing on, because it is already a form of backpressure. When all server connections are busy, clients do not get new backends, they wait, and after query_wait_timeout they are rejected. That is admission control at the database boundary, and it is one of the main reasons a central pooler protects the database better than a hundred independent pools.

And there is one more detail that is easy to miss: the limit is only global if PgBouncer itself is central. PgBouncer is single-threaded, so at high throughput teams run several instances, and each one maintains its own pools. Five PgBouncer instances with default_pool_size = 40 means up to 200 server connections, not 40. Some teams also deploy PgBouncer as a sidecar in every application pod, which is convenient but brings us right back to the original multiplication problem, just one layer further down. Where the pooler runs is as much an architectural decision as whether we run one at all.

If you are on a managed PostgreSQL offering, a lot of this may already be available as a service. Amazon RDS Proxy, Supabase’s Supavisor, and the built-in PgBouncer in Azure Database for PostgreSQL are all variations of the same idea, and the same questions about pool modes, session state, and timeouts apply to them.

PgBouncer sits at the PostgreSQL protocol boundary. It does not need to know whether the application using it is written in Java, Go, or Python. As long as the client speaks the PostgreSQL protocol, PgBouncer can sit in the middle. That is one of the reasons this architecture works so well in organisations where many services and languages share the same PostgreSQL infrastructure.

Other PostgreSQL proxies take this idea further. PgCat, for example, combines connection pooling with health checking, load balancing, failover, and read/write splitting. PgDog, written by one of PgCat’s original authors, is a multi-threaded proxy that provides pooling and routing, including support for sharding. These are not alternative implementations of HikariCP. They are PostgreSQL-aware infrastructure components that sit at the database boundary and can make decisions based on the topology and behaviour of the PostgreSQL cluster.

Here, the proxy is no longer only a mechanism for reducing the number of database connections. It becomes part of the database topology, deciding where different types of traffic should go and reacting when one of the nodes becomes unavailable.

For a long time, these were usually the options people compared and combined when making this decision. Recently, though, Open J Proxy (OJP) reached version 1.0.0, its first production-ready release (JavaPro article).

OJP approaches the problem from a different direction. Rather than being a proxy for one database protocol, it is built around the Java/JDBC boundary. The application uses the OJP JDBC driver, which talks to an OJP server over gRPC, and the OJP server becomes the component responsible for the actual database connections. Because it sits at the JDBC level instead of the wire protocol, it is not tied to PostgreSQL; the same server can front PostgreSQL, MySQL, MariaDB, Oracle, SQL Server, DB2, and others with a JDBC driver, which is a real difference from everything else in this article. The flip side is that it only helps Java clients. A Python batch job or a Go service hitting the same database will go around it.

The OJP server uses HikariCP by default (DBCP is also available), so the underlying pooling mechanism has not disappeared; what has changed is where the pool lives. With a normal deployment, every application instance has its own HikariCP pool; with OJP, the physical pool is centralised in the OJP server. The OJP driver’s connections are virtual, and putting HikariCP on top of them would add a queue that knows nothing about the real limit, so it makes sense to remove the application-side pool and let the OJP server be the single place where connections are managed.

So OJP is not replacing HikariCP with another implementation of the same abstraction. It moves the physical connection pool out of the application instances and into a shared service. The application still talks to the database through JDBC, but the physical connection becomes an infrastructure concern rather than something each JVM manages on its own. In that sense, OJP and HikariCP are not mutually exclusive either; it is closer to OJP wrapping HikariCP than to OJP versus HikariCP.

This lets us look at database connectivity as an admission-control problem. Suppose the database can safely sustain a certain amount of concurrent work. It doesn’t really matter whether that work comes from ten application instances or one hundred; what matters is the total amount reaching the database, and a central component can enforce that limit across the whole fleet.

It is fair to say, though, that PgBouncer already gives us that global point, as we saw with its wait queue. So the interesting question is not whether OJP can limit concurrency (both can) but what it can do because it sits at the JDBC level rather than at the wire protocol. Because the server sees JDBC operations with the context of the datasource they come from, it can make decisions a protocol-level pooler has a harder time making. One example is its slow query segregation, which separates slow and fast operations into different lanes so that a handful of heavy reporting queries cannot occupy every connection and starve the quick transactional ones. Features like backpressure, concurrency limits, circuit breaking, and query monitoring are not different ways of pooling connections; they are mechanisms for controlling the relationship between an elastic application layer and a comparatively constrained database layer, and the layer we put them in determines how much they know.

OJP deserves the same critical look we gave PgBouncer, though. Every database call now makes an additional network hop, from the application to the OJP server over gRPC, before it even reaches the database, and that latency is paid on every statement, not only when connections are opened. The OJP server is also now on the critical path for every Java service using it, so it needs to be deployed with redundancy, scaled, monitored, and upgraded like any other piece of shared infrastructure. None of this is unique to OJP; a central PgBouncer has exactly the same concerns, but it is important not to forget that centralising the pool also means centralising a failure point.

The important change is not that OJP has a pool and PgBouncer has a pool. Both do. The question is what the component surrounding that pool knows about and what responsibilities it has. HikariCP manages a pool inside one Java process. PgBouncer manages PostgreSQL connections at the protocol boundary. PgCat and PgDog add topology and routing concerns. OJP puts the pool behind a Java-aware service and can provide application-level controls around database access.

Do we actually need several layers? There is no universal answer. HikariCP plus PgBouncer is a good fit when we want local pooling in each service and a central limit that works for any language. OJP provides a similar centralisation model for Java applications while adding controls at the JDBC boundary. A PostgreSQL-specific proxy such as PgCat or PgDog may be preferable when the central problem is topology, such as routing reads to replicas or distributing work across shards.

What I would avoid is stacking all of them blindly. It is technically possible, but probably not a good idea. For example, we could build something like Java -> OJP -> PgBouncer -> PostgreSQL, but now there are two components managing connection limits, queueing, timeouts, and connection lifetimes. There may be valid reasons for it, for example, if OJP provides application-level controls while a PostgreSQL proxy handles routing to replicas, but adding another pool does not automatically improve the system. Each additional layer is another place where requests can queue, another set of timeouts to understand, and another failure mode to reason about. The same applies to PgBouncer, PgCat, and PgDog among themselves. They occupy roughly the same boundary, so in most architectures we would choose the one whose capabilities match our problem rather than chaining them.

The deeper problem behind all of these technologies is not connection pooling. Connection pooling is the mechanism we started with because database connections are an expensive and finite resource. Once the application becomes a distributed system, however, we discover that the resource is shared by many independent processes, while the original pool was designed to manage only one.

That leads to a more general distributed-systems problem: how do we govern a shared, finite resource when the clients consuming it are elastic?

  • For a small Java application, the answer can remain entirely local: Application -> HikariCP -> Database.
  • As the fleet grows, we may introduce a PostgreSQL-level pooler: Application Fleet -> HikariCP pools -> PgBouncer -> Database.
  • If database topology becomes part of the problem, a PostgreSQL-aware proxy can take responsibility for routing and failover: Application Fleet -> PgCat / PgDog -> [Primary | Replicas].
  • And if the environment is predominantly Java, and we want database access itself to be a centrally managed capability, we can move the physical pool behind OJP: Application Fleet -> OJP (HikariCP) -> Database.

The architectural decision is not which connection pool to use. It is where we want the responsibility for database resource management to live. HikariCP places it inside each application instance. PgBouncer places it at the PostgreSQL protocol boundary. PgCat and PgDog extend that boundary into routing and topology management. OJP moves it into a shared, Java-aware infrastructure layer while still using HikariCP underneath.

Once we look at it this way, these technologies stop looking like competing connection pools and start looking like different answers to the same scaling problem: the application layer wants to scale independently, while the database remains a shared and finite resource.

The connection pool is simply the first place where that tension becomes visible.

Connection Pooling Is a Distributed Systems Problem

Cloud Security: Principles and Concepts

Time passes, and almost every single minute of the day there is a cybersecurity event. If you follow any news sources in the Infosec space, you’ll see multiple events happening. To see how visible and prominent simple events are, we can try to provision a server with a public IP, and the moment it’s provisioned, we’ll see automated bots and port scanners probing it, often within minutes or even seconds. Obviously, not all traffic will come from malicious attackers, some of it will be just continuous IP harvesting from security researchers, search engines like Shodan or Censys, but some of it will be malicious botnets constantly monitoring especially cloud providers IP ranges. And that is just the beginning. This example is not to scare anyone, it’s just to raise your awareness.

Despite the prominence, and the obvious potential damaging consequences an attack can have, there are still plenty of organisations where security teams have tight budgets, and security is treated as a second-level citizen that wastes money without return of investment (ROI).

To deal with that situation, given those restricted budgets, defenders and cybersecurity professionals need to focus on the most important and efficient parts to optimise the implementation of defence controls. As such, we are going to review a few principles and concepts that will help us to decide where to focus our resources.

Before we start creating our security program, and buying any tools or hiring another analyst, the first questions we should ask ourselves should always be “what are we protecting” and “what happens if we lose it”. While it sounds trivial and obvious, there is a surprising number of organisations running security programs without any idea of what the answers to those two questions are.

This is where asset inventory and data classification comes in play. And yes, I know that it’s not glamorous, but it is essential work. We cannot defend what you don’t know we have, and we cannot prioritise defences for assets we haven’t ranked by importance. For example, a defaced marketing page, and a breached customer database do not have the same importance, especially from a business continuity point of view. Without classification, they can end up competing for the same level of attention and resources.

Another thing that consumes resources unnecessarily is trying to pursue the latest threats. for example, the latest zero-day, a new ransomware, the latest nation-state APT technique, etc… All of them are curious, we should be aware of them, and, I’m not going to lie, they are fun, but probably they are not worth it part of our resources. We should be focus on risks and not just threats.

The Risk is a function of likelihood and impact, not just how alarming something sounds. A well-run risk register, even a simple one, forces the conversation from “this could happen” to “how likely is this to happen to us, and how bad would it actually be”. That shift alone reorients a lot of spending decisions away from fear and toward evidence.

Another thing that I see often in organisation is try to create and add more and more controls when they should, first, try to reduce their Attack Surface. The addition of controls is an attempt of fixing every gap and every crack we perceived with a new addition (e.g., a product, a methodology, a practice). But often the most cost-effective move is to remove something, not add something. For example, decommissioning unused servers, closing unnecessary ports, removing stale accounts, retiring shadow IT, and turning off default services nobody remembers enabling. Every unnecessary system, open port, or forgotten admin account is one more thing a defender has to monitor and one more thing an attacker can try. Shrinking that surface is often cheaper than defending it, and it directly reduces the noise our teams have to process.

One principle that it seems to be, finally after many years, in the spotlight is the Principle of Least Privilege. An unimaginable number of breaches turn into major incidents not because of the initial foothold, but because of what the attacker could reach afterward. Accounts with more access than needed, or flat networks are examples of misconfigurations that can turn a small compromise into a large one.

Least privilege isn’t a one-time project, it’s a discipline. Regularly reviewing who has access to what, removing permissions that accumulate over time (have you heard of “permission creep”), and segmenting networks so a breach in one area doesn’t cascade into another, all pay for themselves the first time they contain an incident that would otherwise have spread.

Another one that has been with us forever, and it feels that is never going away is patching and configuration management. Most successful intrusions don’t rely on rare or exotic zero-days. They rely on known vulnerabilities that were never patched, or misconfigurations that were never caught. A disciplined, prioritised patch management process, one that ranks patches by exploitability and exposure rather than trying to patch everything everywhere immediately, consistently delivers more risk reduction per dollar (euro or pound) than flashier initiatives. And yes, it’s boring, but from a return of investment point of view is difficult to match.

The same goes for configuration hardening, examples such as unchanged default credentials, overly permissive cloud storage buckets, loose rules on firewalls and VPCs, and unnecessary services running on servers are the kind of low-effort, high-impact fixes that we’ll never write an interesting article about, but can make a big difference in an incident report.

The next one is a tad contra intuitive. Organisation should invest in Detection and Response. A very naive approach is to think that we are fully protected, and that nothing can happen to us, but the sad truth is that no set of preventive controls is perfect, and treating prevention as the only goal leads to unpleasant surprises. Mature security programs plan for the moment prevention fails. This is why detection and response capabilities, logging, monitoring, alerting, and a tested incident response plan, deserve a seat at the budget table alongside preventive tools. An attacker who gets in but is detected and contained within hours causes a fundamentally different kind of damage than one who has weeks of undetected access.

The next one is the one that we can control the least, the human element. Have you heard that one that says that “the biggest cybersecurity risk is located between the chair and the monitor”? That is an immutable truth, unfortunately, for our systems to exists, they need to have users, and those users are prone to errors, intentional or not. In this way, technical controls matter, but a large share of incidents still start with a person (e.g., a phishing email, a reused password). Companies need to invest in Security Awareness training. And ideally, not consider it just a checkbox exercise to comply with this or that policy, but make it relevant, useful, and, if possible, engaging. When it’s practical, recurring, and tied to real examples rather than generic slideshows, it measurably reduces the number of incidents that reach our technical defences in the first place.

In parallel to the implementation of a training program, it is equally important to make easy and low-friction for employees to report suspicious activity without fear of blame. A culture where people flag a suspicious email quickly is worth more than a bunch of security tools.

One thing that we should be careful about, and never lose sight of it, it’s that our cybersecurity programs is not the reason why our organisation exist. The organisation is there to serve a purpose, to sustain a business, we are just there to protect business continuity and the organisation. As such we should always implement guardrails and use frameworks as a compass and not as a cage. A radical examples if we do not follow this principle will be to turn our computers and services off, we will have one of the most secure organisations, but it will be impossible to run the business.

Frameworks like the Critical Security Controls (CIS), NIST Cybersecurity Framework, or ISO 27001 exist because this prioritisation problem is universal. They won’t tell you exactly what your organisation needs, but they offer a tested starting point, and a way to benchmark maturity over time, so we’re not reinventing prioritisation logic from scratch or relying purely on gut feeling. If used well, a framework becomes a shared vocabulary between security teams and leadership, making it easier to justify budget requests in terms leadership already understands: maturity levels, coverage gaps, and risk reduction, rather than technical jargon.

And let’s never forget that this is not a one control fits all situation. Not every system needs the same level of protection. A public-facing payment system processing millions of pounds should not necessarily have the same controls as an internal application used by five people. This sounds obvious, but security programs sometimes end up applying controls uniformly because uniformity is easier to measure. The problem is that this can result in spending too much protecting things that don’t matter very much, while the truly critical assets receive no additional attention.

And finally, none of these principles work in isolation, and none of them can offer perfect security, because that doesn’t exist. What they offer instead is a way to make deliberate, defensible choices about where limited time and budget should go.

Security will never be “finished”, and budgets and resources will rarely feel like enough. There will always be another vulnerability, another threat, another product and another control that someone will tell us we absolutely need. The challenge is not to eliminate all of them. The challenge is to understand what matters to our organisation, understand the risks we face, and spend our limited resources where they can make the biggest difference.

Because ultimately, a good security program isn’t the one with the most controls. It’s the one that gives the organisation the best protection for the resources it can realistically afford.

Cloud Security: Principles and Concepts

CPU Cache Hierarchy

Let’s imagine that we can write a little program that is mathematically perfectly optimised, O(n) complexity, zero heap allocations, and beautiful logic, but once executed found that it is still crawling during a benchmark. This would be what would happen without one of the many CPU optimisations. While there are many other (e.g., instruction pipelines, branch prediction), we are going to be focusing on the memory side of the CPU, an its hierarchy. Let’s put on our thinking hat, grab a coffee, and start digging.

In modern computing, there is a massive performance gap between the speed of our processor, and the speed of our main memory (DRAM). A CPU can perform operations in less than a nanosecond, but a trip to DRAM can take upwards of 100 nanoseconds. If our processors had to wait for DRAM for every single instruction, they would spend 99% of their time sitting idle. This gap between how fast a core can compute and how fast memory can feed it is old enough to have a name, the memory wall, and it has only gotten wider over the decades since clock speeds and core counts grew a lot faster than DRAM latency ever did.

The cache hierarchy is the answer to that gap. Instead of one flat pool of memory, we get several small and blazingly fast ones right next to the core, then progressively larger and progressively slower the further out we go, until you finally hit DRAM. The bet the whole design makes is that programs tend to reuse the same data and the same neighbourhoods of data over and over again in a short window of time, and if that bet holds, most loads never need to leave the chip at all.

The design is bounded by the laws of physics: you can have memory that is extremely fast and small, or memory that is large and slow, but you cannot have both. Every level in the hierarchy is trading capacity for speed, and the tradeoff gets more extreme the closer you get to the core.

  • L1 Cache (Level 1): This is the fastest and smallest. It is usually split into two parts: L1i (instructions) and L1d (data). It sits directly inside the CPU core and operates at the same clock speed as the processor itself. It is tiny, often only 32KB to 64KB, but it’s the first line of defence against a memory stall. It is private to a single core.
  • L2 Cache (Level 2): The next tier up. It is larger than L1 (typically 256KB to 1MB) but slightly slower. It is still usually private per core (though some designs share it across a pair of cores), and it’s the safety net for working sets that overflow L1 but are still actively in use.
  • L3 Cache (Level 3): The L3 is the “big” cache. It is significantly larger (several MBs) and is typically shared across all cores on a single CPU die. While much slower than L1 or L2, it is still orders of magnitude faster than DRAM. It acts as the final staging area before a request must be sent out to the system bus.

There are two approaches that can be taken when implementing this hierarchy:

  • In an inclusive hierarchy, anything sitting in L1 is guaranteed to also have a copy in L2 and L3, which makes cross-core coherence checks cheap (a core can ask “is this line anywhere?” by only checking L3) at the cost of wasting capacity on duplicated data.
  • Exclusive hierarchies keep each line in exactly one level, trading a cheaper coherence check for a more complex eviction path when a line gets promoted from L2 into L1.

This coherence is the whole reason multi-core caching is hard. The moment two cores can each hold their own private copy of the same line (see definition below), you need a protocol (MESI and its many descendants) to make sure a write by one core is visible to, or invalidates, the copies held by everyone else.

None of the three levels move data one byte, or even one word, at a time. The unit of transfer between every level of the hierarchy, and between L3 and DRAM, is the cache line, almost universally 64 bytes on modern x86 and ARM designs. Ask for a single int at address 0x1000, and the hardware doesn’t fetch 4 bytes, it fetches the entire 64-byte-aligned block containing that address, and every level between DRAM and the core now holds a copy of all 64 bytes, not just the four you asked for.

This is a deliberate bet on spatial locality. If you touched byte N, there’s a good chance you’re about to touch byte N+1, N+8, or N+40, and if that data already rode along for free, subsequent accesses become hits instead of misses. It’s also why data layout has a real, measurable performance cost that has nothing to do with algorithmic complexity. Two structurally identical loops can have wildly different runtimes purely based on whether consecutive iterations touch the same cache line or hop to a new one every time.

The down side of this is false sharing. If two cores are writing to two different variables that happen to live on the same 64-byte line, the coherence protocol treats it as if they’re fighting over the same data, bouncing that line back and forth between cores’ private caches on every write, even though the two threads never touch each other’s variable. It’s one of the more counter-intuitive bugs to diagnose precisely because the source code looks perfectly correct and free of any shared state.

A cache hit means the line the core asked for is already present at that level, resolved in that level’s fixed latency, done. A cache miss means it isn’t there, and the request has to fall through to the next level out, paying that level’s latency on top of everything already spent. Misses aren’t surprising, and there is a well-known taxonomy for why a miss happened, usually called the four C’s:

  • Compulsory: the very first time a line is ever touched, there was never a chance for it to be cached already. Sometimes called a cold miss.
  • Capacity: the working set is bigger than the cache, so lines get evicted before you come back to reuse them, even with a perfect access pattern.
  • Conflict: caches aren’t fully associative in practice, they’re organised into sets, and two lines that happen to map to the same set can evict each other even when there’s plenty of free space elsewhere in the cache. Higher associativity (8-way, 12-way) exists specifically to make this less likely.
  • Coherence: unique to multi-core, a line you already had cached and hadn’t touched gets invalidated because another core wrote to it, forcing you to re-fetch data you technically never evicted yourself.

Miss rate alone doesn’t tell the whole story, what matters for actual runtime is the average memory access time, roughly hit_time + miss_rate * miss_penalty, and because the miss penalty compounds across levels (an L1 miss that’s also an L2 miss that’s also an L3 miss pays all three latencies plus DRAM), even a small miss rate at the outer levels can dominate the total time spent waiting on memory.

Theory is fine, but the effect is much easier to internalise by making the same core do the same amount of work in two different orders. The classic version of this is traversing a 2D array in row-major order versus column-major order, in a language like C where a 2D array is really laid out as one contiguous block of memory, row after row. To get a better picture of all of this, let’s do a little bit of coding to see it more clearly.

#include <stdio.h>
#include <stdlib.h>
#include <time.h>
#define N 4096
// Row-major traversal: for a fixed row, consecutive j's are
// consecutive addresses in memory, exactly what a 64-byte
// cache line contains. Every 16 int accesses (64 bytes / 4
// bytes per int) share the same line that already got pulled
// in by the first access.
long sum_row_major(int (*matrix)[N]) {
long total = 0;
for (int i = 0; i < N; i++)
for (int j = 0; j < N; j++)
total += matrix[i][j];
return total;
}
// Column-major traversal over the same row-major layout: each
// step of the inner loop jumps N*4 bytes ahead, 16KB with
// N=4096. That's almost certainly a new cache line, and often
// a new page, on every single access.
long sum_col_major(int (*matrix)[N]) {
long total = 0;
for (int j = 0; j < N; j++)
for (int i = 0; i < N; i++)
total += matrix[i][j];
return total;
}
int main(void) {
int (*matrix)[N] = malloc(sizeof(int[N][N]));
for (int i = 0; i < N; i++)
for (int j = 0; j < N; j++)
matrix[i][j] = i + j;
clock_t start = clock();
long a = sum_row_major(matrix);
printf("row-major: %ld, %.3fs\n", a,
(double)(clock() - start) / CLOCKS_PER_SEC);
start = clock();
long b = sum_col_major(matrix);
printf("col-major: %ld, %.3fs\n", b,
(double)(clock() - start) / CLOCKS_PER_SEC);
free(matrix);
return 0;
}

Both functions do exactly N*N additions, in a different order. On my machine, the row-major version consistently comes in several times faster than the column-major one. Nothing about the arithmetic changed, only the order in which the same 64-byte lines got reused versus discarded and re-fetched. To prevent the compiler to optimise under the hood, compile with -O1 or -O0 flags. The result after execution should look similar to the following.

row-major: 68702699520, 0.020s
col-major: 68702699520, 0.107s

It’s tempting to file this away as a funny trick for array traversal, but the same principle shows up anywhere data is scanned in a tight loop: iterating a std::vector of small structs is cache-friendly, chasing a linked list of heap-allocated nodes scattered across the address space is not, even when both hold the same logical data. It’s a big part of why “data-oriented design” and struct-of-arrays layouts keep coming up in performance-sensitive codebases, and why a hash map with open addressing can outperform one built on separately-allocated buckets, despite having ostensibly worse theoretical collision behaviour. The algorithmic complexity on the whiteboard didn’t change in any of these examples. What changed is how well the access pattern lines up with 64 bytes at a time and a handful of megabytes of on-chip memory.

As we can see, the concepts here aren’t difficult once you sit with them for a bit, cache lines, hits, misses, three tiers of shrinking latency. The hard part, as usual, is remembering they exist while you’re actually writing the loop.

CPU Cache Hierarchy

How Object-Storage-Native LSM-Trees Work Under the Hood

In the last few years, data stores have changed, and “a lot” feels like an understatement. Around 2013 RocksDB was shipped, and it assumed a POSIX filesystem underneath it, because back then that was just what a fast key-value engine ran on. A decade or so later, key-value engines like SlateDB, and analytical table formats like Iceberg or Delta Lake, are running LSM-style structures straight against S3, where that assumption doesn’t hold anymore. This post is about what actually has changed to make that work.

An object-storage-native LSM tree isn’t a fundamentally new database architecture, it’s the same Log-Structured Merge-tree RocksDB has shipped for over a decade, but with the traditional POSIX filesystem replaced by immutable objects and a manifest updated via compare-and-swap (CAS). By running directly on cloud storage like AWS S3 or Google Cloud Storage, it adapts the classical design around three core characteristics:

  • Separation of compute & storage: State is persisted in scalable, low-cost object stores rather than local NVMe drives.
  • Immutable file alignment: Because LSM trees naturally write data sequentially into immutable files (Static Sorted Tables – SSTs), they natively match object storage’s write-once, read-many design.
  • Cloud-optimised I/O: Compaction and read paths are optimised to handle object store latency, high GET/PUT bandwidth, and explicit API call costs.

RocksDB’s write path depends on three things that don’t really have anything to do with LSM trees, they’re just what a POSIX filesystem gives us for free: we can append a few bytes to an existing file, we can fsync those bytes and know they survived a crash, and we can atomically rename a temp file over a real one to make a change visible in one step. The WAL leans on the first two. The MANIFEST, RocksDB’s own record of “which SSTs currently exist”, leans on the third.

If any one of those is taken away, the system does not degrade gracefully, it just stop working. S3 takes away all three. There’s no append, a PUT replaces the whole object. There’s no fsync, durability is whatever the object store’s replication does behind the scenes, and it happens on a timescale of tens to hundreds of milliseconds instead of the low single digits a local NVMe fsync costs. And there’s no rename, only, since August 2024, a conditional PUT: “create this key, but only if it doesn’t already exist“.

S3 is limited by its design, we cannot make it faster, so what is the smallest change to an LSM tree’s write path that survives losing append, fsync, and rename, while leaving everything else, memtables, sorted runs, compaction, exactly as it was? To figure out the answer we need to inspect what each missing primitive was actually protecting.

Append and fsync existed to make a partial write durable, a handful of bytes, safely, before the file that holds them is complete. If we can’t do that cheaply anymore, the fix isn’t to find a workaround, it’s to stop needing partial durability at all: buffer writes in memory until we have a whole, complete, self-contained object worth writing, and pay one network round trip for the whole thing instead of one round trip per record. This is exactly what a memtable already is. Object storage doesn’t force a new component into the design here; it just makes the memtable’s flush threshold matter for latency in a way it never did locally.

Rename existed to make a set of files change atomically, so a reader never observes “half the new SSTs, half the old ones”. Once we can’t rename, the only way to keep that guarantee is to never let the set of files change in place at all: every SST, once written, is permanently immutable, and the only thing that ever changes is a single small pointer, a manifest, listing which immutable SSTs are currently live. Updating that pointer is now the one operation in the entire system that needs an atomic primitive, and it’s small and infrequent enough that a conditional PUT and a retry loop can carry it.

RocksDB’s SSTs were already immutable once flushed, that part isn’t new. What’s new is that immutability stops being an implementation detail and becomes a first-level citizen of the entire system. Locally, “immutable” mostly meant “compaction rewrites files instead of editing them“, a convenience for concurrency control inside one process. On object storage, immutability is the only reason concurrent readers and writers can share a table at all without coordinating. A reader holding an old manifest can keep reading old SSTs indefinitely, safely, even while a compaction job somewhere else is busy writing brand-new ones, because nothing the reader is looking at will ever be touched again. Nobody has to lock anything. Nobody has to tell the reader to wait. The old SSTs just sit there, unreferenced eventually, garbage, but never wrong.

This is the part worth a deep consideration, because it’s a genuine inversion of where durability lives. In RocksDB, the WAL is the thing standing between us and data loss, and the MANIFEST is comparatively an afterthought, a bookkeeping file rebuilt from the WAL if it ever gets confused. In an object-storage-native LSM, that hierarchy flips. The manifest, one small JSON or Avro object, updated by Compare-And-Swap (or Compare-And-Set, CAS), is the database. It’s the single point that defines “what does this table currently contain“, and every SST it doesn’t list, however durably it sits in S3, might as well not exist.

That’s why the commit path collapses to one operation: read the current manifest, compute the new one, try to write it at the next version number with put_if_absent. If someone else got there first, the write fails, not with data loss, with a clean, detectable rejection, and we retry against their version instead of ours. This is optimistic concurrency control, the same pattern MVCC databases have used internally for decades, except here the granularity is “the whole table’s file list” instead of “one row“, and the retry cost is a network round trip instead of a spinlock.

SlateDB applies this architecture directly to low-latency key-value workloads by flushing memtables and conditional-PUTing manifest updates directly to S3. Analytical table formats like Iceberg and Delta Lake aren’t point-lookup KV engines, but they apply this exact same paradigm to columnar datasets: raw data is stored in immutable Parquet objects, while state changes (like Merge-on-Read or Copy-on-Write updates) are committed by racing to swap a manifest pointer, either native in S3 or via an external catalog (e.g., Hive metastore, Glue, REST catalog). The underlying engine goals differ, but the storage mechanics are identical: never mutate in place, always write new immutable objects, and guard the active manifest with compare-and-swap.

The design reads cleanly on paper, but as always the best way to learn is hands-on. Below is a minimal engine, in memory buffering, immutable SSTs flushed as whole objects, a manifest committed by CAS-and-retry, and a compaction pass that merges and swaps atomically, against a fake object store that only exposes what S3 actually gives us: put, put_if_absent, get, list.

The example is going to be in Python for convenience using only built-in standard library modules. A simple python script.py should suffice to run it.

import bisect
import json
import time
from dataclasses import dataclass, field
class ConditionalWriteFailed(Exception):
"""The S3 analogue of a failed compare-and-swap on If-None-Match."""
class FakeObjectStore:
"""No append, no in-place edits. The only concurrency primitive is
'create this key, but only if it doesn't already exist.'"""
def __init__(self, put_latency_ms=80):
self._objects: dict[str, bytes] = {}
self.put_latency_ms = put_latency_ms # simulated network cost
def put(self, key: str, data: bytes) -> None:
time.sleep(self.put_latency_ms / 1000)
self._objects[key] = data
def put_if_absent(self, key: str, data: bytes) -> None:
time.sleep(self.put_latency_ms / 1000)
if key in self._objects:
raise ConditionalWriteFailed(key)
self._objects[key] = data
def get(self, key: str) -> bytes:
return self._objects[key]
def list(self, prefix: str) -> list[str]:
return sorted(k for k in self._objects if k.startswith(prefix))
@dataclass
class SSTable:
"""Written once, never touched again. A real SST carries a sparse
block index and a Bloom filter so a miss doesn't cost a full fetch.
This one is small enough that the whole thing is the index."""
sst_id: str
entries: list[tuple[str, str | None]] # (key, value); None = tombstone
def get(self, key: str) -> str | None:
i = bisect.bisect_left([k for k, _ in self.entries], key)
if i < len(self.entries) and self.entries[i][0] == key:
return self.entries[i][1]
return None
def to_bytes(self) -> bytes:
return json.dumps(self.entries).encode()
@classmethod
def from_bytes(cls, sst_id: str, data: bytes) -> "SSTable":
return cls(sst_id, [tuple(e) for e in json.loads(data)])
@dataclass
class Manifest:
"""The one thing in this whole system that ever changes. Everything
it doesn't list might as well not exist."""
version: int
sst_ids: list[str] = field(default_factory=list)
def to_bytes(self) -> bytes:
return json.dumps({"version": self.version, "sst_ids": self.sst_ids}).encode()
@classmethod
def from_bytes(cls, data: bytes) -> "Manifest":
d = json.loads(data)
return cls(d["version"], d["sst_ids"])
class ObjectStoreLSM:
def __init__(self, store: FakeObjectStore, table: str = "t1"):
self.store = store
self.table = table
self.memtable: dict[str, str | None] = {}
self.local_cache: dict[str, SSTable] = {}
self._cached_manifest: Manifest | None = None
self._ensure_manifest_exists()
def _manifest_key(self, version: int) -> str:
return f"{self.table}/manifest/{version:06d}.json"
def _ensure_manifest_exists(self):
if not self.store.list(f"{self.table}/manifest/"):
self.store.put_if_absent(self._manifest_key(0), Manifest(0, []).to_bytes())
def _current_manifest(self) -> Manifest:
"""Real systems cache this pointer and only re-fetch on a CAS
conflict, rather than paying LIST+GET on every read; that's what
the cache below is for. (They still need some way to notice a
*different* writer moved the pointer without telling this
process: a poll, a watch, or a version check on some other
operation. This toy has exactly one writer, so it never has to
solve that half of the problem.)"""
if self._cached_manifest is None:
latest_key = self.store.list(f"{self.table}/manifest/")[-1]
self._cached_manifest = Manifest.from_bytes(self.store.get(latest_key))
return self._cached_manifest
def _commit_append(self, new_sst_ids: list[str], retries: int = 5) -> Manifest:
"""For flush(): a new SST doesn't depend on anything else that
might land first, so on conflict it's always safe to replay it
on top of whatever the latest version turns out to be."""
for _ in range(retries):
current = self._current_manifest()
candidate = Manifest(current.version + 1, current.sst_ids + new_sst_ids)
try:
self.store.put_if_absent(
self._manifest_key(candidate.version), candidate.to_bytes()
)
self._cached_manifest = candidate
return candidate
except ConditionalWriteFailed:
self._cached_manifest = None # someone else landed that version -> re-read
raise RuntimeError("manifest commit did not converge. contention too high")
def _commit_replace(self, sst_ids: list[str], based_on: Manifest) -> Manifest | None:
"""For compact(): the merged output was computed from a specific
snapshot of SSTs (`based_on`). If someone else committed in the
meantime, that output may already be missing data. It can't be
patched by appending, only discarded. Returns None on conflict
so the caller redoes the merge from scratch, rather than risking
a manifest that silently drops or resurrects files."""
candidate = Manifest(based_on.version + 1, sst_ids)
try:
self.store.put_if_absent(
self._manifest_key(candidate.version), candidate.to_bytes()
)
self._cached_manifest = candidate
return candidate
except ConditionalWriteFailed:
self._cached_manifest = None
return None
def put(self, key: str, value: str | None) -> None:
self.memtable[key] = value # None = delete
def flush(self) -> None:
"""The entire durability cost of a batch: one PUT for the SST,
one CAS for the manifest. Compare that to a local WAL paying a
network-grade fsync on every single write (this is the whole
reason batching stopped being optional)."""
if not self.memtable:
return
sst_id = f"sst-{int(time.time() * 1_000_000)}"
sst = SSTable(sst_id, sorted(self.memtable.items()))
self.store.put(f"{self.table}/data/{sst_id}.json", sst.to_bytes())
self.local_cache[sst_id] = sst
self._commit_append([sst_id])
self.memtable.clear()
def get(self, key: str) -> str | None:
if key in self.memtable:
return self.memtable[key]
manifest = self._current_manifest()
for sst_id in reversed(manifest.sst_ids): # newest first
sst = self.local_cache.get(sst_id)
if sst is None:
data = self.store.get(f"{self.table}/data/{sst_id}.json")
sst = SSTable.from_bytes(sst_id, data)
self.local_cache[sst_id] = sst
if key in dict(sst.entries):
return sst.get(key)
return None
def compact(self, retries: int = 5) -> None:
"""Merge, write once, swap the manifest. A reader never sees a
half-merged table, only the manifest before this call, or after.
Note this can't reuse flush's retry strategy. Flush's new SST is
independent of whatever else lands first, so replaying it on top
of the latest version is safe. Compact's merged SST is a snapshot
of a specific set of inputs. If a conflicting write landed in
between, that merge might already be missing an SST's worth of
data, or about to make an already-superseded one look live again.
Patching the manifest instead of redoing the merge is how a
compaction can silently resurrect files it just made obsolete."""
for _ in range(retries):
manifest = self._current_manifest()
if len(manifest.sst_ids) < 2:
return
merged: dict[str, str | None] = {}
for sst_id in manifest.sst_ids: # oldest to newest, newer wins
sst = self.local_cache.get(sst_id) or SSTable.from_bytes(
sst_id, self.store.get(f"{self.table}/data/{sst_id}.json")
)
merged.update(dict(sst.entries))
new_id = f"sst-compacted-{int(time.time() * 1_000_000)}"
new_sst = SSTable(new_id, sorted(merged.items()))
self.store.put(f"{self.table}/data/{new_id}.json", new_sst.to_bytes())
self.local_cache[new_id] = new_sst
if self._commit_replace([new_id], based_on=manifest) is not None:
return
# someone else committed first. this merge is stale, redo it
del self.local_cache[new_id]
raise RuntimeError("compaction did not converge. contention too high")
# old SSTs are now garbage, unreferenced, but still sitting in
# S3 until something is confident no reader still needs them

Now let’s execute a simple example to see how it works. If everything goes as expected we will see the number ’43’ listed twice.

store = FakeObjectStore(put_latency_ms=20)
lsm = ObjectStoreLSM(store)
lsm.put("alice", "42")
lsm.put("bob", "17")
lsm.flush() # one PUT, one CAS (that's the whole commit)
lsm.put("alice", "43") # overwrite, still just sitting in memory
lsm.put("carol", "9")
lsm.flush() # manifest now references two SSTs
print(lsm.get("alice")) # "43" -> the newer SST wins
lsm.compact() # merge both, swap the manifest atomically
print(lsm.get("alice")) # still "43", now from one merged SST

Two methods carry the entire idea:

  • flush, which turns “durable” from a per-write cost into a per-batch one
  • the pair of commit strategies that replace fsync-then-rename with compare-and-swap-then-retry

“Retry on conflict” isn’t a single reusable pattern, it depends on whether the operation retrying is additive (safe to replay on top of whatever won), or a function of a specific snapshot (unsafe to replay, has to be redone). Sorted runs, tombstones, newest-wins reads, merge-based compaction, all of it was present in the RocksDB approach. The object-storage-native part is entirely contained in how visibility gets established, not in the data structure.

Once the write path is solved, what’s left is making the read path fast, and that’s where Arrow, Flight SQL, and multi-tier caching actually earn their place as answers to problems the manifest-and-immutable-objects design creates on the read side.

Immutable SSTs mean a reader fetches whole objects or byte ranges from S3 constantly, so whatever format those objects are stored in had better not cost us a deserialisation pass on every fetch. That’s what Arrow buys: a columnar layout specified exactly enough, down to the byte, that a process can operate on a block pulled straight off the wire without constructing row objects first. Flight extends that further, its wire format is the in-memory format, so shipping a batch of results to a client skips the usual serialise-deserialise-reserialise round trip entirely.

And because every SST is now a network fetch away instead of a disk seek away, the cache in front of it has to be shaped for the access pattern that actually dominates, range scans, not point lookups. A three-tier hierarchy, RAM block cache, local NVMe as a cache of raw S3 bytes, and S3 itself as the source of truth is the same shape a buffer pool always had. What changes is the eviction policy: a point-lookup cache scores blocks mostly by recency, because a miss costs roughly the same either way. A range-scan cache has to score by how expensive a given fetch was relative to its size, and prefetch ahead of the scan cursor, because S3 rewards large sequential GETs far more than it rewards many small ones. Get that scoring wrong and every cache miss becomes a synchronous network stall sitting directly on our p99, no amount of clean manifest design upstream saves us from that.

None of this is a new architectural idea, though. It’s engineering effort spent making the consequences of “the filesystem is gone” fast, once we have already accepted the one substitution that made the whole thing possible in the first place.

Note on the toy engine: It skips a real local WAL for the sub-flush window (something still has to survive a crash between “written to the memtable” and “SST flushed”), garbage collection of orphaned SSTs after compaction, and Bloom filters, all things a production SlateDB or Iceberg deployment can’t skip. None of them change the substitution this whole post has been trying to explain for: an object-storage-native LSM tree is a normal LSM tree with the filesystem replaced by immutable objects and a manifest CAS. Everything harder than that is just making that one idea fast enough to matter.

How Object-Storage-Native LSM-Trees Work Under the Hood

Incremental View Maintenance

If you have spent any real time with RisingWave or Materialize, you have probably had the same reaction most people do the first time a CREATE MATERIALIZED VIEW with three joins in it updates in single-digit milliseconds after an insert. It feels like a magic trick. While it seems like a futuristic idea, it is actually an old idea from database theory combined with a new one. The old concept is called incremental view maintenance, while the new one is called differential dataflow, both combined achieve “magic”, and that is were the brilliantness resides, in the combination of both of them.

This post is about that combination, not the SQL surface, not the operator syntax, but what lays underneath and powers the solution. How a streaming engine represents a changing table, how it decides what to recompute when a single row changes, and how it keeps that state alive across gigabytes of history without either falling over or re-scanning the world. Put on your thinking hat, and grab a coffee because, while I’ll try to exemplify everything in code, the theory is not trivial. If you are one of those software engineers who though that the algebra and statistics were useless in the profession, you are in for a ride. If you prefer to read the code before than the theory, scroll down and come back up later.

A traditional SQL engine treats a query as a pure function from tables to a result. Run it once, get an answer, discard all the intermediate work. If the underlying data changes, you either re-run the whole query, or you don’t notice it at all.

That’s fine for OLAP dashboards refreshed on a scheduled job (e.g., cron), but it is not a very viable solutions if you want to update a visualization a few milliseconds after and update has happened, especially when those updates involve one or more JOIN operations. Re-running the join from scratch on every write is not just slow, it probably means you are trying to solve the wrong problem with the wrong tool. Every update is dealing with the size of the whole dataset to account for a change to one row.

Incremental view maintenance flips the framing. Instead of asking “what is the result of this query”. it asks “given that the input changed by this much, how does the output change”. The query stops being a function evaluated once and becomes a standing computation, a graph of operators sitting between the base tables and the view, permanently subscribed to change.

In this view, everything is a stream of value, time, and multiplicity. The unit of work in a differential dataflow-style engine is not a row. It’s a triple: (v, t, m). That, expresses a value, a logical time, and a multiplicity (sometimes called a weight or a diff). An insert of row v at time t is (v, t, +1). A delete of that same row later is (v, t', -1). An update is not a special case at all, it is treated as a delete and an insert that happen that happen in batch.

This is the single concept that makes the idea work. A “table” is not a set of rows sitting in a heap file, it’s the running sum of every (v, t, m) triple ever seen for it, filtered down to whatever multiplicity is currently nonzero. A join is not an operator that scans two tables, it’s an operator that consumes changes to two tables and produces changes to their join, using the algebraic fact that join distributes over the delta: if A changes by ΔA and B changes by ΔB, then:

Δ(A ⋈ B) = ΔA ⋈ B_old + A_old ⋈ ΔB + ΔA ⋈ ΔB

Everything to the right of the equals sign only touches rows that either changed, or matched something that changed. You never re-derive the parts of A ⋈ B that were already sitting there, unaffected, from an hour ago. This is the same bilinearity trick that DBSP (the theory Feldera and, in spirit, RisingWave’s internals lean on) is built around, and it’s why an incrementally maintained three-way join over a billion-row table can still update in the time it takes to touch the handful of rows that actually moved. If you want to dive deeper into the formal theory, check out the paper DBSP: Automatic Incremental View Maintenance for Rich Query Languages white paper (here).

Knowing the algebra isn’t enough. ΔA ⋈ B_old still requires probing B‘s current state for every row in ΔA. If B‘s current state means “replay every historical delta and sum it up”, you’ve just moved the O(n) cost from the join to the probe. This is what an arrangement solves.

An arrangement is a keyed, indexed trace of a collection: for each key, the full history of (value, time, multiplicity) triples that key has ever seen, laid out so that “give me the consolidated state of this key as of time t” is a cheap, roughly logarithmic lookup instead of a linear scan. It’s the same job a B-tree index does in a traditional database, except the thing being indexed isn’t a static table, it’s a log of changes, and the index has to answer “as of which time” as a first-class question, not an afterthought.

This is also where timestamps as a lattice stops being theoretical and starts being important operationally. If our engine only ever ingests from a single ordered source, a plain integer or a wall-clock timestamp is enough as time only moves forward. The moment you have multiple sources (e.g., a Kafka topic per shard, a CDC stream per upstream table, a join across two independently-progressing inputs), “time” is no longer a single number two events can be compared on unambiguously. You need a partial order: each source has its own frontier, and the only thing you can say for certain is “everything below this frontier has been fully seen”. Differential dataflow represents this as a lattice: timestamps that can be compared, joined (least upper bound), and advanced independently per input precisely so that out-of-order arrival across sources doesn’t produce a nondeterministic or incorrect result. An arrangement’s “as of” query is really “as of this point in the lattice, not this point on a clock”.

The uncomfortable part of all this is that an arrangement is, in the worst case, the entire history of every key. Real workloads don’t let you keep that in RAM forever. This is the part of the internals that separates “toy differential dataflow demo” from “thing that survives a production heavy load”.

Two things to consider:

  • Compaction. Once no live query and no downstream operator can possibly ask “what did this look like at time t” for any t below some frontier, we don’t need to keep the individual deltas that led up to that frontier, we only need their sum. This is structurally identical to what an LSM tree does on compaction, and it’s not a coincidence: engines like Materialize and RisingWave lean directly on LSM-shaped storage (custom implementations, or increasingly embedded engines like RocksDB or the newer disaggregated designs like SlateDB) specifically because “append deltas, periodically merge and prune” is the same access pattern an LSM was built for.
  • Spilling to object storage. Once the working set of “hot” arrangement state exceeds memory, the tail of the trace, the part least likely to be probed again soon, gets pushed to S3-shaped storage, with only the recent, high-churn part of the arrangement kept hot. This is crucial for disaggregated storage architectures, and both RisingWave and Materialize have moved toward: compute nodes stay stateless-ish and cheap to scale, while the arrangement’s durable history lives in cheap object storage behind a caching layer, and only gets pulled back in when a rare late-arriving event or a cold key needs it.

The risk is exactly what you’d expect, if our compaction or our spill path can’t keep up with the delta rate, either memory grows unbounded or every join probe pays a network round trip. Most of the actual engineering effort in these systems’ internals goes into keeping that spill path invisible on the hot path.

But, this is enough theory for today, let’s try to build a smallest version of all of this to better understand the ideas and concepts. We are going to be building a self-contained differential-dataflow-shaped engine, as simple as possible, that implements exactly the pieces above: (value, time, diff) triples, an arrangement with a compaction routine, and an incrementally maintained equi-join between two collections. It deliberately simplifies logical time down to a plain long instead of a full lattice, because the goal is to make the “why is an incremental join hard” idea legible in one sitting, not to reimplement Materialize’s timestamp model.

The code is written using Java 25, and it maintains a live “which customer has which open orders” view as customers and orders are added, cancelled, or renamed, the joined report updates itself in real time instead of being recomputed from scratch on each change.

Java
import java.util.*;
import java.util.stream.*;
public final class MiniDataflow {
// --- Core model ---------------------------------------------------
/** A single change to a collection: `data` appeared/disappeared at
* `time` with weight `diff`. Positive = insert(s), negative =
* retraction(s); an update is just a retraction and an insert that
* happen to share a batch. */
record Update<D>(D data, long time, long diff) {}
/** One entry of an *input* batch for a keyed collection: this
* (key, value) pair changed by `diff` at `time`. */
record Change<K, V>(K key, V value, long time, long diff) {}
/** A generic pair for join output. */
record Pair<A, B>(A first, B second) {}
private record ConsolidateKey<D>(D data, long time) {}
/** Collapse a batch: sum diffs for identical (data, time) pairs and
* drop anything that nets out to zero. Without this, an insert
* followed by a delete in the same batch would linger forever. */
static <D> List<Update<D>> consolidate(List<Update<D>> batch) {
return batch.stream()
.collect(Collectors.groupingBy(
u -> new ConsolidateKey<>(u.data(), u.time()),
LinkedHashMap::new,
Collectors.summingLong(Update::diff)))
.entrySet().stream()
.filter(e -> e.getValue() != 0)
.map(e -> new Update<>(e.getKey().data(), e.getKey().time(), e.getValue()))
.toList();
}
// --- Arrangements ----------------------------------------------------
/** A keyed, indexed trace of a collection's full history — the thing
* a real engine persists as an LSM so it can spill to object
* storage. This is a HashMap of Lists: same idea, minus surviving
* a crash. */
static final class Arrangement<K, V> {
private record HistEntry<V>(V value, long time, long diff) {}
private final Map<K, List<HistEntry<V>>> index = new HashMap<>();
void insert(List<Change<K, V>> batch) {
batch.forEach(c -> index
.computeIfAbsent(c.key(), k -> new ArrayList<>())
.add(new HistEntry<>(c.value(), c.time(), c.diff())));
}
/** The "as-of" view for a key: every (value, multiplicity) pair
* as it stood at logical time `asOf`, already summed. This is
* what a join probes instead of scanning the base table. */
Map<V, Long> snapshot(K key, long asOf) {
return index.getOrDefault(key, List.of()).stream()
.filter(e -> e.time() <= asOf)
.collect(Collectors.groupingBy(HistEntry::value, Collectors.summingLong(HistEntry::diff)))
.entrySet().stream()
.filter(e -> e.getValue() != 0)
.collect(Collectors.toMap(Map.Entry::getKey, Map.Entry::getValue));
}
/** Fold history at or before `frontier` into a single running
* total per (key, value) — the in-memory stand-in for what an
* LSM's compaction buys you. */
void compact(long frontier) {
index.replaceAll((key, history) -> {
var byFrontier = history.stream()
.collect(Collectors.partitioningBy(e -> e.time() <= frontier));
var folded = byFrontier.get(true).stream()
.collect(Collectors.groupingBy(HistEntry::value, Collectors.summingLong(HistEntry::diff)))
.entrySet().stream()
.filter(e -> e.getValue() != 0)
.map(e -> new HistEntry<>(e.getKey(), frontier, e.getValue()));
return Stream.concat(folded, byFrontier.get(false).stream())
.collect(Collectors.toCollection(ArrayList::new));
});
}
}
// --- Incremental join --------------------------------------------------
/**
* Delta(A join B): join a delta batch against the *other* side's
* current arrangement. This is the standard bilinear rule behind
* incremental joins (DBSP calls it the bilinearity of join): process
* the left delta against the right's prior state, then the right
* delta against the left's now-updated state, and between the two
* passes you cover deltaA join B, A join deltaB, and deltaA join
* deltaB exactly once — without ever touching a row that didn't
* change.
*/
static <K, VA, VB> List<Update<Pair<VA, VB>>> joinDelta(
List<Change<K, VA>> delta, Arrangement<K, VB> otherSide) {
return delta.stream()
.flatMap(c -> otherSide.snapshot(c.key(), c.time()).entrySet().stream()
.map(match -> new Update<>(
new Pair<>(c.value(), match.getKey()),
c.time(),
c.diff() * match.getValue())))
.toList();
}
}

The full file, including the IncrementalJoinView orchestrator that wires a customers arrangement and an orders arrangement together into a maintained customers ⋈ orders view, plus the Order record and the driver is in my repo here, and it has been extensively comment to make it more understandable.

The driver feeds five rounds of changes through the view: two customers arrive, one places an order, a second customer places an order while the first’s is cancelled, the first customer is renamed, and a brand new customer places an order in the same round they’re created. Nowhere does the program re-scan customers or orders in full, every round only touches the arrangement entries for the keys that changed.

t=0: two customers arrive, no orders yet
t=1: Alice places order #101 for 42
[DELTA t=1] +1 row=(Alice, order=(101, 42))
t=2: Bob places order #102 for 17, Alice's #101 is cancelled
[DELTA t=2] +1 row=(Bob, order=(102, 17))
[DELTA t=2] -1 row=(Alice, order=(101, 42))
t=3: Alice's account is renamed (retract old name, insert new)
t=4: a brand new customer places an order in the same round
[DELTA t=4] +1 row=(Devon, order=(103, 9))
Compacting history up to t=2 (simulating spilling old deltas)...
Final materialized view (customer -> order):
Bob order #102 amount=17
Devon order #103 amount=9

Two things worth mentioning:

  • t=3 produces no delta at all. Alice’s rename retracts and reinserts her customer row, but by t=3 she has no live orders, her one order was retracted at t=2. The join correctly sees that there’s nothing on the other side to combine with, and emits nothing. That’s not a special case in the code; it falls straight out of snapshot returning an empty map for a key with net-zero multiplicity. This is exactly the kind of thing that’s easy to get wrong by hand and impossible to get wrong once the algebra is doing the work.
  • The final view has no row for Alice, and no error either. The retraction at t=2 and the (empty) rename at t=3 leave her with zero live rows in the join, which currentView correctly filters out by keeping only positive net multiplicities. Nothing needed to be told to “delete” her row, it just stopped having positive weight.

The full source (linked in the repository above) contains the orchestrator that ties the two arrangements together, plus the compaction call that simulates pushing old deltas out of hot memory.

To keep this at a readable size, the mini engine skips a few things a real one can’t:

  • A real timestamp lattice. Multiple independently-progressing sources need partial-order timestamps and per-source frontiers, not a single long.
  • Actual spill-to-storage. compact folds history in memory; it doesn’t push the folded tail out to a storage, or handle bringing it back in on a cold-key probe.
  • Non-equi joins, aggregations, and windowing. The bilinear join rule generalises, but range joins and windowed aggregates each need their own incremental operator with their own state layout.
  • Backpressure and batching policy. A real engine has to decide how large a batch to accumulate before committing a round of deltas. Too small and you thrash the arrangement, too large and you blow your latency budget.

As we can see, with a little bit of reading and thinking the concepts exposed are not difficult to understand, the really hard part is to get their implementation right under load.

Incremental View Maintenance

Compliance Is Becoming a Software Engineering Problem

I try to spend as much time as I can talking to software engineers – some I work with, others I meet through meetups or conferences. Over time, I’ve started to notice something curious.

Most are absolute experts at their craft. I am constantly amazed by the mountain of information and acronyms we carry in our minds without even realising it. But almost all of it revolves around things we find interesting, immediate problems we need to solve, or the “new shiny thing”. Rarely, when talking with engineers, do compliance or governance come up.

I totally get it. Compliance is universally perceived as less fun. But despite that reputation, regulatory shifts are happening right now that will directly affect how we build software, and we need to be aware of them.

Many engineers can elegantly explain Raft consensus, debate the merits of eBPF, or spend hours discussing the subtleties of eventual consistency. But mention a CVE, and while most engineers will recognise the term, usually because they encounter it through scanner reports, Dependabot alerts, Renovate PRs, or Jira tickets, go a step further and mention a CWE or a CRE, and familiarity drops dramatically even with experienced developers.

For many engineers, vulnerability identifiers belong to security teams, auditors, or compliance departments. They are things that appear in scanner reports, Jira tickets, Dependabot alerts, or Renovate PRs. They are somebody else’s problem.

That perception is becoming increasingly difficult to sustain.

Modern software development is inseparable from security. Every application is built atop layers of frameworks, libraries, containers, operating systems, and cloud services. Vulnerabilities emerge continuously throughout that supply chain. More importantly, governments and regulators have started treating vulnerability management not as a polite recommendation, but as a legal obligation.

In Europe, the Cyber Resilience Act (CRA) marks a massive shift in thinking. Beginning in September 2026, manufacturers of digital products face mandatory reporting obligations for actively exploited vulnerabilities and severe incidents. In the years that follow, broader security requirements will become legally enforceable, with penalties reaching millions of euros or a significant percentage of global turnover.

We can no longer live in a restricted world where our sole purpose is to resolve problems in clever, performant ways. Understanding vulnerabilities is no longer exclusively the domain of security specialists – it is becoming a core part of software engineering itself.

After years of reading literature around shift-left, DevOps, DevSecOps, and SRE, I used to assume everyone was on the same page. I’ve since realised that was just my personal bias as a cybersecurity hobbyist.

Before we can dive into the legislation that will soon affect us, we need to understand the vocabulary that underpins modern vulnerability management.

Common Vulnerabilities and Exposures (CVE)

Imagine trying to coordinate a fix for a defect without a common naming system. One security researcher publishes a blog post describing a flaw in OpenSSL, a cloud provider releases an advisory using a different title, vulnerability scanners invent their own proprietary identifiers, and operating system vendors create yet another naming scheme. Chaos follows.

The CVE program exists to prevent precisely that. A CVE identifier answers a simple question: “Which specific vulnerability are we talking about?”. For example, Heartbleed became CVE-2014-0160, and Log4Shell became CVE-2021-44228.

These identifiers create a shared language used by researchers, vendors, scanners, incident response teams, and governments. Once a vulnerability receives a CVE number, everyone can refer to exactly the same issue. However, CVEs describe symptoms, not causes. They tell us what is broken, but they don’t explain why.

Common Weakness Enumeration (CWE)

The CWE represents a category of software weakness rather than a specific instance of a bug. For example, CWE-89 refers to SQL Injection, or CWE-79 represents Cross-Site Scripting.

Individual CVEs always map, not without controversy sometimes, back to one or more CWEs. Log4Shell, for example, was ultimately traced back to unsafe lookup behaviour and improper handling of untrusted input. The specific vulnerability was unique, but the underlying engineering weakness had been known for decades.

Unfortunately, this is where the gap between vulnerability management and daily engineering lives. Software engineers rarely think in numeric identifiers; they think in architectural practices: input validation, output encoding, authentication, authorisation, and dependency management.

Common Requirements Enumeration (CRE)

The CRE attempts to close that exact gap. Where CVEs answer “what happened” and CWEs explain “why it happened”, CREs focus on the most practical question of all: “What should engineers do to prevent it from happening again?“

When you put the three frameworks together, they form a complete defensive picture:

  • A CVE describes an incident
  • A CWE describes the underlying weakness.
  • A CRE describes the engineering practices that reduce the probability of that weakness appearing in the first place.

Organisations trapped entirely at the CVE layer spend their lives reactively patching. Organisations that understand CWEs start eliminating recurring technical debt. But organisations that embrace CREs and secure engineering requirements begin preventing vulnerabilities before they ever exist – which is exactly what regulators are now expecting.

For decades, security failures were primarily treated as business risks. Companies suffered reputational damage, customers temporarily lost trust, and major breaches led to lawsuits or expensive remediation. But general regulatory intervention remained light.

That world is gone. The European Union’s Cyber Resilience Act is one of the most ambitious pieces of software legislation ever written. Its premise is straightforward: products containing digital components must be secure by design and maintained throughout their entire lifecycle.

Under this framework, non-compliance penalties can reach up to €15 million or 2.5% of global annual turnover. While these numbers inevitably grab the attention of executives and legal departments, compliance needs to be solved by software engineers:

  • Lawyers cannot produce a Software Bill of Materials (SBOM).
  • Finance departments cannot determine if an exposed container image contains a vulnerable dependency.
  • Executives cannot decide whether a newly reported exploit affects production infrastructure.

Only engineering organisations possess that knowledge. Because of that, engineers must understand the language used by scanners, advisories, regulators, and vulnerability databases.

Additionally, with the arrival of advanced AI tools, discovering vulnerabilities is no longer the hard part. Automated scanners can identify thousands of CVEs in minutes. The actual challenges moving forward are entirely context-driven:

  • Which vulnerabilities actually matter to our architecture?
  • Which systems are genuinely exposed?
  • Which structural weaknesses keep recurring in our codebase?
  • Which engineering practices need to change?
  • Which incidents legally require regulatory reporting?

Organisations are no longer overwhelmed by an absence of information; they are overwhelmed by an abundance of it. A single application may depend on thousands of open-source packages, each potentially introducing risk. Security has transformed from a periodic, pre-release gate into a continuous operational discipline.

As engineers, we have to evolve with it. Understanding a CVE is no longer just a task for a security researcher. Vulnerability management is becoming a core part of everyday software engineering, and for many organisations operating in Europe, it will soon become a legal obligation.

Compliance Is Becoming a Software Engineering Problem

Implementing Durable Execution

Reading and writing about topics we are learning is great, but there is nothing better than some hands-on approach, as such, let’s build a couple implementations of the pizza example described in yesterday’s article.

The first implementation is using the Temporal SDK to allow us to get more familiar with the details of Durable Execution, and consolidate a bit better what we are reviewing. In the second example, we will try to implement our incredible tiny very reduced version of the whole thing.

First example: Using the Temporal SDK

The full implementation of this example can be found in GitHub in the repository pizza-durable-execution. The code and the repository have been heavily documented, which what I think should be enough information to understand the example. But some quick overview is:

  • PizzaActivities: Activities are the side-effecting operations in Durable Execution.
  • PizzaActivitiesImpl: Concrete implementation of the activities.
  • PizzaOrderStarter: Starts a new pizza order workflow instance.
  • PizzaOrderWorkflow: The workflow interface defines the durable process contract.
  • PizzaOrderWorkflowImpl: The durable workflow implementation with pure orchestration, zero side effects.
  • PizzaWorker: The Worker process.

Once we run it, we should see something like:

The pizza worker

=================================================
Pizza Worker started. Polling: pizza-order-queue
Now run PizzaOrderStarter to place an order.
=================================================
10:50:09.222 [workflow-method-pizza-order-margherita-001-019e3031-78f6-7222-9627-e9f11a3367de] INFO d.b.pizza.PizzaOrderWorkflowImpl - [WORKFLOW] Starting pizza order for: margherita
10:50:09.247 [Activity Executor taskQueue="pizza-order-queue", namespace="default": 1] INFO d.b.pizza.PizzaActivitiesImpl - [ACTIVITY] Taking order for pizza: margherita
10:50:09.247 [Activity Executor taskQueue="pizza-order-queue", namespace="default": 1] INFO d.b.pizza.PizzaActivitiesImpl - [ACTIVITY] Order created → ORDER-MARGHERITA
10:50:09.256 [workflow-method-pizza-order-margherita-001-019e3031-78f6-7222-9627-e9f11a3367de] INFO d.b.pizza.PizzaOrderWorkflowImpl - [WORKFLOW] Order accepted → ORDER-MARGHERITA
10:50:09.259 [Activity Executor taskQueue="pizza-order-queue", namespace="default": 1] INFO d.b.pizza.PizzaActivitiesImpl - [ACTIVITY] Kitchen preparing pizza for order: ORDER-MARGHERITA
10:50:09.260 [Activity Executor taskQueue="pizza-order-queue", namespace="default": 1] INFO d.b.pizza.PizzaActivitiesImpl - [ACTIVITY] Pizza ready → PIZZA-ORDER-MARGHERITA
10:50:09.262 [workflow-method-pizza-order-margherita-001-019e3031-78f6-7222-9627-e9f11a3367de] INFO d.b.pizza.PizzaOrderWorkflowImpl - [WORKFLOW] Pizza prepared → PIZZA-ORDER-MARGHERITA
10:50:09.262 [workflow-method-pizza-order-margherita-001-019e3031-78f6-7222-9627-e9f11a3367de] INFO d.b.pizza.PizzaOrderWorkflowImpl - [WORKFLOW] Waiting 5 seconds for delivery window (durable timer)...
10:50:14.291 [Activity Executor taskQueue="pizza-order-queue", namespace="default": 1] INFO d.b.pizza.PizzaActivitiesImpl - [ACTIVITY] Dispatching delivery for: PIZZA-ORDER-MARGHERITA
10:50:14.292 [Activity Executor taskQueue="pizza-order-queue", namespace="default": 1] INFO d.b.pizza.PizzaActivitiesImpl - [ACTIVITY] Delivery confirmed → DELIVERED-PIZZA-ORDER-MARGHERITA
10:50:14.297 [workflow-method-pizza-order-margherita-001-019e3031-78f6-7222-9627-e9f11a3367de] INFO d.b.pizza.PizzaOrderWorkflowImpl - [WORKFLOW] Pizza delivered → DELIVERED-PIZZA-ORDER-MARGHERITA
10:50:14.301 [Activity Executor taskQueue="pizza-order-queue", namespace="default": 1] INFO d.b.pizza.PizzaActivitiesImpl - [ACTIVITY] Sending receipt for delivery: DELIVERED-PIZZA-ORDER-MARGHERITA
10:50:14.301 [Activity Executor taskQueue="pizza-order-queue", namespace="default": 1] INFO d.b.pizza.PizzaActivitiesImpl - [ACTIVITY] Receipt sent. Workflow complete.
10:50:14.305 [workflow-method-pizza-order-margherita-001-019e3031-78f6-7222-9627-e9f11a3367de] INFO d.b.pizza.PizzaOrderWorkflowImpl - [WORKFLOW] Workflow complete for order: ORDER-MARGHERITA

The pizza order starter

=================================================
Starting pizza order workflow...
=================================================
=================================================
Workflow finished. Pizza delivered!
=================================================

We can check the Temporal UI, and see our execution:

All necessary instructions for running it are present in the README of the project.

Second example: Implementing our own

The full implementation of this example can be found in GitHub in the repository mini-durable-execution-platform. The code and the repository have been heavily documented, which what I think should be enough information to understand the example. But some quick overview is:

  • MiniTemporal: Single class containing the whole project.
  • PizzaWorkflow: Durable orchestration of a pizza order.
  • WorkflowContext: The replay engine: the heart of durable execution.

Once we run it, we should see something like:

The pizza worker

worker started
[EXECUTING] prepare-dough
[EXECUTING] add-toppings
[EXECUTING] bake-pizza
[EXECUTING] prepare-dough
[EXECUTING] add-toppings
[EXECUTING] bake-pizza
[EXECUTING] deliver-pizza

The pizza order starter

[WAITING] prepare-dough
[WAITING] add-toppings
[WAITING] bake-pizza
=== JVM CRASH SIMULATED ===
workflowId=d2a8dc13-e428-4bc2-996e-c721d71592f7
...
[WAITING] prepare-dough
[WAITING] add-toppings
[WAITING] bake-pizza
[WAITING] deliver-pizza
===== ORDER COMPLETED =====
dough=dough-ready
toppings=toppings-added
baked=pizza-baked
delivery=pizza-delivered
=== WORKFLOW COMPLETED ===

All necessary instructions for running it are present in the README of the project.

Implementing Durable Execution

Durable Execution: The Runtime for Distributed Systems

Note: This article has two main sections. The first one is an abstract explanation of the Durable Execution concept. The second one is a simple workflow example to try to reduce the abstraction, and show a more realistic view to anchor the explanation. Depending on how you like to learn, feel free to read the explanation first, the example first, or even alternate between them while reading.


There is a quiet but important change happening in how we build software. It isn’t a sudden “revolution”, but rather a new way of thinking about how programs run across multiple servers. We call this Durable Execution.

Durable Execution could be described as a simple inversion of responsibility: instead of treating failure as something applications must anticipate and recover from, durable execution systems assume that failure is constant and design the execution model itself to survive it. Which means we are no longer just coordinating work across services, we are starting to treat execution itself as a persistent, stateful entity.

Most modern backend systems are built on an architecture that, on paper, looks clean and composable. A request enters a system, an orchestrator decomposes it into tasks, and a fleet of stateless workers executes those tasks independently. For example, when you click “buy” on a website, a request goes to a server, which then talks to a database, a payment processor, a storage system, and a shipping service among other things. In this set up, each service is responsible for doing one thing well, and persistence is delegated to databases and queues.

This model scales remarkably well in terms of throughput and organisational clarity. It is the backbone of microservices architecture. But as systems grow in complexity, something inevitable and subtle happens: workflows begin to leak across boundaries.

A “simple” business process such as processing a payment, or fulfilling an order, quietly evolves into a distributed orchestration of services. Each step is straightforward in isolation, yet the overall process becomes fragile not because any single component is complex, but because no single component owns the lifecycle of the workflow itself. The orchestration of these systems eventually involves:

  • retry policies embedded in clients and workers
  • state stored in databases with evolving schemas
  • queues that act as implicit progress trackers
  • compensating logic scattered across services
  • and operational heuristics encoded in dashboards and alerts

Eventually, the boundaries of the workflow start blurring, the workflow itself ceases to exist, and becomes highly entangled with the infrastructure. This is the root tension that durable execution tries to address.

To understand the shift, it helps to contrast two models of thinking about distributed systems.

In the traditional worker-based architecture, the system guarantees delivery of work. A message will eventually reach a worker. A job will eventually be retried. A task will eventually be processed. A completed or failed result will be published. But what is not guaranteed is the continuity of execution. If a process begins, partially completes, and then fails mid-way, nothing in the infrastructure inherently remembers what step it was on or what should happen next. That responsibility is pushed upward into application code and external state stores.

Durable execution tries to flip this assumption by, instead of treating an execution as ephemeral, it treats it as a persistent object. A workflow is not something that is “run”; it is something that exists over time. It has a history, a state, and a deterministic progression that can be paused, resumed, replayed, or migrated. This is the core idea behind systems such as Temporal, which model workflows as durable state machines whose execution history is recorded and reconstructed as needed. The runtime becomes responsible not only for executing steps, but for preserving the identity of the execution itself.

At the centre of durable execution lies a constraint that initially feels unnatural to most engineers: workflow code must be deterministic. This does not mean the system itself is deterministic in the mathematical sense. It means that given the same recorded history of events, the workflow must always reconstruct the same state and make the same decisions. This requirement exists because durable systems often rely on replay. When a workflow resumes after a failure, the runtime does not “continue” execution in the traditional sense. Instead, it reconstructs the workflow by replaying prior decisions and rehydrating state from a persisted event history. This has an important consequence: side effects cannot be executed freely during replay. External interactions such as API calls, database writes, message emissions, must be carefully separated from the logical flow of the workflow. In practice, this introduces a separation between:

  • activities (the side-effecting operations performed by workers)
  • workflow logic (the durable orchestration layer)

This separation is what allows executions to be safely paused and resumed without ambiguity. While this may feel restrictive, it is precisely this constraint that enables durability.

Let’s try to put side by side some of the characteristics of a traditional system, and the characteristics of Durable Execution systems.

In traditional systems, retries are usually an implementation detail. A worker fails, a message is requeued, and eventually the task is attempted again. But retries quickly become more complex than they first appear, especially when failures happen mid-workflow rather than at the boundaries of tasks. What should happen if a payment succeeds but inventory reservation fails? Should the system retry the inventory step, or compensate the payment? What if compensation itself fails? What if the system crashes between deciding to compensate and actually doing so? These questions are not edge cases; they are the natural consequence of long-running distributed coordination.

Durable execution systems turn retries into a first-class runtime concept. Instead of scattering retry logic across services, the workflow engine tracks execution attempts as part of its history. Time itself becomes a managed dimension of the system, with timers, delays, and waiting periods becoming durable constructs rather than external scheduling hacks. Even waiting for days becomes structurally simple, because the workflow state is persisted independently of process memory. In this sense, durable execution is not just about reliability under failure. It is about treating time as a durable resource.

Moving deeper, once a workflow spans multiple services, failure is no longer binary, it is often partial. Some steps succeed, others fail, and the system must reconcile an inconsistent reality. This is where compensation logic enters the picture, often through patterns such as sagas.

If you don’t know, a saga is essentially a distributed transaction without atomicity. Instead of rolling everything forward or backward as a single unit, the system defines compensating actions that attempt to undo completed steps when later steps fail. In traditional architectures, sagas are notoriously difficult to implement correctly because their logic is distributed across services and tightly coupled to operational state.

Durable execution brings sagas into the workflow layer itself. Compensation is no longer a scattered concern but part of the execution model. The workflow runtime knows what has succeeded, what has failed, and what needs to be undone. This does not eliminate complexity, but it changes where complexity lives. Instead of being embedded in infrastructure glue code, it becomes explicit in the structure of the workflow.

Perhaps the most important conceptual shift introduced by durable execution is the normalisation of long-running processes. In traditional request-driven systems, time is implicitly assumed to be short. A request is expected to complete within milliseconds or seconds. Anything longer is pushed out of band into queues, schedulers, or cron jobs. But many real-world processes do not fit this model. They are inherently extended in time:

  • waiting for human approval
  • integrating with third-party systems
  • coordinating multi-stage financial flows
  • handling asynchronous physical-world processes

Durable execution embraces this directly. A workflow can span seconds, hours, or weeks without requiring external orchestration mechanisms to simulate persistence. These systems treat long-running execution not as an anomaly, but as a primary use case.

An increasingly important driver of interest in durable execution comes from the domain of AI agents. Modern agentic systems are not stateless request handlers. They maintain evolving context, interact with external tools, retry operations, and often run through multi-step reasoning and action loops that can span long periods of time. Without durable execution, these systems are fragile in predictable ways:

  • a crash loses context
  • a timeout breaks continuity
  • partial tool execution leads to inconsistent state
  • retries can duplicate side effects

What is emerging is a recognition that AI agents are, structurally, workflows. They are not single computations, they are long-running, stateful processes that require persistence, replayability, and controlled side effects. This is why systems such as durable workflow engines are increasingly being explored as the underlying runtime for agent orchestration, not just business process automation.

It is tempting to view durable execution systems as simply better workflow engines, but that framing underestimates what is actually changing. Traditional orchestrators coordinate tasks across workers. They are, in essence, message routers with state tracking bolted on. Durable execution systems, by contrast, begin to resemble execution environments. They manage:

  • state persistence
  • execution history
  • scheduling and timers
  • retries and failure recovery
  • deterministic replay
  • coordination across services

This is why a useful mental model is to think of them as a kind of distributed operating system. Not for hardware resources, but for business logic execution over time. The comparison is not perfect, but it is instructive. Just as operating systems abstracted away hardware complexity to allow applications to run reliably on unstable machines, durable execution abstracts away distributed failure to allow workflows to run reliably on unstable networks.

Durable execution is still an emerging paradigm. It is not yet the default model for building distributed systems, and many teams continue to rely successfully on queues, workers, and stateless services. But the pressure that led to its development is becoming more pronounced. Systems are becoming more distributed, workflows are becoming longer-lived, and AI systems are introducing a new class of stateful computation that does not fit neatly into request/response paradigms.

What durable execution offers is not a new tool, but a new boundary. It draws a line around execution itself and says: this is something worth making persistent. If databases taught us how to make data reliable, durable execution is attempting something analogous for computation. And if that trajectory continues, workflow engines may evolve from infrastructure components into something closer to what operating systems became for the hardware era: a foundational layer that quietly defines how everything else runs.

Simple workflow example

If you are like me, after reading this you are probably thinking “Yeah, that sounds good, but what does it actually look like?“. For that reason, let’s try to make the abstraction more concrete, and look at a single workflow.

Imagine a simple order fulfilment process: ordering a pizza. In a traditional system, this would be split across services, queues, and callbacks. In a durable execution system, it is expressed as a single continuous workflow, but importantly, it does not execute like a function call. It behaves more like a stateful timeline.

The workflow (conceptual view)

Workflow: PizzaOrder
Step 1 → Take order ("margherita")
Step 2 → Prepare pizza (orderId)
Step 3 → Wait for 5 seconds (or 5 hours, or 5 days)
Step 4 → Deliver pizza
Step 5 → Send receipt

At first glance, this looks trivial. The important part is what happens under the hood.

What actually happens at runtime

When the workflow starts, the runtime does not “execute everything”. It begins a controlled sequence of recorded decisions.

1. Start of execution

The system creates a durable record:

WorkflowInstance: PizzaOrder-128
State: STARTED
History: []

It then executes Step 1.

→ Schedule activity: TakeOrder("margherita")

This is not executed inline. It is dispatched to a worker. The workflow pauses.

2. First suspension point (important idea)

At this moment, the workflow is not “running”. It is:

  • persisted
  • waiting
  • fully safe to crash
  • resumable from history
WorkflowInstance: PizzaOrder-128
State: WAITING_FOR(TakeOrder)
History:
- Scheduled TakeOrder("margherita")

A worker eventually responds:

Result: ORDER-MARGHERITA

3. Replay begins (the non-obvious part)

Now comes the key durable execution concept. If the workflow needs to continue (or recover from failure), the runtime does not simply “continue execution”. Instead, it replays the workflow from the beginning using recorded history:

Replay:
- Step 1: TakeOrder → already completed (from history)
- Step 2: PreparePizza(orderId=ORDER-MARGHERITA)

This is where determinism matters, the workflow code must behave consistently during replay.

4. Another suspension

The system schedules the next activity:

→ Schedule activity: PreparePizza(ORDER-MARGHERITA)

Again, execution pauses.

State: WAITING_FOR(PreparePizza)
History:
- TakeOrder completed
- PreparePizza scheduled

A worker completes it:

Result: PIZZA-ORDER-MARGHERITA

5. Time becomes a first-class construct

Now the workflow reaches an unusual step:

WAIT 5 seconds

In a traditional system, this would require:

  • cron jobs
  • timers
  • sleep threads
  • external schedulers

In a durable system, time itself is persisted:

Workflow state:
Timer scheduled: +5 seconds

The workflow is now completely idle, but still alive as a durable object. It can survive:

  • process crashes
  • machine restarts
  • deployments
  • network partitions

Nothing is lost.

6. Resumption after time passes

When the timer fires:

→ Resume workflow

Replay reconstructs state again:

History:
- TakeOrder ✓
- PreparePizza ✓
- Wait(5s) ✓

Now execution continues:

→ Schedule activity: DeliverPizza(PIZZA-ORDER-MARGHERITA)

7. Completion

Finally:

→ Schedule activity: SendReceipt(deliveryId)
→ Workflow COMPLETE

At the end, the system has not just executed steps. It has produced a durable execution trace:

Workflow History:
1. TakeOrder → ORDER-MARGHERITA
2. PreparePizza → PIZZA-ORDER-MARGHERITA
3. Timer → 5s elapsed
4. DeliverPizza → DELIVERED
5. SendReceipt → RECEIPT SENT

The key insight this example is trying to surface is that what matters is not the steps themselves, any system can execute steps. What matters is that the workflow is not stored as state we manage, but as history the runtime owns. This is why durable execution feels different from:

  • queues
  • cron jobs
  • orchestration services
  • worker pools

It is not just “a better way to coordinate tasks”, it is a system where execution itself becomes a recoverable data structure. Once we internalise this model, several earlier concepts become much clearer:

  • retries are not logic → they are replay
  • state is not stored → it is reconstructed
  • time is not external → it is persisted
  • workflows are not running → they are waiting, resuming, continuing

And this is why systems like Temporal feel less like task schedulers and more like a runtime layer for distributed computation over time.

Durable Execution: The Runtime for Distributed Systems

Charging for the ink, not the ideas

I am sure most of you are familiar with the story of Charles Steinmetz, or one of its many variations. Steinmetz was a brilliant engineer at General Electric, and the story goes like this.

Henry Ford was having trouble with a massive generator and called in Steinmetz. After listening to the machine for a few moments, Steinmetz took out a piece of chalk and made a small ‘X’ on a specific metal casing. Ford’s engineers opened it up and found the defect exactly where the mark was. When Steinmetz sent a bill for $1,000, Ford, ever the businessman, asked for an itemised invoice. Steinmetz replied:

  • Making chalk mark: $1
  • Knowing where to mark: $999

Ford paid the bill without further question. He understood that the physical act was trivial; the value lay in the decades of experience required to know exactly where that one-dollar mark belonged.

We are living through a remarkably similar moment, yet we may be heading in the opposite direction. Many AI companies are attempting to persuade us that the act of generation is more valuable than what is generated. They are moving away from simple subscription models towards pay-per-token billing. Even those that have not fully transitioned are clearly pivoting that way. In doing so, they are, perhaps unintentionally, asking us to pay for the weight of the chalk.

A token is a mathematical fragment, the raw material of a response. Billing by the token reflects real computational costs, but it also participates in a market that can prize volume over validity. Intelligence starts to resemble a metered utility, much like water or electricity. But intelligence is not a liquid; it is a coordinate. In software and engineering, the most elegant solution is rarely the longest. A thousand lines of generated code may solve a problem, but they are often a liability, a burden of technical debt. Ten lines of precise logic can be a masterpiece.

Under the current model, however, the thousand-line mess can end up being priced as though it were a hundred times more valuable than the ten-line stroke of genius. We can find ourselves paying for the stuttering of the machine, the computational friction it incurs while searching for an answer, rather than for the answer itself. We are, quite literally, paying for the ink and ignoring the idea.

This push towards ever greater output is accelerating beyond the limits of human review. Where an engineer might once have produced a single page of clear, concrete documentation, a model now generates thousands. This is not necessarily progress; it can become a flood. It creates an artificial demand for even more powerful and expensive models, just to process the noise produced by earlier ones. We can end up in a loop in which we need AI to summarise the verbosity of other AI.

This begins to reveal a structural tension in the current AI arms race. The incentives do not always point towards efficiency; they often reward scale. By flooding the ecosystem with information, the need for larger contexts and more powerful reasoning becomes easier to justify. These, in turn, support higher price points and more ambitious positioning. The result is not a map to the ‘X’, but an ever-growing supply of chalk, along with the tools to manage it.

This creates a perverse incentive for the future of technology. If we measure the value of AI by the number of tokens it produces, we encourage a digital world of bloat. We risk being buried under a mountain of cheap, generated noise, where quantity is mistaken for quality. It is a system that can reward the machine for being chatty rather than correct.

The true revolution of artificial intelligence should not be that it makes the chalk mark easier to produce. The revolution is that it should help us find the ‘X’ faster. But as long as the price remains closely tied to the token, the industry will tend to focus on the tool rather than the result.

We should eventually demand a different kind of invoice from the architects of these models. We should stop subsidising the cost of digital ink and start valuing the precision of knowledge. Until we shift our perspective, we are not fully purchasing intelligence; we are still largely paying for the act of writing. Steinmetz’s ‘X’ was valuable because it was singular and precise. If he had covered the entire generator in chalk, his bill would not have been worth a penny, regardless of how much chalk he used.

A deeper question follows. What if AI does not consistently deliver what is promised, not because it cannot, but because the incentives are misaligned? At present, many incentives favour the production of large volumes of content that require millions of tokens and iterations to process. They do not always favour the creation of systems that can deal with complexity in a genuinely intelligent way. It is worth asking whether the pursuit of revenue might, at times, be steering us away from the outcome we actually want.

Charging for the ink, not the ideas

Fact-checking GitHub controversy

In recent days, a familiar kind of narrative has swept across developer circles: the claim that GitHub is ‘dying’. It has appeared in videos, threads, and blog posts with the usual hallmarks of online virality: strong opinions, selective facts, and a tone of urgency that suggests an imminent collapse. Given how central GitHub remains to modern software development, it is hardly surprising that such claims attract attention. Yet, as is often the case, the reality is both less dramatic and more interesting than the headline.

The current wave of criticism did not emerge from nowhere. It was catalysed, in part, by a public critique from Mitchell Hashimoto, a figure whose opinions carry weight within the developer community. His frustration centred on reliability: repeated outages, degraded performance, and an overall experience that, in his estimation, fell well short of what one expects from a platform of GitHub’s stature. His decision to move his terminal project, Ghostty, away from GitHub was not merely a personal choice but a symbolic gesture that resonated with others who had experienced similar issues. It is important, however, to interpret this moment with care. High-profile departures can shape perception disproportionately; they signal discontent, certainly, but they do not in themselves constitute evidence of a broader exodus.

Reliability concerns are nonetheless real. GitHub has experienced intermittent instability in recent months, and while no large-scale platform is immune to outages, expectations in this domain are exceptionally high. Developers rely on such services not only for storage but for collaboration, automation, and deployment pipelines. When interruptions occur, they ripple through entire workflows. The perception that reliability has slipped, even if only temporarily, can therefore have an outsized impact on trust. What remains less clear is whether these incidents represent a systemic decline or a series of unfortunate but ultimately transient issues. At present, the evidence supports frustration, but not collapse.

Source: GitHub’s own status logs for late April 2026

Alongside these operational concerns sits a more subtle, yet arguably more consequential, shift: changes to data usage policies for GitHub Copilot. As of late April 2026, GitHub moved to an opt-out model for certain forms of data collection used in training its AI systems. In practical terms, this means that interactions, such as prompts and generated code, may be used to improve models unless the user explicitly disables this behaviour. For many developers, particularly those working across personal and professional contexts, this introduces a degree of ambiguity that did not previously exist. The concern is not that entire repositories are being indiscriminately absorbed into training datasets, as some commentary has suggested, but rather that the boundary between private work and aggregated learning has become less immediately transparent.

It is worth noting that these concerns are not without qualification. Organisational and enterprise tiers are excluded from such data usage, and the opt-out mechanism remains available. Nevertheless, defaults matter. In software, as in many other domains, what is enabled by default often defines the practical reality of a system. The shift, therefore, is less about technical risk in the strictest sense and more about a recalibration of expectations around control and consent.

A further source of unease arises from changes to pricing. GitHub’s move towards a more usage-based model for Copilot reflects a broader industry trend: the recognition that AI-assisted tooling incurs substantial and uneven costs. Under such a model, light users may see little difference, while heavier users, those running extended sessions or integrating AI deeply into their workflows, may encounter higher and less predictable expenses. It is not difficult to understand why this has been received with scepticism. Developers tend to value clarity and stability in pricing, and any departure from that can feel, rightly or wrongly, like a shifting of the goalposts.

Compounding these issues was a temporary pause on new Copilot sign-ups, justified by GitHub as a measure to maintain service quality. Although this decision was framed in pragmatic terms, it inevitably fuelled speculation about underlying capacity constraints. Whether such speculation is warranted remains unclear; what can be said is that the optics of limiting access, even temporarily, sit uneasily alongside narratives of rapid expansion and technological progress.

Taken together, these developments form the basis of the current backlash. Yet it is equally important to consider what has been overstated or misrepresented in the process. Claims of a mass departure from GitHub, for instance, are not supported by credible evidence. The platform continues to dominate its space, and while alternatives such as GitLab attract periodic attention, there is no indication of a large-scale migration. Similarly, the suggestion that GitHub’s focus on artificial intelligence has directly caused reliability issues remains speculative. Correlation, in this case, has been readily interpreted as causation without sufficient proof.

Adding to the general unease, a number of more sensational claims began to circulate online, including suggestions that Copilot had, at one point, inserted what resembled promotional content into pull requests. These reports are difficult to substantiate and have not been confirmed by reliable sources, yet their spread is telling in itself. They reflect a growing suspicion among some developers that the platform’s priorities may be shifting in ways that are not entirely aligned with user interests. Even when such claims prove unfounded, they tend to gain traction in an environment where trust is already under strain.

What, then, should one make of the situation as a whole? Rather than signalling decline, it seems more accurate to view this moment as a period of adjustment. GitHub, like many technology platforms, is navigating the complex transition towards AI-integrated workflows while attempting to balance cost, performance, and user trust. Each of the current points of contention, reliability, data usage, or pricing, reflects a facet of that broader challenge.

For developers, the appropriate response is neither alarm nor indifference, but informed attention. The concerns being raised are not trivial; they touch on fundamental aspects of how tools are built, maintained, and monetised. At the same time, the more dramatic narratives obscure as much as they reveal. GitHub is not ‘dying’, but it is changing, and not all of those changes will be universally welcomed.

In the end, the significance of this episode lies less in any single policy or outage, and more in what it reveals about the relationship between developers and the platforms they depend on. Trust, once established, can be surprisingly resilient, but it is not immutable. Moments like this serve as a reminder that even the most entrenched tools must continually justify that trust, not through promises or positioning, but through consistent, transparent, and reliable behaviour over time.

Fact-checking GitHub controversy