Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
13 changes: 13 additions & 0 deletions .changeset/kafka-table-engine.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,13 @@
---
"@chkit/core": minor
"@chkit/clickhouse": patch
"@chkit/plugin-pull": minor
"chkit": minor
---

Support Kafka engine tables without MergeTree key clauses, with escaped literal
settings, pull round trips, and normalized drift/check comparisons. Preserve the
existing setting-string contract for other engines. Reject unsupported Kafka
changes before writing migration artifacts and document an explicit, destructive-
gated replacement workflow. Handle quoted delimiters and escaped trailing
backslashes in introspection and migration statement splitting.
48 changes: 48 additions & 0 deletions .github/workflows/kafka.yml
Original file line number Diff line number Diff line change
@@ -0,0 +1,48 @@
name: Kafka integration

on:
pull_request:
paths:
- 'packages/**'
- 'chkit_python/**'
- 'test/kafka/**'
- '.github/workflows/kafka.yml'
- 'bun.lock'
push:
branches: [main]

permissions:
contents: read

jobs:
kafka:
runs-on: blacksmith-8vcpu-ubuntu-2404
strategy:
fail-fast: false
matrix:
clickhouse: ['25.3', '26.3']
env:
CLICKHOUSE_VERSION: ${{ matrix.clickhouse }}
steps:
- uses: actions/checkout@v4
- uses: ./.github/actions/setup
- uses: actions/setup-python@v5
with:
python-version: '3.12'
- name: Build TypeScript and install Python
run: |
bunx turbo run build --filter=chkit... --filter=@chkit/plugin-pull
python -m venv .venv
.venv/bin/pip install -e './chkit_python[dev]'
- name: Start Kafka and ClickHouse
run: docker compose -p chkit-issue203 -f test/kafka/docker-compose.yml up -d --wait
- name: TypeScript Kafka workflow
run: bun test test/kafka/kafka.e2e.test.ts
- name: Python Kafka workflow
run: .venv/bin/python -m pytest test/kafka/test_python_e2e.py -q
- name: Service logs
if: failure()
run: docker compose -p chkit-issue203 -f test/kafka/docker-compose.yml logs --tail=150
- name: Clean up
if: always()
run: docker compose -p chkit-issue203 -f test/kafka/docker-compose.yml down -v
6 changes: 6 additions & 0 deletions apps/docs/src/content/docs/guides/clickhouse-compatibility.md
Original file line number Diff line number Diff line change
Expand Up @@ -13,6 +13,12 @@ The continuous test suite runs against ObsessionDB, so the `SharedMergeTree`/`Sh

## Version-gated features

[Kafka engine tables](/schema/kafka/) support creation, pull, and schema drift,
validated on self-hosted ClickHouse 25.3 and 26.3. Existing queue changes require
explicit replacement; generic Kafka column/settings ALTERs are not generated.
The server must provide the Kafka engine and broker connectivity. Distributed and
other integration engines remain outside this support scope.

A few schema features depend on the ClickHouse version of your target:

| Feature | Requirement |
Expand Down
7 changes: 7 additions & 0 deletions apps/docs/src/content/docs/schema/dsl-reference.mdx
Original file line number Diff line number Diff line change
Expand Up @@ -168,6 +168,10 @@ Creates a table definition.

### Optional fields

For [Kafka tables](/schema/kafka/), omit `primaryKey` and `orderBy`. Kafka also
rejects storage clauses (`partitionBy`, `uniqueKey`, `ttl`, indexes, projections)
and column defaults. Its setting strings are escaped SQL literals.

| Field | Type | Description |
|-------|------|-------------|
| `partitionBy` | `string` | Partition expression, e.g. `'toYYYYMM(created_at)'` |
Expand Down Expand Up @@ -907,6 +911,9 @@ chkit validates schema definitions and throws a `ChxValidationError` if any issu

## Structural vs. alterable properties

The table rules below apply to MergeTree-family tables. Kafka changes require an
[explicit replacement](/schema/kafka/#changing-a-queue); chkit refuses generic ALTERs.

When a property changes, chkit determines whether the table can be altered in place or must be dropped and recreated.

**Structural** (drop + recreate): `engine`, `primaryKey`, `orderBy`, `partitionBy`, `uniqueKey`
Expand Down
155 changes: 155 additions & 0 deletions apps/docs/src/content/docs/schema/kafka.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,155 @@
---
title: Kafka tables
description: Define Kafka queues and materialized-view ingestion pipelines in schema code.
sidebar:
order: 4
---

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.

```ts
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:

