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
17 changes: 16 additions & 1 deletion README.md
Original file line number Diff line number Diff line change
Expand Up @@ -53,9 +53,11 @@ your repository, and cdclint reads it:
| A column is `0`, `''`, `NULL` or `1970-01-01` on every row in ClickHouse, BigQuery, Snowflake or Iceberg, with nothing in the logs | the column is not on the connector's `column.include.list` (or is on its exclude list), so Debezium drops it before Kafka | `sink-column-not-captured` |
| A column added in Postgres never shows up downstream | the migration added it and nobody added it to the include list | `schema-before-connector` (with `--base`), `sink-column-not-captured` |
| Listing the columns of one table made another table's columns disappear | `column.include.list` is one list for every captured table | `sink-column-not-captured` |
| "`table.include.list` not working", or the connector is `RUNNING` and no topic appears | the entry has no schema (`orders` for `public.orders`), or is a glob (`public.bg_*`) where Debezium expects a regex (`public\.bg_.*`); Debezium matches each entry against the whole `schema.table` name, never a substring | `sink-table-not-captured`, when a sink reads the table |
| "`table.include.list` not working", or the connector is `RUNNING` and no topic appears | the entry has no schema (`orders` for `public.orders`), or is a glob (`public.bg_*`) where Debezium expects a regex (`public\.bg_.*`); Debezium matches each entry against the whole `schema.table` name, never a substring | `captured-table-missing` names the entry and the fix; `sink-table-not-captured` when a sink reads the table |
| A Kafka-engine table or sink receives nothing, or the wrong table fills | the topic it reads is not the one the connector produces (prefix, schema, `RegexRouter`) | `topic-table-mapping` |
| A typo in the include list, and a column quietly missing | the pattern matches no column in the source | `captured-column-missing` |
| A typo in `table.include.list`, and a topic that never appears | the entry matches no table in the source | `captured-table-missing` |
| The Kafka Connect sink writes a row of `0` and empty strings for every change, and no error ("sink replicates zero values or nulls") | a `Flatten` transform keeps Debezium's envelope and names every field `after.<column>`, while the sink table names its columns as Postgres does; sinks match fields to columns by name | `sink-column-flattened` |
| A ClickHouse materialized view writes defaults | a refreshable view's `SELECT` order differs from the target's (it matches by position), or a streaming view's names differ (it matches by name) | `mv-column-match` |

It reads files only: Postgres sources today, no type checks, no connection to
Expand All @@ -70,7 +72,9 @@ anything running.
| `sink-column-unknown` | the sink expects a field the source table does not have (warning; renames and computed fields are legitimate) | v0.1 |
| `source-column-not-captured` | a source column nothing captures and nothing reads yet, so the day something asks for it is the day it is found missing (info) | v0.1 |
| `captured-column-missing` | the include list names a column the source does not have (warning) | v0.1 |
| `captured-table-missing` | a `table.include.list` entry matches no table: a typo, a missing schema (`orders` for `public.orders`), or a shell glob (`public.bg_*`) where Debezium reads a regular expression; the last two get the corrected entry as the fix (warning) | next |
| `topic-table-mapping` | a Kafka-engine table reads a topic the connector will not produce | v0.1 |
| `sink-column-flattened` | a sink table names its columns as the source does while a `Flatten` transform on the way delivers the envelope as `after.<column>`, so every column stays at its default | next |
| `mv-column-match` | ClickHouse streaming materialized views match by name, refreshable ones by position. ClickHouse 25.4 and later reject a streaming view that writes a column the target lacks when it is created; a refreshable view's order mismatch was loud on 24.8 and is silent on 26.8 | v0.1 |
| `schema-before-connector` | this change adds a column to a captured table and leaves it out of the stream without deciding to, the trap itself, judged on the diff (warning: leaving PII off is right, so it asks for the decision) | v0.2 |
| `replica-identity` | a captured table's replica identity cannot supply what the sink reads | next |
Expand Down Expand Up @@ -213,6 +217,17 @@ same name, and the finding says so. Debezium's `RegexRouter` transform is
applied when computing topic names. Other sources (MySQL, SQL Server) and
sinks are packages behind the same two interfaces.

