connectors/nats registers the nats sink (alpha). It publishes one JSON
message per input row to an existing JetStream stream. Both bounded runs and
continuous execution use the same bounded asynchronous publisher. Each record
remains an individual NATS message.
Configuration
Source and sink connections share these authentication settings. Choose only one
method: username/password, JWT/NKey seed, or token. Select
none for an unauthenticated server. Selected authentication fields are required.
An omitted auth_mode means no authentication; credentials require an explicit mode. JWT authentication
uses the seed to sign the server challenge; a JWT alone is insufficient.
Username/password and JWT/seed require no files mounted on workers. Passwords,
JWTs, seeds, and tokens use Filament’s secret storage.
Both stream and subject are templates. The {resource} placeholder, or
its Go template spelling {{resource}}, is replaced with the destination
resource name, so one configuration can route
every resource to its own subject, or its own stream, or both:
For example, provision a stream named
WAREHOUSE_EVENTS capturing
warehouse.>, then configure:
create_stream enabled, a missing stream is created with file storage,
server defaults for retention and limits, and one subject filter derived from
the subject template: events.{resource} yields events.>, {resource}.events
yields *.events, a fixed subject yields itself, and a per-resource stream
captures that resource’s exact subject. Anything beyond that, such as retention
or replicas, belongs in server-side stream management. Disable creation on
production pipelines where streams are provisioned deliberately.
The sink never modifies an existing stream. It must enable publish
acknowledgments and its subject filter must capture the expanded subject. A shared stream is checked at open using * for the resource, before
any write. A per-resource stream is resolved when that resource’s first batch
is routed, and a missing or non-capturing stream fails the batch before any
message is sent. Note that a filter of events.> matches events.orders but
not the bare subject events. Each publish carries Nats-Expected-Stream, so a
subject routed to another stream fails rather than silently writing there.
Workers need permission to inspect the stream, publish to its subjects, and
receive publish acknowledgments. Connection testing also requires JetStream
account-info access.
Write modes and message format
Supported bounded ingestion types arefull_append, incremental_append, and
cdc_append. Continuous execution supports append for rows, messages, and
change events. Replace, upsert, merge, and destination deletion are unsupported.
Each message body is a JSON object containing the Arrow row’s columns, without
a trailing newline. The shared JSON encoder renders bytes as base64 and temporal
values as text. Every message has:
Content-Type: application/jsonFilament-Operation: insert,update, ordelete
Apply, messages are submitted in row order over one reused
connection. The publisher pauses whenever either outstanding limit is reached.
SDK acknowledgment callbacks can arrive in any order and release capacity for
more messages. Shared completion tracking avoids one waiting goroutine per message.
Apply waits for every acknowledgment before returning success. Set
max_in_flight to 1 for sequential publication.
The byte limit includes encoded payloads, headers, and subjects, excluding reply
subjects and protocol framing. An individual message must fit within this limit;
otherwise the entire batch is rejected before any messages are sent. The limit
bounds outstanding publications, not total process memory: the complete encoded
batch is also held until Apply returns. Use the run’s BatchMaxRows and
BatchMaxBytes options to control input batch size; these do not combine records
into a single NATS message.
Concurrent Apply calls remain serialized, so the configured limits apply to
the whole sink session. No global ordering is promised for concurrently produced
batches. Raw-message aggregation and connection pooling are not required for
asynchronous publication and are not implemented.
Durability, replay, and failures
Apply succeeds only after JetStream acknowledges every message. Durability
therefore follows successful Apply, using the destination stream’s configured
storage, replication, and retention settings. Commit finalizes a bounded run;
CommitEpoch returns per-resource totals for an already acknowledged epoch.
Continuous source progress advances through Filament’s normal epoch certification
flow, after destination success.
Messages become visible individually. Neither batches nor epochs are atomic.
An abort stops work and clears local accounting; it never deletes acknowledged
messages. If a batch fails partway through, any accepted messages remain visible, including
messages submitted after the one whose acknowledgment failed.
When a row includes a nonempty _filament_event_id, the sink derives a stable
Nats-Msg-Id from the event ID, tenant, pipeline, sink connection, stream,
subject, and destination resource. The ID remains stable across attempts and
rebatching. Duplicate acknowledgments count as successfully handled input rows.
Rows without event IDs are published without deduplication IDs.
Delivery is at least once. JetStream duplicate suppression lasts only for the
stream’s configured duplicate window. Retries outside that window, or retries
of ordinary rows without event IDs, can produce duplicates. Reusing an event ID
for different content can suppress the later publication.
A lost acknowledgment leaves the write outcome uncertain. The sink stops
publishing, rejects the epoch commit, and reports unproven destination quiescence
on session closure. Any other session failure is also retained and returned by
CloseSession; an epoch still active at closure is dropped like an abort.
A bounded Abort does not repeat a failure that Apply already returned. On failure or cancellation, it stops submissions, closes the connection,
detaches completion tracking, and cleans up SDK publications before returning.
Late callbacks cannot complete or fail a subsequent batch.
It does not reconnect or retry publishes in the background.
It advertises no destination owner fencing, isolated epochs, or server-enforced
bound on residual effects; uncertain shutdown follows Filament’s conservative
attempt recovery path.
Write receipts report payload bytes and an encoded CRC32C over concatenated JSON
message bodies. This integrity evidence excludes NATS headers and protocol framing.