Skip to content

design: Materialize SDK for SUBSCRIBE consumption and custom sinks - #37463

Open
bobbyiliev wants to merge 19 commits into
MaterializeInc:mainfrom
bobbyiliev:subscribe-sdk-design
Open

bobbyiliev wants to merge 19 commits into
MaterializeInc:mainfrom
bobbyiliev:subscribe-sdk-design

Conversation

@bobbyiliev

@bobbyiliev bobbyiliev commented Jul 6, 2026 •

Copy link
Copy Markdown
Contributor

Proposes the Materialize SDK: a Rust protocol core with no I/O, compiled into thin per-language packages that use their native drivers, with subscribe and sink modules. It answers the Sink SDK product requirements, includes the subscribe semantics review checked against the code, plans the move onto durable subscriptions (#38468), and makes turbopuffer the first sink. Prototypes: #37483 (Rust), #37517 (Python). Diagrams (Materialize-internal): https://claude.ai/artifact/Q6193A4KskJpfWgUuYNY6v


TL;DR

  • What: an official SDK for building your own sink from Materialize to anything: search indexes, caches, your own Postgres, webhooks. Rust first, with Python and Node to follow.
  • Why customers care: no lost or duplicate data after a crash or restart. Several views can be written at one consistent point in time, which log-based CDC tools cannot do from offsets alone. No Kafka needed.
  • How it's built: every hand-written sink we found had a different bug in the same logic: when a batch is safe to write, how to resume without gaps, how to line up several views. So that logic lives once, in a small Rust library (the protocol core) compiled into each language package. It runs inside the customer's app, not as a service, and never opens a connection. Each language keeps its usual Postgres driver. Python, Node and later languages therefore behave the same, and one fix reaches every language in one release.
  • What a sink author writes: open the target, write a batch, commit. The SDK handles the snapshot, resume after restart, retries, dead-lettering and clear errors.
  • First proof: a turbopuffer sink running on our internal context graph in October.
  • Input wanted: the name (deferred question 1), whether the SDK can rely on subscribing by id (question 2), and what a sink should do when a blue/green deploy swaps the view it reads (question 3).

@sjwiesman

Copy link
Copy Markdown
Contributor

I would suggest some sort of sdk agnostic test suite so we can validate the implementations.

@bobbyiliev bobbyiliev changed the title design: Add design doc for Subscribe SDKs and external sink framework design: Materialize SDK for SUBSCRIBE consumption and custom sinks Oct 7, 2026
@bobbyiliev

Copy link
Copy Markdown
Contributor Author

Agreed. The Testing section now makes a shared conformance suite the contract: JSON test files every package runs through its binding, plus an end-to-end suite in this repo that drives each SDK through the same scenarios in the nightlies. Does that cover what you had in mind?

Comment on lines +579 to +590
### Repository layout and releases

```
misc/materialize-sdk/
core/ protocol core (Rust library)
spec/ behavior spec, versioned
conformance/ vectors, shared by every package
python/ package: PyO3 binding, psycopg transport, sinks
node/ package: WebAssembly build, node-postgres transport, sinks
rust/ package: protocol core directly, tokio-postgres transport, sinks
test/ mzcompose end-to-end suite, run in the nightlies
```

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

@jubrad I know we discussed one repo per language SDK, but could we keep everything in the materialize repo and release the way dbt-materialize does? A version bump PR merges to main and the deploy pipeline publishes any version that is not yet on PyPI or npm, so there are no tags and every package stays on one core version.

@bobbyiliev
bobbyiliev marked this pull request as ready for review October 7, 2026 11:35

@antiguru antiguru left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Posted by Claude Code on behalf of @antiguru.

Reviewed against the durable subscribe design (#38468) and main at 2d2a689. The direction is right: a pure core, the target-side checkpoint, the narrow scope, and splitting commit() from next(). The items below are gaps in error classification, size limits, and turbopuffer fencing, plus places where the durable-subscription section diverges from #38468.

Blocking

Snapshots and catch-up larger than max_result_size cannot be delivered. The doc says a large snapshot "will arrive as chunks marked partial", but the server fails before any chunk reaches the client. PendingSubscribe::stash (src/compute-client/src/service.rs:537-556) sums the stashed updates of an unclosed timestamp and errors past max_result_size (1 GiB default), and the compute sink holds the whole snapshot in rows_to_emit. A resume after long downtime hits the same limit when the frontier jumps in one step. The failure is deterministic, so retry repeats it. The doc needs to state the bound and give it a typed error. Durable subscriptions do not change this, because they resume at timestamp granularity and a snapshot is a single timestamp.

Dataflow errors are XX000, not "varies". Dataflow errors and the max_result_size error both reach the client as AdapterError::Unstructured, which maps to INTERNAL_ERROR (src/adapter/src/error.rs:1063). StreamPoisoned therefore cannot be separated from genuine internal errors or the size error without matching message text, so the HistoryLost match would not be the only message matching. 53200 is also shared by SubscribeFellBehind and ResultSize (error.rs:1014-1015), so FellBehind keyed on the code alone resumes into a deterministic failure.

turbopuffer: a stale worker can resurrect deleted keys. The per-namespace checkpoint document protects replay by the live worker, which drops changes below the checkpoint. A stale worker read its checkpoint before it was fenced, and turbopuffer conditions are per document, so its data writes cannot be conditioned on the checkpoint document's epoch. Its upsert of a key that the live worker has since deleted finds no document and is applied unconditionally. The determinism argument covers older rows replacing newer ones, but not a write after a delete. This makes tombstones (open question 4) a correctness requirement, not a usability tradeoff.

The durable attach should keep AS OF F - 1. The section says the SDK attaches with USING DURABLE SUBSCRIPTION in place of AS OF F - 1 and filters below its own position. #38468 recommends that a library holding T position the read with AS OF T - 1 instead, because that fails loudly when a recreate or RESET has moved H past T, while a filter hides the gap whenever the opening progress check is missed. The SDK already confines the - 1 to one function, so dropping AS OF buys nothing. "Removes the AS OF F - 1 arithmetic on the default path" should go from the list of benefits.

The SDK scope is wider than what #38468 accepts. #38468 targets storage collections only (tables, materialized views, sources), rejects views and indexes, and rejects temporal filters. The SDK accepts any object, and a non-materialized view is inlined, so its resume needs history on every input and reads their snapshots, as the "Snapshot elision" section notes. Accepting views or mz_now() filters now turns the later move into a narrowing, which is the API break the section wants to avoid. Restrict durable mode to storage collections without temporal filters from the start.

Strong suggestions

  • Retention sizing changes shape under durable subscriptions rather than disappearing. ACKNOWLEDGE WITHIN is wall-clock time since the last acknowledgement, so the retry budget (up to ten attempts at up to 1m each) plus the operator response to a stall has to fit inside it. The margin check becomes time since last acknowledgement against the deadline. Empty batches in quiet periods must still acknowledge.
  • The indexed-object check runs only at startup. An index created later on the subscribing cluster makes the next resume fail AS OF against the index's 1s since with the 22000 timestamp-selection error, which the SDK classifies as HistoryLost and stalls on. Re-check for an index before classifying history loss. Under durable subscriptions the attach reads storage, so IndexedTarget only applies to the pre-durable path.
  • Blue/green cutover is a name swap. ALTER SCHEMA ... SWAP leaves ids untouched, so the running stream keeps reading the decommissioned collection until teardown drops it. A Transient reconnect in between re-resolves the name to the new object with the same columns, and the fingerprint matches, so the SDK switches objects silently. Include the object id in the fingerprint. With durable subscriptions the teardown DROP needs CASCADE, and refollow means creating a new subscription START AT the checkpoint, which needs the new object's history to cover it.
  • Multi-view recovery after a partial expiry. "Resume every member at the stored cut" fails once one member has expired. #38468 describes the recovery: reset the expired members, re-establish the cut at the largest reset position, and buffer every stream to it. Also reject members outside the EpochMilliseconds timeline, since the cut is unsound across timelines.
  • R8 does not hold on the durable path. next() runs only after the target commit, so the backlog accumulates in environmentd while the target writes. During catch-up this can exceed subscribe_max_buffered_bytes, retire the subscribe, resume, and repeat. Each cycle still advances, but it pays a dataflow restart every time. Use the bounded fetch loop here too, or back off on FellBehind.
  • The security section understates the privileges. The Materialize-table checkpoint store needs INSERT and UPDATE, and durable subscriptions need privileges on the subscription object and someone to create it. #38468 expects subscriptions to be provisioned per logical consumer, because creation is a catalog transaction, so the SDK should not create one on each start.
  • The Materialize-table checkpoint store is not atomic with the target. State that it offers at_least_once only.

Nits

  • mz_timestamp is a u64. Current values fit a JS number, but the Node package should expose a BigInt.
  • The UP TO workaround for SQL-528 should be gated on server version, so it does not release data twice once the fix lands.
  • ENVELOPE UPSERT emits key_violation rows, which the problem table cites as the mz-redis-sync crash, but the sink module does not say how the SDK handles them. Mapping them to reject or a stall would close that.
  • "Out of scope: ENVELOPE DEBEZIUM" is fine for now, but #38468 supports both envelopes across a resume, so it can move in once the flag lifts.

For what it is worth, the per-timestamp ordinal in the idempotency key holds: process_response (src/adapter/src/active_compute_sink.rs:266-279) merges and consolidates each batch across workers, and a timestamp never straddles a batch, so each timestamp's output is canonical across a resume.

@bobbyiliev
bobbyiliev requested a review from martykulma October 8, 2026 14:36
@bobbyiliev

Copy link
Copy Markdown
Contributor Author

Thanks, these all held up in the code, and I applied them along with fixes from a second review pass. turbopuffer now uses tombstones with conditional writes (swept by write time), durable mode is limited to storage collections, the attach keeps AS OF F - 1, the error section is honest about the shared SQLSTATEs, and a new Object identity section uses the storage shard so a replacement resumes and a name swap stops. Could you look at the turbopuffer and Object identity sections again, since they now carry the correctness argument?

@bobbyiliev
bobbyiliev requested a review from antiguru October 8, 2026 14:42

@antiguru antiguru left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Posted by Claude Code on behalf of @antiguru.

Re-checked 7a8343c with a focus on the turbopuffer and Object identity sections. The earlier points are addressed. The remaining items are below.

Object identity

A replacement takes its retention from the replacement, not the original. MaterializedView::apply_replacement (src/catalog/src/memory/objects.rs:1504-1566) copies custom_logical_compaction_window and refresh_schedule from the replacement. The shard and catalog id parts of the table hold: new_entry.id = replacement_id in transact.rs, and the new GlobalId is appended to the existing collections on the same shard. But a replacement created without RETAIN HISTORY drops the shard to the one-second default at apply, so the "resume" row then fails with HistoryLost. Deploy tooling has to carry RETAIN HISTORY onto the replacement, and the SDK should run its margin check again after it sees a new catalog id.

The name lookup and the SUBSCRIBE can race. The SDK resolves the name, compares the fingerprint, and then subscribes by name. A swap between those two steps subscribes to the new object after the check has passed. mz_internal.mz_subscriptions.referenced_object_ids lists the catalog ids a running subscribe actually reads, so checking it right after DECLARE against the fingerprint closes the window.

Store the new catalog id after a replacement. Otherwise every later resume compares against the pre-replacement id and goes down the replacement path again, including the margin re-check above.

turbopuffer

The generation sweep can move a timestamp backwards. The patch filter is generation < t_s AND deleted = false and it sets mz_timestamp = t_s. A stale worker from before the re-snapshot can have written a key at ts > t_s, and by determinism that content is correct. The snapshot's upsert of that key is then skipped, since the stored timestamp is newer, so the document keeps its old generation, and the sweep tombstones it at t_s. The key reads as deleted until the live worker reaches ts, and in the meantime the invariant that "the conditions only ever let a newer timestamp win" is broken. Adding mz_timestamp < t_s to the filter keeps the sweep from moving any document backwards.

Make the frontier commit monotone as well. The commit patch is conditional only on the epoch. A retried commit for F1 that times out but lands after the commit for F2 moves the frontier back. That is safe on its own, because replaying from F1 reapplies the deletes. It is worth adding a condition that the stored frontier is lower anyway, because the tombstone sweep's safety argument reads the frontier.

The grace-period argument for stale workers is sound as written. The doc says the risk is bounded rather than eliminated, which is the honest framing.

Durable subscriptions

#38468 now exempts a subscription that has caught up (H + 1 at the write frontier) from the deadline (9b39a88). Its last-acknowledgement time advances while it has nothing left to acknowledge, so a caught-up sink on a REFRESH EVERY view, a paused source, or a cluster with no replicas no longer expires. Two changes follow:

  • In "History loss and retention margin", drop "the longest time the object's frontier can stand still".
  • Workstream item 4 can go.

The deadline still has to exceed the retry budget plus an operator's response to a stall. Renewing on an acknowledgement at the current position was rejected, because a sink stuck behind a poisoned batch could then hold history forever.

@bobbyiliev

Copy link
Copy Markdown
Contributor Author

Update: we will build the Rust package first, with Python and Node as stretch goals, so the October turbopuffer sink will be written in Rust against turbopuffer's HTTP API. I updated the delivery plan and open question 2 to match, and the Linear project will follow once this doc is agreed. Does anyone see a problem with a Rust turbopuffer sink for October?

@bobbyiliev

Copy link
Copy Markdown
Contributor Author

Thanks, all of these check out and I applied them. The replacement path re-runs the margin check and stores the new catalog id, the SDK now subscribes by the catalog id it checked (SUBSCRIBE [u123 AS ...]) so a swap cannot slip in between the lookup and the subscribe, the generation sweep only touches documents older than the snapshot, and frontier commits only move forward. I also dropped the frontier-standing-still rule and workstream item 4 now that #38468 exempts caught-up subscriptions. A name swap now re-snapshots from the new object, into a new namespace for targets without transactions, since resuming it would keep keys the new object never had. Can the SDK rely on the id-based reference syntax in SUBSCRIBE (open question 13), and does the swap handling match what you would expect from deploy tooling?

@bobbyiliev
bobbyiliev requested a review from antiguru October 8, 2026 18:09
@jubrad

jubrad commented Oct 8, 2026

Copy link
Copy Markdown
Member

The SDK scope is wider than what #38468 accepts. #38468 targets storage collections only (tables, materialized views, sources), rejects views and indexes, and rejects temporal filters. The SDK accepts any object, and a non-materialized view is inlined, so its resume needs history on every input and reads their snapshots, as the "Snapshot elision" section notes. Accepting views or mz_now() filters now turns the later move into a narrowing, which is the API break the section wants to avoid. R

This call out seems accurate, but I think we might want to differentiate between what we could use for subscribes and what we could expect for something that could power a sync. It's possible that when we're powering a sink like turbopuffe, postgres, etc... we'd actually want to res on MV's (or tables), something more directly backed by a persist shard, and further we could subscribe <object> rather than via a query that references the object. I believe there's a lot of room for optimization here, and it matches our requirements for kafka sinks.

@bobbyiliev

Copy link
Copy Markdown
Contributor Author

@jubrad Agreed. I narrowed the sink path to match Kafka sinks: a sink reads one table, source, or materialized view with SUBSCRIBE <object>, with an optional envelope and no projection or filter, so a resume never reads the snapshot. Live streams that never resume still accept any query. Does that split match what you had in mind?

@antiguru antiguru left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

The second-round fixes check out. Subscribing by id resolves by CatalogItemId and uses the bracketed name for display only (src/sql/src/names.rs:1585), so a swap cannot redirect the stream. A FETCH without a timeout is WaitOnce, and it returns empty only when the channel closes, so the end-of-stream check is sound. The commit contract with an expected frontier, the mz_timestamp < t_s sweep clause, and the counter generation all close the holes from the last round. Narrowing sinks to SUBSCRIBE <object> matches CREATE SINK and keeps a later widening compatible.

Remaining points, none blocking:

  • Reader switch after a swap. refollow writes turbopuffer data into new namespaces and "switches readers", but the doc does not say how. Readers need a pointer to the active namespace, for example a field in the checkpoint document or a name derived from the generation. That indirection belongs in the sink's reader contract. The "next cleanup" of old namespaces needs a durable list of them, and with several namespaces the switch itself is not atomic.
  • InvalidId is XX000 at DECLARE. That is the same SQLSTATE the doc says would be misread as StreamPoisoned. Consider stating the rule by statement rather than by code: any DECLARE failure goes back to the name lookup, and only FETCH failures count as StreamPoisoned.
  • nit: While the re-snapshot runs, readers of the old namespaces see data frozen at the swap. One sentence in the sink docs would cover it.

Posted by Claude Code.

@bobbyiliev

Copy link
Copy Markdown
Contributor Author

@antiguru Thanks. I applied all three. The checkpoint document now records the active namespace incarnation and the retired ones, readers derive namespace names from it with a small helper, and switching readers is one patch to that document, so it is atomic across namespaces. Retired incarnations stay on the list, so cleanup also catches a namespace a stale write recreated. Errors are now classified by statement: any DECLARE failure goes back to the name lookup, and only FETCH failures count as StreamPoisoned. The sink docs will also say that old namespaces show data frozen at the swap until the switch.

@jubrad jubrad left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

One potential correctness issue caught by claude. I have a few questions, but otherwise, this looks relaly good — 1 blocker, 2 considerations.

written by claude on behalf of @jubrad

because conditions are evaluated per document, so its data writes cannot be
conditioned on the epoch in a checkpoint.

The sink will write tombstones in place of deletes. A delete will upsert the

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

blocking — Keeping tombstones forever still leaves a hole during history-loss recovery: worker A pauses before inserting k@90, k is deleted upstream at 95, and worker B re-snapshots at 100 and finishes its sweep while k has never existed in the target. When fenced A resumes, the missing-document upsert is unconditional, and no tombstone or post-snapshot delete removes k, so the namespace stays wrong. Please extend the new-namespace isolation to these re-snapshots too, or provide fencing that covers absent keys, and include this case in the fencing test.

written by claude on behalf of @jubrad

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Good catch, you're right: tombstones only protect keys the namespace has seen. Every re-snapshot on a target without transactions now goes into a new namespace incarnation, after history loss as well as after a swap, and the fencing test includes your exact case.

Comment on lines +62 to +67
`RETAIN HISTORY` sets `since = upper - window`. After an environment outage or
upgrade the upper catches up and the window slides with it, so a one-hour
window can expire during a one-hour outage. A `REFRESH EVERY` view's upper
jumps to the next refresh, so a sink that is down across a refresh can lose
history however short the downtime. Retention belongs to the view's owner, and
any `ALTER` changes every sink's margin without telling it.

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

It is worth calling out that retain history is not GA and I believe there are mixed feelings about it internally. It does seem like the right thing to rely on here, but we should also understand this will encourage customers to ask for a that feature. @antiguru, you've approved this PR so I assume this is reasonable to you.

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Good point. My understanding is that this design is compatible with whatever we're going to build to replace retain history, but it certainly raises the question of what do we do until then. GA'ing retain history is not an option, because of the rough edges around it, so it'll need to stay focused for the moment.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Agreed. The doc now says it's in private preview and that the SDK depends on it until durable subscriptions land, which is one more reason to get those in.

| `AS OF` below the readable frontier | `22000` (`DATA_EXCEPTION`, shared) | "could not find a valid timestamp for the query" |
| Subscribed object or cluster dropped | `42704` | "relation 'x' was dropped" and similar |
| Client fell behind | `53200` (shared with the adapter's result-size error) | `SubscribeFellBehind` |
| Result over `max_result_size` | `XX000` (`INTERNAL_ERROR`, shared) | "total result exceeds max size of ..." |

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

consider — Max result size really is going to be an issue for snapshots. There has been some work to support large result sizes, but I do not believe they work on persist, there is other poc work that would support larger results and move snapshot processing off envd. If subscribes become a serious mechanism for sink creation I believe we will need to resolve this.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Agreed, this is the biggest server-side risk. I linked #38789 in the server-side chunking item. Until that lands, the doc states the limit and gives it its own typed error.

| --- | --- | --- |
| R1 | Exactly-once state in the target, across restarts and crashes | PRD |
| R2 | Back off and retry while the target is unavailable | PRD |
| R3 | Dead-letter a batch after too many retries, and reject single rows the target cannot take | PRD |

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

@antiguru recently proposed (ignore errors) for subscribes, what would be a truly cool feature here is to be able to subscribe to both the Ok stream and the Err stream. - feel free to ignore this comment for this PR, mostly just thinking out loud.

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

There's a bigger problem: Errors taint subscribes, and it'll be impossible to get a consistent view without re-snapshotting after an error. This limitation we can't work around today.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Love this. I added it to the non-poisoning errors item: a subscribe that emitted both data and errors would let a sink dead-letter bad rows instead of stopping.

Comment on lines +213 to +214

This document proposes **Materialize SDK**, with `subscribe` and `sink` as its

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

@maheshwarip, @sjwiesman, it would be good to have a thumbs up from either of you on this point.

Comment on lines +336 to +337
### Live streams

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This is the correct behavior. But also where moving subscribes to being process by clusters and more direct reads from persist seems like a big win.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Agreed. I tied it into the same workstream item as #38789, since sinks only read storage-backed objects now.

Comment on lines +385 to +386
stored epoch equals the worker's and the stored frontier equals `expected`, the
frontier the batch started from. A stale worker gets a typed `Fenced` error.

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

consider — Do we want this to be a "name" or an "id"

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

It's a user-chosen id for the deployment, like a consumer group id, not the object's name. I renamed it to sink_id to make that clear, which is also the term mz-sink-sdk uses.

Comment on lines +420 to +421
When delivery fails, a policy will return `retry` (back off and deliver the
same batch again), `skip` (the user has dead-lettered it, advance past it), or

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Do we need a commit size here as well?

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Yes, added max_commit_bytes. Every timestamp in a batch is closed on its own, so the SDK can split a large batch at timestamp boundaries and commit at each split. One timestamp over the limit, like a snapshot, goes through a generation or a new incarnation.

Comment on lines +498 to +500
users, so deferred question 2 asks whether the SDK can rely on it.

A replacement keeps the storage shard and gives the view a new catalog id

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Do we want to provide a knob here that allows us to overwirte when a the id of the object we're subscribing to changes, or maybe just if the columns change, or should we let the consumers handle this by creating new subscribes / sinks. For instnace we could ignore new columns, or allow some transform between subscribe and sink modules.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Good idea. The default stays stop, but I added a deferred question on a policy for added columns, either dropping them on the client or passing rows through a transform between the subscribe and sink modules.

Comment on lines +546 to +547

A multi-view subscription will hold one connection per member. Any member error

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

If we have ~10 objects we want to subscribe to through this, and each hits their max_result_size, I think we'd still be very likely to put environmentd at risk of OOMing without batching snapshots. Nothing to fix here, the method is right, but it is something we're going to have to think about.

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Elsewhere in this doc we describe trying to push subscribes to materialized views so we could later back this directly from persist or via clusterd -- is that the idea?

Will serving not out of environmentd become a blocker for rollout?

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Agreed. The doc now says a multi-view sink multiplies this, and the SDK will start members one at a time at the shared AS OF, so only one snapshot is in environmentd at once. The real fix is still the server-side work.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Yes, that's the idea: sinks only read storage-backed objects, so serving from persist or clusterd later needs no API change. Not a blocker for the October prototype, but the doc now makes snapshot size a GA requirement.

@antiguru

antiguru commented Oct 9, 2026

Copy link
Copy Markdown
Member

(My approval means "I don't have more immediate comments," but please wait for others, too.)

| R10 | Several languages with identical behavior, checked by machines | Design |
| R11 | Sink authors can test without a live environment | Design |

R4 is met in full by targets that can write the whole cut atomically, such as

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Is "cut" general terminology? What does it mean?

### Naming

This document proposes **Materialize SDK**, with `subscribe` and `sink` as its
first modules. This is the decision we most want reviewers to push on.

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

nit: we can probably delete this line

client. The scope rules that out: drivers stay the way to run queries, and the
SDK only adds modules for things drivers get wrong.

Package names will be `materialize-sdk` on PyPI, `@materializeinc/sdk` on npm,

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I share the concern here that this communicates a general sdk. materialize-sink-sdk and @materializeinc/sink-sdk seem better to me at first glance.

Just to clarify, the reason you want to go with the more general name here is to more accurately name a library that covers both sinks and general subscribe utilities, right? Would we add more things (i.e. a source SDK) to this package too in the future?

The SDK will refuse an indexed object at startup. It will check whether the
subscribed object has an index on the subscribing cluster and fail with an error
that names the remedy: a dedicated subscribe cluster with no index on the object.
An index created later makes the next resume fail with the timestamp-selection

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

What's the purpose of this behavior? Is this because retain history on an index is expensive?

stored frontier already past `expected`, rolls back, and the SDK reads the
checkpoint and treats the batch as committed instead of applying it twice.

Every durable subscription will require an explicit `sink_id`, which is its

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

These need to be unique to avoid fencing/conflicts right? Do we recommend UUIDs?

and `mz_internal.mz_object_transitive_dependencies`.

Each member's snapshot is collected in `environmentd` up to `max_result_size`
(see "Buffering limits"), so N members starting together can hold N times that

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Any concerns with causing environmentd to oom with too many tables being sinked?

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Yes, real concern. Both limits are per subscribe with no total, so N sinks can hold N times max_result_size during snapshots. I added this to the doc, and GA now needs server-side chunking or a documented per-sink size limit.

Comment on lines +546 to +547

A multi-view subscription will hold one connection per member. Any member error

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Elsewhere in this doc we describe trying to push subscribes to materialized views so we could later back this directly from persist or via clusterd -- is that the idea?

Will serving not out of environmentd become a blocker for rollout?

has no Kafka, so the existing Kafka-based sink does not fit there. The sink will
write only to namespaces it creates, so every document carries `mz_timestamp`.

The October sink will be written in Rust on the Rust package and will call

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

nit: October timeline detail feels a little random here.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Fair, I moved it out. The timeline now only lives in the delivery plan.

conditioned on the epoch in a checkpoint.

The sink will write tombstones in place of deletes. A delete will upsert the
key's document with `deleted = true`, the delete's `mz_timestamp`, a wall-clock

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

will you keep track of multiplicity for multiple rows/documents with the same key and only mark as deleted when that hits 0?

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Turns out no, nice catch. The upsert envelope only looks at one timestamp at a time, so if the key isn't unique, a delete can wipe a key that still has rows. I updated the doc to require a unique key like CREATE SINK does, and added a server item to check it in SUBSCRIBE.

Comment thread doc/developer/design/20260706_materialize_sdk.md Outdated

This branch has not been deployed

No deployments
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

5 participants