Skip to main content
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:
Expanded subjects must be concrete: wildcards, whitespace, and empty subject tokens are rejected, and stream names cannot contain dots or wildcards. Resource names are not sanitized. Use destination resource mapping when source names are unsuitable for NATS subjects or stream names. With 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 are full_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/json
  • Filament-Operation: insert, update, or delete
Updates and deletes append change-history messages; they do not modify previous messages. Message-source envelope columns remain in the JSON object. Raw payload and header passthrough are not implemented. The sink validates and serializes the complete batch before publishing. Messages exceeding the configured stream or server size limit fail; they are not split. Within each 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.