> ## Documentation Index
> Fetch the complete documentation index at: https://filament.getgalaxy.io/llms.txt
> Use this file to discover all available pages before exploring further.

# Kafka

> Consume Kafka message envelopes with certified partition checkpoints

`connectors/kafka` registers an **alpha** bounded and continuous message source. Topics are
resources, and topic partitions have independent next-offset checkpoints. The
first implementation uses direct partition assignment to one Filament owner.
It does not join a consumer group or write Kafka group offsets.

## Configuration

```json theme={"theme":{"light":"github-light-default","dark":"github-dark-default"}}
{
  "brokers": ["localhost:9092"],
  "topics": ["orders", "customers"],
  "start_position": "earliest"
}
```

`topics` supplies default resources when none are explicitly selected. Discovery
lists accessible noninternal topics. `start_position` is `earliest` (default) or
`latest` for partitions without a checkpoint. Existing certified offsets always
win. Initial frontiers, including idle partitions, are emitted for certification
before reading records. Until that first certificate is durable, an unsuccessful
startup may resolve the initial frontier again.

| Connection field | Description |
| - | - |
| `brokers` | Required list of bootstrap `host:port` addresses |
| `client_id` | Kafka client identifier; default `filament` |
| `tls_enabled` | Enable certificate-verified TLS |
| `tls_ca_pem` | Optional pasted PEM CA certificates; stored as a secret |
| `tls_cert_pem`, `tls_key_pem` | Optional pasted PEM client certificate chain and matching private key; both stored as secrets |
| `sasl_mechanism` | Authentication dropdown: `none` (default), `PLAIN`, `SCRAM-SHA-256`, or `SCRAM-SHA-512` |
| `sasl_username`, `sasl_password` | SASL credentials; password is a secret |

No files need to be mounted on workers. Enable TLS to use the system trust store;
provide `tls_ca_pem` only for a private CA. For mutual TLS, paste both the client
certificate chain and private key. PEM inputs preserve multiline pastes and use
Filament's secret storage. Certificate fields require `tls_enabled=true`.

Selecting a SASL mechanism shows required username/password fields. Use TLS with
PLAIN authentication. The integration test brokers use neither TLS nor authentication.
Existing `tls_ca_file`, `tls_cert_file`, and `tls_key_file` settings are rejected:
replace them with the corresponding `_pem` fields containing the actual contents.

Connection testing authenticates and
requests metadata without joining a consumer group or publishing records.
Workers need topic metadata and read access; consumer-group permissions are not
required for this direct-assignment profile.

## Integration tests

Run `just test-integration` from the repository root. The Kafka tests provision
Kafka through Testcontainers, create topics, publish test records,
and clean up their containers. No separately running broker or publisher is needed.

## Bounded reads

Choose bounded execution with **Full** reads and an append or replace destination
write policy. Each run captures the retained beginning and the read-committed end
(last stable offset) of every selected partition before extraction. It reads that
fixed window and finishes, even if producers keep publishing. Transactions still
open when the window is captured and messages appended afterward are excluded.
This is a set of partition windows, not an atomic snapshot across topics.

Bounded full reads always begin at the retained start; `start_position` applies
only to continuous execution. Each new run rereads the retained log, so append
runs can duplicate earlier output. Bounded runs do not advance continuous
checkpoints or Kafka consumer-group offsets and do not support incremental or
checkpointed resume. A retry captures a new window. Compaction may remove records
within the window; the run reads the records still retained when fetched.

Resource planning infers payload columns before typed destination tables are
prepared. The first retained message fixes each topic's schema, including the
same tombstone/raw-payload behavior described below. Empty topics finish without
waiting for new records. Retention loss or topic recreation affecting unread
partitions fails the run instead of silently resetting offsets. Partitions added
after planning are outside the captured window.

## Schema and checkpoints

Kafka and NATS sources decode JSON object payloads into top-level, source-owned
columns. For example, `{"event_id":"order-1","amount":100}` produces `event_id`
and `amount` columns alongside Filament's `_filament_event_id` and checkpoint
identity. There is no `_filament_payload` column. Nested objects and arrays remain
JSON values; numbers retain their JSON precision.

Kafka adds `topic`, `partition`, `offset`, and `tombstone` transport columns; NATS
adds `subject`. Both retain `_filament_event_id`, `_filament_event_identity`,
`_filament_event_ts`, `_filament_headers`, and `_filament_key` event metadata.
Payload names must not collide with transport fields (case-insensitive) or use
the reserved `_filament_` prefix. Duplicate JSON keys are rejected.

The first message fixes the column schema for each resource during a run.
Missing fields become null; new fields and incompatible string/bool type changes
fail before progress is certified. JSON-typed fields accept any JSON value.
Discovery exposes transport/event metadata; payload columns appear on the first
read. Non-object JSON and binary messages use a source-owned binary `payload`
column. Mixing raw and object payloads in a resource requires a new compatible
schema and is rejected within a run.

Kafka tombstones set `tombstone=true` and leave decoded payload columns null.
If the first message is a tombstone, the resource uses the raw `payload` schema.
Keys, timestamps, and ordered duplicate headers are preserved. Existing output
tables may need a new destination or schema migration for this column layout.

The source uses `read_committed` isolation. A record at offset `n` has progress
position `n+1`, encoded by `kafka.partition.next_offset` version 0. Offset zero is
a valid initial frontier. Domain incarnation combines the stable Filament source
connection ID, Kafka cluster ID, and topic ID. Brokers without stable cluster and
topic IDs are rejected. Recreated topics cannot reuse an old checkpoint.

Filament certificates own durable progress. `Acknowledge` checks exact coverage
and current Filament authority; there is no broker acknowledgment to advance.
The source refuses another read while coverage is pending. Retries resume from
certified offsets, including when the previous worker stopped before its local
acknowledgment completed.

Offset gaps from compaction or transactional records are valid. Replay covers
retained records, not every event ever produced. Positions outside the retained
range fail instead of resetting to the beginning or end. The connector also
rejects topic recreation and partition membership changes during a session.
These checks cannot reconstruct records lost through broker log corruption or
unclean election followed by regrowth; configure broker durability accordingly.

## Initial release limits

The source supports bounded full reads and continuous execution. Consumer groups,
bounded incremental/resumable reads, distributed partition handoff, and Schema
Registry decoding remain follow-up work.

The continuous source processes one record per epoch initially and advertises no end-to-end
ordering guarantee. Normal fetches use a 4 MiB request limit, 1 MiB partition
limit, and one concurrent fetch. Kafka may return an oversized first record, so
these are not hard total-memory limits. Topic membership changes require stopping
and restarting the stream; saved domains must remain present.

## Tested profiles

The connector integration suite tests Kafka **4.0.2** using single-broker plaintext fixtures. It checks tombstones, source resume, retention loss, topic
recreation, and live partition changes. Bounded tests cover fixed windows, JSON
schemas, limits, empty partitions, transactions, and retention loss. TLS/SASL configuration has unit coverage;
secured-cluster and multi-broker fault tests remain follow-up work.


This documentation is built and hosted on [Mintlify](https://mintlify.com), a developer documentation platform.