```python
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

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](https://clickhouse.com/docs/engines/table-engines/integrations/kafka).

## 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:

1. 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`.
2. Re-add the updated queue and views. Generate a second migration with
`chkit generate --name restart-events-queue`.
3. 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

`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.
1 change: 1 addition & 0 deletions apps/docs/src/content/docs/schema/overview.mdx
Original file line number Diff line number Diff line change
Expand Up @@ -69,6 +69,7 @@ chkit discovers schema files using the `schema` glob in your [configuration](/co

- [DSL Reference](/schema/dsl-reference/) — every function, option, and column type, in both languages.
- [Refreshable Views](/schema/refreshable-views/) — using ClickHouse refreshable materialized views from chkit.
- [Kafka tables](/schema/kafka/) — queue → materialized view → storage pipelines, pull, drift, and explicit replacements.

## Related

Expand Down
11 changes: 11 additions & 0 deletions chkit_python/CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -2,10 +2,21 @@

## Unreleased

### Added
- Add `SkipIndexText` for full-text index generation, introspection, pull, and drift.
Preserve quoted SQL literals, normalize ClickHouse’s fixed granularity, and reject
malformed or unsupported metadata. Exercise adversarial round trips and actual
indexed search results on ClickHouse 26.3 and 26.8.
- Native Kafka table definitions without sorting keys, escaped literal settings,
pull round trips, and normalized live drift/check comparisons.
- Validation and migration guards reject unsupported Kafka changes before
writing artifacts. Offline drift/check report changes requiring replacement;
explicit drop/create migrations preserve the existing destructive-operation gate.
- Kafka ingestion and replacement integration tests on ClickHouse 25.3 and 26.3.

### Fixed
- Preserve quoted clause names, delimiters, whitespace, and escaped trailing
backslashes in table introspection and migration statement splitting.

## 0.2.0 — 2026-08-10

Expand Down
17 changes: 6 additions & 11 deletions chkit_python/src/chkit/cli/commands/check.py
Original file line number Diff line number Diff line change
Expand Up @@ -22,6 +22,7 @@
from chkit.cli.commands.drift_compare import summarize_drift_reasons
from chkit.cli.commands.drift_payload import build_drift_payload
from chkit.cli.commands.migrate_scope import filter_pending_by_scope
from chkit.cli.commands.snapshot_drift import plan_snapshot_drift
from chkit.cli.config_loader import load_config
from chkit.cli.journal_store import JournalStore
from chkit.cli.migration_store import (
Expand All @@ -32,14 +33,12 @@
from chkit.cli.plugin_runtime import load_plugin_runtime
from chkit.cli.schema_loader import load_schema
from chkit.cli.table_scope import (
filter_plan_by_table_scope,
resolve_table_scope,
table_keys_from_definitions,
)
from chkit.clickhouse.client import ClickHouseClient
from chkit.core.canonical import canonicalize_definitions
from chkit.core.model import ChxConfigEnv
from chkit.core.planner import plan_diff
from chkit.core.validate import validate_definitions
from chkit.plugins import ChxOnCheckContext, ChxPlugin

Expand Down Expand Up @@ -106,13 +105,11 @@ def run( # noqa: PLR0912, PLR0915

drift_ops: list[str] = []
drift_reason_counts: dict[str, int] = {}
replacement_drift = False
if snapshot is not None:
plan = plan_diff(snapshot_defs, schema_defs)
if table_scope.enabled:
filtered = filter_plan_by_table_scope(
plan, set(table_scope.matched_tables)
)
plan = filtered.plan
plan, replacement_issues = plan_snapshot_drift(snapshot_defs, schema_defs, table_scope)
replacement_drift = bool(replacement_issues)
issues.extend(issue.model_dump(mode="json") for issue in replacement_issues)
drift_ops = [op.key for op in plan.operations]

live_drifted = False
Expand Down Expand Up @@ -157,9 +154,7 @@ def run( # noqa: PLR0912, PLR0915
failed_checks.append("pending_migrations")
if fail_on_mismatch and mismatches:
failed_checks.append("checksum_mismatch")
drift_fired = (
bool(drift_ops) if not live else live_drifted or bool(drift_ops)
)
drift_fired = bool(drift_ops) or replacement_drift or (live and live_drifted)
if fail_on_drift and drift_fired:
# Match the TS finding code: ``schema_drift`` (was ``drift``).
failed_checks.append("schema_drift")
Expand Down
Loading
Loading