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
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.
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
Runjust 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.