Skip to main content
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

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.