Kafka tables
Use table({ engine: 'Kafka', ... }) to manage a Kafka queue together with its
materialized view and storage table. Kafka tables have columns and engine settings,
but no primaryKey, orderBy, partitioning, TTL, indexes, or projections.
import { schema, table, materializedView } from '@chkit/core'
const columns = [ { name: 'id', type: 'String' }, { name: 'event_time', type: 'DateTime64(3)' },]
const queue = table({ database: 'analytics', name: 'events_queue', engine: 'Kafka', columns, settings: { kafka_broker_list: 'kafka01:9092,kafka02:9092', kafka_topic_list: 'events', kafka_group_name: 'analytics_events', kafka_format: 'JSONEachRow', kafka_num_consumers: 1, input_format_skip_unknown_fields: true, },})
const events = table({ database: 'analytics', name: 'events', engine: 'MergeTree', columns, primaryKey: ['event_time', 'id'], orderBy: ['event_time', 'id'], partitionBy: 'toYYYYMM(event_time)',})
const consumer = materializedView({ database: 'analytics', name: 'events_consumer', to: { database: 'analytics', name: 'events' }, as: 'SELECT id, event_time FROM analytics.events_queue',})
export default schema(queue, events, consumer)The Python DSL supports the same pipeline natively:
from chkit import materialized_view, schema, table
columns = [{"name": "id", "type": "String"}, {"name": "event_time", "type": "DateTime64(3)"}]queue = table( database="analytics", name="events_queue", engine="Kafka", columns=columns, settings={ "kafka_broker_list": "kafka01:9092,kafka02:9092", "kafka_topic_list": "events", "kafka_group_name": "analytics_events", "kafka_format": "JSONEachRow", "kafka_num_consumers": 1, "input_format_skip_unknown_fields": True, },)events = table( database="analytics", name="events", engine="MergeTree", columns=columns, primary_key=["event_time", "id"], order_by=["event_time", "id"], partition_by="toYYYYMM(event_time)",)consumer = materialized_view( database="analytics", name="events_consumer", to={"database": "analytics", "name": "events"}, as_="SELECT id, event_time FROM analytics.events_queue",)definitions = schema(queue, events, consumer)Run chkit generate, review the SQL, then chkit migrate --apply. Tables are
created before materialized views. Attaching the view starts background consumption.
The Kafka engine must be available on the target server, which must be able to
reach the brokers. This workflow is tested on self-hosted ClickHouse 25.3 and 26.3.
Settings and validation
Section titled “Settings and validation”Kafka settings are literal values: strings are quoted and escaped, numbers remain
numbers, and booleans render as 1 or 0. Pass kafka_format: 'JSONEachRow', not
a string containing SQL quotes. Existing MergeTree setting strings keep their
previous raw SQL behavior; this addition does not reinterpret older snapshots.
With engine: 'Kafka' or 'Kafka()', supply nonempty kafka_broker_list,
kafka_topic_list, kafka_group_name, and kafka_format settings. Positional
arguments such as Kafka('broker:9092', 'topic', 'group', 'JSONEachRow') and
Kafka(named_collection) are also retained on pull. The server validates those
arguments and named collections.
Kafka columns cannot have DEFAULT values. Compute defaults in the consuming
materialized view instead. Format-specific settings are passed through; check
their availability on your ClickHouse version.
kafka_auto_offset_reset is not a standard ClickHouse Kafka table setting. Set
auto_offset_reset in the server’s extended Kafka configuration. See the
ClickHouse Kafka reference.
Changing a queue
Section titled “Changing a queue”generate refuses changes to an existing Kafka table’s columns, settings, engine,
or comment with kafka_change_requires_replacement. It leaves migrations and the
snapshot untouched. Kafka does not support the generic column/settings ALTER
operations used for MergeTree tables, and automatically replacing an active queue
could disrupt ingestion.
For an intentional replacement:
- Remove the queue and all consuming materialized views from your schema.
Keep the destination storage table. Generate a migration with
chkit generate --name stop-events-queue. - Re-add the updated queue and views. Generate a second migration with
chkit generate --name restart-events-queue. - Review both migrations, the interruption window, and the consumer group and
offset behavior. Apply with
chkit migrate --apply --allow-destructive.
The generated drop migration removes views before synchronously dropping the queue; the second migration creates the queue before its views. The storage table is retained, and both generated snapshots remain consistent with the migration history. This is an explicit replacement, not a zero-downtime operation. Existing consumer-group offsets, broker retention, and in-flight batches determine whether messages resume, replay, or are unavailable. chkit does not promise exactly-once delivery or reset offsets. Coordinate other consumers and unmanaged views yourself.
generate --empty still provides manual SQL, but intentionally does not update
schema snapshots. Use the two generated migrations above when replacing a managed
queue so the snapshot follows the change.
Pull, drift, and scope
Section titled “Pull, drift, and scope”chkit pull schema emits a Kafka definition without sorting keys and decodes its
settings into literal values. drift and check compare the engine, columns, and
settings declared in the snapshot. Numeric/boolean metadata and SQL string quoting
are normalized without removing meaningful whitespace from strings.
In Python, use chkit pull --database analytics --out-file schema.py,
chkit drift --live, and chkit check --live. Python’s default drift/check also
inspect local schema changes against the snapshot and report Kafka changes that
require replacement, including when scoped with --table.
Consumer offsets, lag, assignments, topic contents, and server configuration or named-collection contents are runtime/external state, outside schema drift. Server settings absent from the desired schema are not treated as drift.
Prefer server-side configuration for credentials. Pull warns when credential
settings are returned in plaintext or redacted; [HIDDEN] credentials must be
resolved before generating SQL. Values evaluated from environment variables in
schema code are still written into snapshots and migration files.
Distributed and other integration engines are outside this release’s scope.