Transforms on either connector decide what the sink receives, and cdclint
follows the documented ones in order: Debezium's `ExtractNewRecordState` (or the
Iceberg `DebeziumTransform`) unwraps the envelope into the row, and Kafka
Connect's `Flatten` turns the envelope into `after.<column>`, `before.<column>`,
`source.<field>`, `op` and `ts_ms` (its `delimiter` included). A sink table
named that way, as in ClickHouse's own CDC example, is matched through
`after.` and `before.` to the source columns. Routers, key transforms and
`Filter` change no field name. Any other transform that touches the value, or a
modelled one applied under a predicate, leaves the shape unknown, and cdclint
then reads the sink's column names as the row's, as it always has.

## Design

- **One job.** Lint the contract. Not a CDC platform, not a migration runner,
Expand Down
6 changes: 6 additions & 0 deletions corpus/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -28,6 +28,12 @@ The rest exercise one path each:
| `include-list-typo` | an include-list pattern that matches no column, and the read it silently breaks |
| `golang-migrate-down-files` | Mattermost v10.11.0's `000092_add_createat_to_teammembers.down.sql` drops a column from the wrong table; read in name order it ran just before its up file and deleted `reactions.createat`. Down files are skipped |
| `goose-down-section` | a goose file's down section follows its up section; applied whole, every table was created and dropped again. Down sections are skipped |
| `include-table-no-schema` | Stack Overflow 74103659: `table.include.list` names `ipaddrs` without its schema, the connector runs and no topic appears |
| `include-table-glob` | Stack Overflow 51345636: `public.bg_*` written as a shell glob; as a regular expression it names no table |
| `include-table-typo` | a misspelled `table.include.list` entry, which Debezium only logs as a warning (debezium/dbz#872) |
| `flatten-bare-columns` | ClickHouse/clickhouse-kafka-connect discussion 182: a `Flatten` transform on the Debezium connector sends `after.<column>`, the ClickHouse table names its columns as Postgres does, and every change lands as a row of zeros |
| `flatten-envelope-columns` | ClickHouse's own Postgres CDC example (ClickHouse/examples `cdc/postgresql` at ae417ae): `before.` and `after.` columns resolve to the source columns, so the capture rules see `postcode2`, which the include list added for this entry leaves out. The sink connector config there has no `connector.class`; it is added so the file names its sink |
| `flatten-after-unwrap` | a guard: unwrapping first leaves a flat row, so a `Flatten` after it renames nothing |
| `diff-adds-column-connector-untouched` | the diff rule: `base/` holds the migrations and connector before the change; the change adds a column and leaves the connector alone |
| `diff-connector-captures-one-of-two` | #963 with a hurried fix: two of three new columns go on the include list in the same change; the third is raised, since adding two says nothing about it |
| `diff-connector-touched-other-table` | RefuseRadar #962 and #963 in one range: the connector gains report_validations columns, the migration adds reports columns; the reports ones are raised |
Expand Down
13 changes: 13 additions & 0 deletions corpus/flatten-after-unwrap/connector.json
Original file line number Diff line number Diff line change
@@ -0,0 +1,13 @@
{
"name": "shop-source",
"config": {
"connector.class": "io.debezium.connector.postgresql.PostgresConnector",
"plugin.name": "pgoutput",
"topic.prefix": "shop",
"table.include.list": "public.orders",
"transforms": "unwrap,flatten",
"transforms.unwrap.type": "io.debezium.transforms.ExtractNewRecordState",
"transforms.flatten.type": "org.apache.kafka.connect.transforms.Flatten$Value",
"transforms.flatten.delimiter": "_"
}
}
1 change: 1 addition & 0 deletions corpus/flatten-after-unwrap/expected.txt
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
ok: source, connector and sink agree
6 changes: 6 additions & 0 deletions corpus/flatten-after-unwrap/migrations/0001_orders.sql
Original file line number Diff line number Diff line change
@@ -0,0 +1,6 @@
-- A guard: unwrapping first leaves a flat row, so a flatten after it changes
-- no field name and the plain column names are right.
CREATE TABLE orders (
id BIGSERIAL PRIMARY KEY,
total NUMERIC(12, 2) NOT NULL
);
6 changes: 6 additions & 0 deletions corpus/flatten-after-unwrap/sink-connector.json
Original file line number Diff line number Diff line change
@@ -0,0 +1,6 @@
{
"connector.class": "com.clickhouse.kafka.connect.ClickHouseSinkConnector",
"database": "shop",
"topics": "shop.public.orders",
"topic2TableMap": "shop.public.orders=orders"
}
7 changes: 7 additions & 0 deletions corpus/flatten-after-unwrap/sink/0001_orders.sql
Original file line number Diff line number Diff line change
@@ -0,0 +1,7 @@
CREATE TABLE shop.orders
(
`id` Int64,
`total` String
)
ENGINE = MergeTree
ORDER BY id;
21 changes: 21 additions & 0 deletions corpus/flatten-bare-columns/connector.json
Original file line number Diff line number Diff line change
@@ -0,0 +1,21 @@
{
"config": {
"connector.class": "io.debezium.connector.postgresql.PostgresConnector",
"plugin.name": "pgoutput",
"database.server.name": "postgres",
"table.include.list": "public.testradiomeas",
"topic.prefix": "testprefix",
"publication.name": "myint_logical_replication",
"value.converter": "org.apache.kafka.connect.json.JsonConverter",
"value.converter.schemas.enable": "false",
"snapshot.mode": "never",
"tombstones.on.delete": "false",
"decimal.handling.mode": "double",
"transforms": "set_topic,flatten",
"transforms.flatten.type": "org.apache.kafka.connect.transforms.Flatten$Value",
"transforms.flatten.delimiter": ".",
"transforms.set_topic.type": "org.apache.kafka.connect.transforms.RegexRouter",
"transforms.set_topic.regex": ".*",
"transforms.set_topic.replacement": "testradiomeas"
}
}
5 changes: 5 additions & 0 deletions corpus/flatten-bare-columns/expected.txt
Original file line number Diff line number Diff line change
@@ -0,0 +1,5 @@
error sink-column-flattened flatten-bare-columns/sink/0001_testradiomeas.sql:4
default.testradiomeas (is fed topic testradiomeas by flatten-bare-columns/sink-connector.json) reads measDate, measLongitude, measLatitude, measNetNoCoverage, measNetType by their names in public.testradiomeas, but transforms.flatten in flatten-bare-columns/connector.json flattens each change event, so every one arrives as after.<column>
every row will carry the columns' defaults, with no error anywhere
fix: rename them after.measDate and so on, or unwrap the event with io.debezium.transforms.ExtractNewRecordState in place of transforms.flatten
1 error(s), 0 warning(s), 0 info
11 changes: 11 additions & 0 deletions corpus/flatten-bare-columns/migrations/0001_testradiomeas.sql
Original file line number Diff line number Diff line change
@@ -0,0 +1,11 @@
-- Reduced from ClickHouse/clickhouse-kafka-connect discussion 182 and issue
-- 183 ("ClickHouseSinkConnector & Debezium PostgresConnector creates rows
-- zero values or nulls", 2023-09): every insert, update and delete produced
-- a ClickHouse row of zeros and empty strings, with no error anywhere.
CREATE TABLE testradiomeas (
"measDate" INTEGER NOT NULL,
"measLongitude" DOUBLE PRECISION,
"measLatitude" DOUBLE PRECISION,
"measNetNoCoverage" INTEGER,
"measNetType" TEXT
);
9 changes: 9 additions & 0 deletions corpus/flatten-bare-columns/sink-connector.json
Original file line number Diff line number Diff line change
@@ -0,0 +1,9 @@
{
"connector.class": "com.clickhouse.kafka.connect.ClickHouseSinkConnector",
"tasks.max": "1",
"database": "default",
"exactlyOnce": "false",
"topics": "testradiomeas",
"value.converter": "org.apache.kafka.connect.json.JsonConverter",
"value.converter.schemas.enable": "false"
}
13 changes: 13 additions & 0 deletions corpus/flatten-bare-columns/sink/0001_testradiomeas.sql
Original file line number Diff line number Diff line change
@@ -0,0 +1,13 @@
-- The table as first written: the columns named as they are in Postgres.
-- The flatten transform names every field after.<column>, and the sink
-- matches fields to columns by name, so nothing matched.
CREATE TABLE default.testradiomeas
(
`measDate` Int32,
`measLongitude` Float64,
`measLatitude` Float64,
`measNetNoCoverage` Int32,
`measNetType` String
)
ENGINE = MergeTree
ORDER BY `measDate`;
23 changes: 23 additions & 0 deletions corpus/flatten-envelope-columns/connector.json
Original file line number Diff line number Diff line change
@@ -0,0 +1,23 @@
{
"connector.class": "PostgresCdcSource",
"name": "uk_price_paid_changes",
"database.server.name": "postgres",
"database.dbname": "postgres",
"publication.name": "dbz_publication",
"publication.autocreate.mode": "filtered",
"table.include.list": "public.uk_price_paid",
"column.include.list": "public\\.uk_price_paid\\.(id|price|postcode1|town)",
"snapshot.mode": "never",
"tombstones.on.delete": "false",
"plugin.name": "pgoutput",
"output.data.format": "JSON",
"after.state.only": "false",
"output.key.format": "STRING",
"tasks.max": "1",
"transforms": "flatten,set_topic",
"transforms.flatten.type": "org.apache.kafka.connect.transforms.Flatten$Value",
"transforms.flatten.delimiter": ".",
"transforms.set_topic.type": "io.confluent.connect.cloud.transforms.TopicRegexRouter",
"transforms.set_topic.regex": ".*",
"transforms.set_topic.replacement": "uk_price_paid_changes"
}
9 changes: 9 additions & 0 deletions corpus/flatten-envelope-columns/expected.txt
Original file line number Diff line number Diff line change
@@ -0,0 +1,9 @@
error sink-column-not-captured flatten-envelope-columns/sink/0001_uk_price_paid_changes.sql:6
public.uk_price_paid.postcode2 is read by default.uk_price_paid_changes (is fed topic uk_price_paid_changes by flatten-envelope-columns/sink-connector.json) but is not matched by column.include.list in flatten-envelope-columns/connector.json
every row will carry the column's default, with no error anywhere
fix: add public.uk_price_paid.postcode2 to column.include.list, deploy the connector, then apply the sink schema; rows already written need a snapshot
error sink-column-not-captured flatten-envelope-columns/sink/0001_uk_price_paid_changes.sql:11
public.uk_price_paid.postcode2 is read by default.uk_price_paid_changes (is fed topic uk_price_paid_changes by flatten-envelope-columns/sink-connector.json) but is not matched by column.include.list in flatten-envelope-columns/connector.json
every row will carry the column's default, with no error anywhere
fix: add public.uk_price_paid.postcode2 to column.include.list, deploy the connector, then apply the sink schema; rows already written need a snapshot
2 error(s), 0 warning(s), 0 info
12 changes: 12 additions & 0 deletions corpus/flatten-envelope-columns/migrations/0001_uk_price_paid.sql
Original file line number Diff line number Diff line change
@@ -0,0 +1,12 @@
-- Reduced from ClickHouse's own Postgres CDC example (ClickHouse/examples,
-- cdc/postgresql at ae417ae, and the blog post "ClickHouse PostgreSQL change
-- data capture, part 2"): the change events keep Debezium's envelope, a
-- flatten transform names the fields before.<column> and after.<column>, and
-- the ClickHouse table names its columns the same way.
CREATE TABLE uk_price_paid (
id BIGSERIAL PRIMARY KEY,
price INTEGER NOT NULL,
postcode1 TEXT,
postcode2 TEXT,
town TEXT
);
10 changes: 10 additions & 0 deletions corpus/flatten-envelope-columns/sink-connector.json
Original file line number Diff line number Diff line change
@@ -0,0 +1,10 @@
{
"connector.class": "com.clickhouse.kafka.connect.ClickHouseSinkConnector",
"database": "default",
"exactlyOnce": "false",
"schemas.enable": "false",
"topics": "uk_price_paid_changes",
"value.converter": "org.apache.kafka.connect.json.JsonConverter",
"value.converter.schemas.enable": "false",
"key.converter": "org.apache.kafka.connect.storage.StringConverter"
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,17 @@
CREATE TABLE default.uk_price_paid_changes
(
`before.id` Nullable(UInt64),
`before.price` Nullable(UInt32),
`before.postcode1` Nullable(String),
`before.postcode2` Nullable(String),
`before.town` Nullable(String),
`after.id` Nullable(UInt64),
`after.price` Nullable(UInt32),
`after.postcode1` Nullable(String),
`after.postcode2` Nullable(String),
`after.town` Nullable(String),
`op` LowCardinality(String),
`ts_ms` UInt64
)
ENGINE = MergeTree
ORDER BY tuple();
9 changes: 9 additions & 0 deletions corpus/include-table-glob/connector.json
Original file line number Diff line number Diff line change
@@ -0,0 +1,9 @@
{
"name": "shop-source",
"config": {
"connector.class": "io.debezium.connector.postgresql.PostgresConnector",
"plugin.name": "pgoutput",
"topic.prefix": "shop",
"table.include.list": "public.bg_*"
}
}
9 changes: 9 additions & 0 deletions corpus/include-table-glob/expected.txt
Original file line number Diff line number Diff line change
@@ -0,0 +1,9 @@
warning captured-table-missing include-table-glob/connector.json
public.bg_* in table.include.list matches no table in the source schema
Debezium reads each entry as a regular expression, where * repeats the character before it
fix: write public\.bg_.* to capture public.bg_items and public.bg_orders, or remove it
error sink-table-not-captured include-table-glob/sink/0001_bg_orders.sql:1
kafka_bg_orders reads topic shop.public.bg_orders, but public.bg_orders is not matched by table.include.list in include-table-glob/connector.json
nothing will ever arrive on that topic
fix: add public.bg_orders to table.include.list and deploy the connector before the sink schema
1 error(s), 1 warning(s), 0 info
21 changes: 21 additions & 0 deletions corpus/include-table-glob/migrations/0001_tables.sql
Original file line number Diff line number Diff line change
@@ -0,0 +1,21 @@
-- Reduced from Stack Overflow question 51345636 (5.9k views): the include
-- list was written as a shell glob, public.bg_*, meaning every table whose
-- name starts with bg_. Debezium reads it as a regular expression, where *
-- repeats the character before it, so it names public.bg, public.bg_,
-- public.bg__ and no real table.
CREATE TABLE bg_orders (
id BIGSERIAL PRIMARY KEY,
amount NUMERIC(12, 2) NOT NULL,
created_at TIMESTAMPTZ NOT NULL DEFAULT now()
);

CREATE TABLE bg_items (
id BIGSERIAL PRIMARY KEY,
order_id BIGINT NOT NULL REFERENCES bg_orders (id),
sku TEXT NOT NULL
);

CREATE TABLE cp_users (
id BIGSERIAL PRIMARY KEY,
email TEXT NOT NULL
);
12 changes: 12 additions & 0 deletions corpus/include-table-glob/sink/0001_bg_orders.sql
Original file line number Diff line number Diff line change
@@ -0,0 +1,12 @@
CREATE TABLE kafka_bg_orders
(
`id` Int64,
`amount` String,
`created_at` String
)
ENGINE = Kafka
SETTINGS
kafka_broker_list = '${KAFKA_BROKERS}',
kafka_topic_list = 'shop.public.bg_orders',
kafka_group_name = 'clickhouse-bg-orders',
kafka_format = 'JSONEachRow';
9 changes: 9 additions & 0 deletions corpus/include-table-no-schema/connector.json
Original file line number Diff line number Diff line change
@@ -0,0 +1,9 @@
{
"name": "testdb",
"config": {
"connector.class": "io.debezium.connector.postgresql.PostgresConnector",
"plugin.name": "pgoutput",
"topic.prefix": "test",
"table.include.list": "ipaddrs"
}
}
9 changes: 9 additions & 0 deletions corpus/include-table-no-schema/expected.txt
Original file line number Diff line number Diff line change
@@ -0,0 +1,9 @@
warning captured-table-missing include-table-no-schema/connector.json
ipaddrs in table.include.list matches no table in the source schema
Debezium matches each entry against the whole schema.table name
fix: write myschema\.ipaddrs, or remove it
error sink-table-not-captured include-table-no-schema/sink/0001_ipaddrs.sql:3
kafka_ipaddrs reads topic test.myschema.ipaddrs, but myschema.ipaddrs is not matched by table.include.list in include-table-no-schema/connector.json
nothing will ever arrive on that topic
fix: add myschema.ipaddrs to table.include.list and deploy the connector before the sink schema
1 error(s), 1 warning(s), 0 info
11 changes: 11 additions & 0 deletions corpus/include-table-no-schema/migrations/0001_ipaddrs.sql
Original file line number Diff line number Diff line change
@@ -0,0 +1,11 @@
-- Reduced from Stack Overflow question 74103659 ("Debezium PostgreSQL
-- connector not creating topic", 6.9k views) and the Debezium Google Group
-- thread Kp_9jbTTEkY: the table lives in a schema, the include list names it
-- without one, the connector reports RUNNING and nothing is ever captured.
CREATE SCHEMA myschema;

CREATE TABLE myschema.ipaddrs (
id BIGSERIAL PRIMARY KEY,
address INET NOT NULL,
seen_at TIMESTAMPTZ NOT NULL DEFAULT now()
);
14 changes: 14 additions & 0 deletions corpus/include-table-no-schema/sink/0001_ipaddrs.sql
Original file line number Diff line number Diff line change
@@ -0,0 +1,14 @@
-- The warehouse side, waiting on the topic Debezium would produce for
-- myschema.ipaddrs once the table is captured.
CREATE TABLE kafka_ipaddrs
(
`id` Int64,
`address` String,
`seen_at` String
)
ENGINE = Kafka
SETTINGS
kafka_broker_list = '${KAFKA_BROKERS}',
kafka_topic_list = 'test.myschema.ipaddrs',
kafka_group_name = 'clickhouse-ipaddrs',
kafka_format = 'JSONEachRow';
9 changes: 9 additions & 0 deletions corpus/include-table-typo/connector.json
Original file line number Diff line number Diff line change
@@ -0,0 +1,9 @@
{
"name": "shop-source",
"config": {
"connector.class": "io.debezium.connector.postgresql.PostgresConnector",
"plugin.name": "pgoutput",
"topic.prefix": "shop",
"table.include.list": "public.customers,public.oders"
}
}
8 changes: 8 additions & 0 deletions corpus/include-table-typo/expected.txt
Original file line number Diff line number Diff line change
@@ -0,0 +1,8 @@
warning captured-table-missing include-table-typo/connector.json
public.oders in table.include.list matches no table in the source schema
fix: remove it, or check the spelling against the migrations
error sink-table-not-captured include-table-typo/sink/0001_orders.sql:1
kafka_orders reads topic shop.public.orders, but public.orders is not matched by table.include.list in include-table-typo/connector.json
nothing will ever arrive on that topic
fix: add public.orders to table.include.list and deploy the connector before the sink schema
1 error(s), 1 warning(s), 0 info
Loading