Repository navigation
design: Materialize SDK for SUBSCRIBE consumption and custom sinks - #37463
bobbyiliev wants to merge 19 commits into
Conversation
|
I would suggest some sort of sdk agnostic test suite so we can validate the implementations. |
|
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? |
| ### 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 | ||
| ``` |
There was a problem hiding this comment.
@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.
antiguru
left a comment
There was a problem hiding this comment.
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 WITHINis 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 astallhas 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 OFagainst the index's 1ssincewith the22000timestamp-selection error, which the SDK classifies asHistoryLostand stalls on. Re-check for an index before classifying history loss. Under durable subscriptions the attach reads storage, soIndexedTargetonly applies to the pre-durable path. - Blue/green cutover is a name swap.
ALTER SCHEMA ... SWAPleaves ids untouched, so the running stream keeps reading the decommissioned collection until teardown drops it. ATransientreconnect 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 teardownDROPneedsCASCADE, andrefollowmeans creating a new subscriptionSTART ATthe 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
EpochMillisecondstimeline, 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 inenvironmentdwhile the target writes. During catch-up this can exceedsubscribe_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 onFellBehind. - The security section understates the privileges. The Materialize-table checkpoint store needs
INSERTandUPDATE, 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_onceonly.
Nits
mz_timestampis a u64. Current values fit a JSnumber, but the Node package should expose aBigInt.- The
UP TOworkaround for SQL-528 should be gated on server version, so it does not release data twice once the fix lands. ENVELOPE UPSERTemitskey_violationrows, 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 torejector 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.
|
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 |
antiguru
left a comment
There was a problem hiding this comment.
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.
|
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? |
|
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 ( |
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 |
|
@jubrad Agreed. I narrowed the sink path to match Kafka sinks: a sink reads one table, source, or materialized view with |
antiguru
left a comment
There was a problem hiding this comment.
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.
refollowwrites 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. InvalidIdis XX000 atDECLARE. That is the same SQLSTATE the doc says would be misread asStreamPoisoned. Consider stating the rule by statement rather than by code: anyDECLAREfailure goes back to the name lookup, and onlyFETCHfailures count asStreamPoisoned.- 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.
|
@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 |
| 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 |
There was a problem hiding this comment.
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
There was a problem hiding this comment.
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.
| `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. |
There was a problem hiding this comment.
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.
There was a problem hiding this comment.
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.
There was a problem hiding this comment.
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 ..." | |
There was a problem hiding this comment.
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.
There was a problem hiding this comment.
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 | |
There was a problem hiding this comment.
@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.
There was a problem hiding this comment.
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.
There was a problem hiding this comment.
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.
|
|
||
| This document proposes **Materialize SDK**, with `subscribe` and `sink` as its |
There was a problem hiding this comment.
@maheshwarip, @sjwiesman, it would be good to have a thumbs up from either of you on this point.
| ### Live streams | ||
|
|
There was a problem hiding this comment.
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.
There was a problem hiding this comment.
Agreed. I tied it into the same workstream item as #38789, since sinks only read storage-backed objects now.
| 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. |
There was a problem hiding this comment.
consider — Do we want this to be a "name" or an "id"
There was a problem hiding this comment.
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.
| 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 |
There was a problem hiding this comment.
Do we need a commit size here as well?
There was a problem hiding this comment.
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.
| 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 |
There was a problem hiding this comment.
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.
There was a problem hiding this comment.
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.
|
|
||
| A multi-view subscription will hold one connection per member. Any member error |
There was a problem hiding this comment.
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.
There was a problem hiding this comment.
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?
There was a problem hiding this comment.
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.
There was a problem hiding this comment.
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.
|
(My approval means "I don't have more immediate comments," but please wait for others, too.) |
…e, and mz-sink-sdk
| | 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 |
There was a problem hiding this comment.
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. |
There was a problem hiding this comment.
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, |
There was a problem hiding this comment.
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 |
There was a problem hiding this comment.
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 |
There was a problem hiding this comment.
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 |
There was a problem hiding this comment.
Any concerns with causing environmentd to oom with too many tables being sinked?
There was a problem hiding this comment.
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.
|
|
||
| A multi-view subscription will hold one connection per member. Any member error |
There was a problem hiding this comment.
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 |
There was a problem hiding this comment.
nit: October timeline detail feels a little random here.
There was a problem hiding this comment.
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 |
There was a problem hiding this comment.
will you keep track of multiplicity for multiple rows/documents with the same key and only mark as deleted when that hits 0?
There was a problem hiding this comment.
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.
…ostgres sink before GA
Proposes the Materialize SDK: a Rust protocol core with no I/O, compiled into thin per-language packages that use their native drivers, with
subscribeandsinkmodules. 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/Q6193A4KskJpfWgUuYNY6vTL;DR