Most dashboards, fraud rules and feature pipelines ask the same SQL question again and again while the data underneath keeps changing. A warehouse answers by recomputing from scratch on a schedule, so results are minutes or hours old. A stream processor such as Flink answers in milliseconds, but you write operators, manage state and reason about watermarks yourself. Materialize sits between them: you write ordinary PostgreSQL-flavoured SQL, including multi-way joins and aggregates, and it keeps the answer correct as each change arrives.
This article explains how that works from first principles, then shows how to build with it: sources, views and indexes, SUBSCRIBE consumers, and running it without running out of memory. Syntax shown here was checked against the Materialize documentation in October 2026; confirm it against the version you run.
The problem incremental view maintenance solves
Suppose an orders table holds 200 million rows and a dashboard shows revenue per region. Between two one-second refreshes perhaps 3,000 orders arrived and 40 were cancelled. The necessary work is proportional to those 3,040 changes, not to 200 million rows.
Incremental view maintenance turns that observation into a system. Instead of storing a result and recomputing it, the engine turns every input change into the exact change of the output. The hard cases are joins where any side can change, outer joins, DISTINCT and aggregates over joins. Materialize compiles SQL into dataflows built on Timely Dataflow and Differential Dataflow, two Rust libraries designed to handle exactly those cases.
The update model: (row, time, diff)
Everything in Materialize is a collection of updates. Each update is a triple: a row, a logical timestamp, and a diff, an integer saying how many copies of that row were added (positive) or removed (negative) at that time. An insert is diff +1. A delete is diff -1. An update is a delete of the old row and an insert of the new one at the same timestamp.
Operators consume and produce these triples. A filter passes them through or drops them. A map rewrites the row and keeps the diff. A join multiplies diffs of matching rows. An aggregate, for each changed group, retracts its old output row and emits the new one. Here is a count per region, followed through one order moving from us to eu:
input changes at time 7: (order 42, region=us) diff -1
(order 42, region=eu) diff +1
count(*) GROUP BY region emits at time 7:
(us, 5) diff -1 -- retract old count
(us, 4) diff +1 -- assert new count
(eu, 2) diff -1
(eu, 3) diff +1Two ideas make this efficient. First, operators that need history, such as joins and aggregates, keep it in an arrangement: an indexed, compacted collection of updates organised by key. A join looks up the other side's arrangement for just the changed keys. Second, timestamps let the engine know when an output at time 7 is complete, because every input has promised it will send nothing earlier than time 8. That frontier is what lets Materialize hand you answers that are consistent rather than half-updated.
Architecture: storage, compute and the adapter
Materialize separates three jobs. Storage ingests data from external systems, decodes it, assigns timestamps, and writes the resulting updates durably to its persistence layer, which keeps data in object storage and coordinates metadata through a consensus store. Compute runs dataflows that read from storage and maintain indexes and materialized views. The adapter speaks the PostgreSQL wire protocol on port 6875, holds the catalog, plans queries and picks the timestamp each query reads at.
Compute and ingestion run in clusters, which the documentation defines as isolated pools of compute resources for sources, sinks, indexes, materialized views and ad-hoc queries. A cluster can have more than one replica running the same work, which gives fault tolerance: if one replica dies, the others keep serving. Because a cluster's state lives in memory, a new or restarted replica must hydrate, rebuilding its arrangements by reading from storage rather than from the upstream database.
Bringing data in: PostgreSQL and Kafka sources
A PostgreSQL source reads the database's logical replication stream, so it sees committed transactions in order without polling. The upstream needs wal_level = logical, a publication listing the tables, and REPLICA IDENTITY FULL on each table so updates and deletes carry the full old row. Materialize creates a replication slot whose name starts with materialize_. Credentials go in secrets, network details in a reusable connection:
-- on the upstream PostgreSQL database
ALTER TABLE orders REPLICA IDENTITY FULL;
ALTER TABLE customers REPLICA IDENTITY FULL;
CREATE PUBLICATION mz_pub FOR TABLE orders, customers;
-- in Materialize
CREATE SECRET pg_password AS '...';
CREATE CONNECTION pg_conn TO POSTGRES (
HOST 'orders-db.internal', PORT 5432, USER 'materialize',
PASSWORD SECRET pg_password, SSL MODE 'require', DATABASE 'shop'
);
CREATE SOURCE shop_pg IN CLUSTER ingest
FROM POSTGRES CONNECTION pg_conn (PUBLICATION 'mz_pub')
FOR TABLES (orders, customers);This is the older source syntax. Newer releases add a CREATE TABLE ... FROM SOURCE form that the documentation says handles upstream column additions and drops without downtime; with the older form, schema changes upstream mean recreating the source.
Kafka sources choose a format and an envelope. The envelope tells Materialize how to interpret records: with no envelope every record is an append-only insert; ENVELOPE UPSERT treats the key as identity, so a new value replaces the old and a null value deletes it; ENVELOPE DEBEZIUM decodes before and after images from change-data-capture topics.
CREATE SECRET kafka_password AS '...';
CREATE CONNECTION kafka_conn TO KAFKA (
BROKER 'broker-1.internal:9092',
SASL MECHANISMS = 'SCRAM-SHA-256',
SASL USERNAME = 'mz', SASL PASSWORD = SECRET kafka_password
);
CREATE CONNECTION csr_conn TO CONFLUENT SCHEMA REGISTRY (URL 'https://registry.internal');
CREATE SOURCE payments IN CLUSTER ingest
FROM KAFKA CONNECTION kafka_conn (TOPIC 'payments')
FORMAT AVRO USING CONFLUENT SCHEMA REGISTRY CONNECTION csr_conn
ENVELOPE UPSERT;Upsert sources must remember the current value for every key to emit retractions, so their state grows with distinct keys, not throughput.
Views, indexes and materialized views
This is the decision that most affects cost. A view is just a named query; nothing is computed until something reads it. An index on a view keeps its results in memory in one cluster, ready for fast lookups and for reuse by other queries in that cluster. A materialized view maintains results and writes them to durable storage, so any cluster can read them and sinks can export them.
| Object | Where results live | Use it when |
|---|---|---|
| View | Nowhere until queried | Naming intermediate logic; it costs nothing on its own |
| Indexed view | Memory of one cluster | Low-latency SELECT and SUBSCRIBE served from that cluster |
| Materialized view | Durable storage, read by any cluster | Sharing results across clusters, feeding sinks, decoupling compute from serving |
A common layout uses three clusters, each created first with CREATE CLUSTER: one for sources, one that maintains materialized views, and one that indexes those views for applications, so a runaway query in serving cannot starve ingestion.
CREATE VIEW order_facts AS
SELECT o.id, o.amount_cents, o.created_at, c.region
FROM orders o JOIN customers c ON c.id = o.customer_id
WHERE o.status <> 'cancelled';
CREATE MATERIALIZED VIEW revenue_by_region IN CLUSTER transform AS
SELECT region, sum(amount_cents) AS revenue_cents, count(*) AS orders
FROM order_facts
WHERE mz_now() <= created_at + INTERVAL '1 hour'
GROUP BY region;
CREATE INDEX revenue_by_region_idx IN CLUSTER serving ON revenue_by_region (region);
Time: temporal filters with mz_now()
The WHERE clause above is a temporal filter. mz_now() returns Materialize's current logical timestamp, and a view may only use it in comparisons like this one, where a row's validity window can be computed from its own columns. Materialize then schedules each row's retraction for the moment its window closes, so the hour-long window slides without scanning anything. Without the filter, the view would hold every order ever placed. Temporal filters are the main tool for bounding memory in views over append-only data; compare them with the window operators in stream windowing.
Consistency: what a SELECT sees
The default isolation level is strict serializable: each query reads one consistent timestamp across all its inputs, and that timestamp respects real-time order between queries. You never see an order counted in one view but missing from another view joined in the same query.
Strict serializability does not by itself mean the newest upstream data is included. By default a query reads at a timestamp Materialize has already ingested. With SET real_time_recency = true, available only at strict serializable, a query waits until Materialize has ingested everything that was visible in the external sources when the query arrived. That gives read-your-writes against the upstream database at the price of latency, so enable it only where it matters.
Reading results: SELECT, SUBSCRIBE and sinks
SELECT against an indexed view is a lookup and returns in milliseconds. SUBSCRIBE streams changes: an initial snapshot, then every update as rows with mz_timestamp and mz_diff. With the PROGRESS option, extra rows marked mz_progressed tell you that a timestamp is complete, which is how a consumer knows it can apply a batch atomically. The documentation recommends wrapping SUBSCRIBE in a cursor so client drivers do not buffer the endless result:
import psycopg # psycopg 3
conn = psycopg.connect(DSN, autocommit=True) # DSN points at port 6875, sslmode=require
cur = conn.cursor()
cur.execute("BEGIN")
cur.execute("DECLARE c CURSOR FOR SUBSCRIBE revenue_by_region WITH (PROGRESS)")
pending = {} # timestamp -> list of (diff, row)
while True:
cur.execute("FETCH ALL c WITH (timeout='1s')")
cols = [d.name for d in cur.description]
for values in cur.fetchall():
rec = dict(zip(cols, values))
ts = rec.pop("mz_timestamp")
if rec.pop("mz_progressed"):
# every update with a timestamp below ts is now final
for t in sorted(k for k in pending if k < ts):
apply_atomically(pending.pop(t)) # e.g. one cache transaction
else:
diff = rec.pop("mz_diff")
pending.setdefault(ts, []).append((diff, rec))The consumer reads columns by name, buffers by timestamp, and applies a timestamp only when progress passes it; applying rows one at a time would briefly show a retracted count with no replacement. For durable export to other systems, create a sink that writes a materialized view's changes to a Kafka topic, and let downstream consumers follow exactly-once processing practice on their side.
Worked example: a live revenue panel
An online shop wants a panel showing revenue and order counts per region for the last hour, refreshed live, plus an alert when a region's payment failure rate passes 5 percent. Orders and customers live in PostgreSQL; payment results arrive on a Kafka topic keyed by payment id.
The team creates three clusters. The ingest cluster holds the PostgreSQL source and the upsert Kafka source. The transform cluster maintains revenue_by_region and a second materialized view joining payments to orders, grouping by region and computing failed over total for the last 15 minutes with a temporal filter. The serving cluster indexes both views. The dashboard backend runs one SUBSCRIBE per view and pushes each completed timestamp to browsers; the alerting service subscribes to rows above 0.05.
Sizing follows the state, not the event rate: one hour of joined orders, every customer, and one value per payment key. The team measures memory after hydration with production-sized data and adds headroom.
Operating Materialize
- Watch memory per cluster. Arrangements are in memory. A join on a high-cardinality key, a missing temporal filter or an unbounded upsert key space shows up as steady growth. The system catalog exposes introspection relations for arrangement sizes and dataflow state; check them before resizing.
- Plan for hydration. A new replica or a newly created index must rebuild state from storage before it serves. Add a replica and let it hydrate before removing the old one.
- Share arrangements. Indexes in a cluster can be reused by other dataflows in that cluster. Index join keys you use repeatedly and read
EXPLAINto confirm reuse. - Watch upstream slots. A PostgreSQL replication slot holds write-ahead log until it is consumed. If the source stalls, the upstream disk fills. Alert on slot lag in PostgreSQL, and drop slots for sources you remove.
Failure modes
| Symptom | Likely cause | Fix |
|---|---|---|
| Cluster memory climbs without bound | View over append-only data with no temporal filter | Add an mz_now() filter or aggregate earlier |
| Upstream disk fills | Stalled source holding a replication slot | Alert on slot lag; fix or drop the source |
| Dashboard flickers wrong totals | Consumer applies SUBSCRIBE rows one by one | Buffer by timestamp and apply on progress |
| Source breaks after an ALTER TABLE upstream | Older source syntax does not follow schema changes | Use the newer table-from-source form or recreate |
Trade-offs against other streaming engines
Compared with Flink, Materialize trades operator-level control for SQL correctness: you do not write watermarks or state TTLs, and joins between changing tables are handled for you, but you cannot drop to custom operators, and state must fit the memory of a cluster. Compared with Kafka Streams, it is a separate service rather than a library in your application, which means another system to run but also one consistent SQL view across many topics and databases. Compared with a warehouse, results are fresh within seconds, but you pay continuously for memory holding state. Choose it when the same queries run constantly over changing data.
What to do next
- List the queries you recompute most often and estimate, for each, the state it needs: distinct keys, window length and joined table sizes.
- Connect one upstream database with a publication and REPLICA IDENTITY FULL, and alert on its replication slot lag from day one.
- Build one materialized view with a temporal filter, index it in a separate serving cluster, and read it with SELECT.
- Write a SUBSCRIBE consumer that buffers by timestamp and applies on progress, and test it by updating rows upstream.
- Measure memory after hydration with production-sized data, then size replicas with headroom.
- Decide per session whether you need real-time recency, and enable it only where read-your-writes matters.