From cf0363a4e13e090a53ad0e4d964802117acfcc28 Mon Sep 17 00:00:00 2001 From: avison9 Date: Fri, 25 Sep 2026 22:25:43 +0100 Subject: [PATCH 1/2] A table.include.list entry that matches no table is raised, with the corrected entry when the mistake is a missing schema or a shell glob Debezium logs a warning when a table.include.list entry matches nothing and keeps running, so the topic is simply never produced. People ask it to fail instead (debezium/dbz#872); a maintainer answered that a table may be created later. cdclint raises it as a warning for the same reason: when a sink actually reads the table, sink-table-not-captured already fails the pull request. The two shapes behind the most-viewed questions get a pointed fix, because both follow from Debezium's documented matching (each entry is a regular expression matched against the whole schema.table name): an entry without its schema (Stack Overflow 74103659: ipaddrs for myschema.ipaddrs) gets "write myschema\.ipaddrs", and a shell glob (Stack Overflow 51345636: public.bg_* names no table as a regex) gets "write public\.bg_.* to capture ..." with the tables it would capture. An entry that already contains .* is read as the regex it is. Exclude-list entries are not checked: one that matches nothing excludes nothing, the same reasoning as captured-column-missing. Three corpus entries, written before the code: include-table-no-schema, include-table-glob, include-table-typo. --- README.md | 4 +- corpus/README.md | 3 + corpus/include-table-glob/connector.json | 9 ++ corpus/include-table-glob/expected.txt | 9 ++ .../migrations/0001_tables.sql | 21 ++++ .../sink/0001_bg_orders.sql | 12 +++ corpus/include-table-no-schema/connector.json | 9 ++ corpus/include-table-no-schema/expected.txt | 9 ++ .../migrations/0001_ipaddrs.sql | 11 +++ .../sink/0001_ipaddrs.sql | 14 +++ corpus/include-table-typo/connector.json | 9 ++ corpus/include-table-typo/expected.txt | 8 ++ .../migrations/0001_tables.sql | 14 +++ .../include-table-typo/sink/0001_orders.sql | 12 +++ internal/capture/debezium/debezium.go | 11 +++ internal/engine/engine.go | 96 ++++++++++++++++++- 16 files changed, 247 insertions(+), 4 deletions(-) create mode 100644 corpus/include-table-glob/connector.json create mode 100644 corpus/include-table-glob/expected.txt create mode 100644 corpus/include-table-glob/migrations/0001_tables.sql create mode 100644 corpus/include-table-glob/sink/0001_bg_orders.sql create mode 100644 corpus/include-table-no-schema/connector.json create mode 100644 corpus/include-table-no-schema/expected.txt create mode 100644 corpus/include-table-no-schema/migrations/0001_ipaddrs.sql create mode 100644 corpus/include-table-no-schema/sink/0001_ipaddrs.sql create mode 100644 corpus/include-table-typo/connector.json create mode 100644 corpus/include-table-typo/expected.txt create mode 100644 corpus/include-table-typo/migrations/0001_tables.sql create mode 100644 corpus/include-table-typo/sink/0001_orders.sql diff --git a/README.md b/README.md index 6ed9833..543e302 100644 --- a/README.md +++ b/README.md @@ -53,9 +53,10 @@ 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` | | 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 @@ -70,6 +71,7 @@ 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 | | `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 | diff --git a/corpus/README.md b/corpus/README.md index c2dddf5..41fdbbe 100644 --- a/corpus/README.md +++ b/corpus/README.md @@ -28,6 +28,9 @@ 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) | | `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 | diff --git a/corpus/include-table-glob/connector.json b/corpus/include-table-glob/connector.json new file mode 100644 index 0000000..96dc386 --- /dev/null +++ b/corpus/include-table-glob/connector.json @@ -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_*" + } +} diff --git a/corpus/include-table-glob/expected.txt b/corpus/include-table-glob/expected.txt new file mode 100644 index 0000000..121f076 --- /dev/null +++ b/corpus/include-table-glob/expected.txt @@ -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 diff --git a/corpus/include-table-glob/migrations/0001_tables.sql b/corpus/include-table-glob/migrations/0001_tables.sql new file mode 100644 index 0000000..67aa027 --- /dev/null +++ b/corpus/include-table-glob/migrations/0001_tables.sql @@ -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 +); diff --git a/corpus/include-table-glob/sink/0001_bg_orders.sql b/corpus/include-table-glob/sink/0001_bg_orders.sql new file mode 100644 index 0000000..3343339 --- /dev/null +++ b/corpus/include-table-glob/sink/0001_bg_orders.sql @@ -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'; diff --git a/corpus/include-table-no-schema/connector.json b/corpus/include-table-no-schema/connector.json new file mode 100644 index 0000000..2714711 --- /dev/null +++ b/corpus/include-table-no-schema/connector.json @@ -0,0 +1,9 @@ +{ + "name": "testdb", + "config": { + "connector.class": "io.debezium.connector.postgresql.PostgresConnector", + "plugin.name": "pgoutput", + "topic.prefix": "test", + "table.include.list": "ipaddrs" + } +} diff --git a/corpus/include-table-no-schema/expected.txt b/corpus/include-table-no-schema/expected.txt new file mode 100644 index 0000000..79c00e0 --- /dev/null +++ b/corpus/include-table-no-schema/expected.txt @@ -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 diff --git a/corpus/include-table-no-schema/migrations/0001_ipaddrs.sql b/corpus/include-table-no-schema/migrations/0001_ipaddrs.sql new file mode 100644 index 0000000..de36077 --- /dev/null +++ b/corpus/include-table-no-schema/migrations/0001_ipaddrs.sql @@ -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() +); diff --git a/corpus/include-table-no-schema/sink/0001_ipaddrs.sql b/corpus/include-table-no-schema/sink/0001_ipaddrs.sql new file mode 100644 index 0000000..1c86740 --- /dev/null +++ b/corpus/include-table-no-schema/sink/0001_ipaddrs.sql @@ -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'; diff --git a/corpus/include-table-typo/connector.json b/corpus/include-table-typo/connector.json new file mode 100644 index 0000000..ef131ba --- /dev/null +++ b/corpus/include-table-typo/connector.json @@ -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" + } +} diff --git a/corpus/include-table-typo/expected.txt b/corpus/include-table-typo/expected.txt new file mode 100644 index 0000000..f19d386 --- /dev/null +++ b/corpus/include-table-typo/expected.txt @@ -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 diff --git a/corpus/include-table-typo/migrations/0001_tables.sql b/corpus/include-table-typo/migrations/0001_tables.sql new file mode 100644 index 0000000..e7c5914 --- /dev/null +++ b/corpus/include-table-typo/migrations/0001_tables.sql @@ -0,0 +1,14 @@ +-- The plain shape: one include-list entry misspelled. Debezium logs a +-- warning when an entry matches nothing and keeps running (Debezium issue +-- dbz#872 asks it to fail instead), so the orders topic is simply never +-- produced. +CREATE TABLE customers ( + id BIGSERIAL PRIMARY KEY, + email TEXT NOT NULL +); + +CREATE TABLE orders ( + id BIGSERIAL PRIMARY KEY, + customer_id BIGINT NOT NULL REFERENCES customers (id), + total NUMERIC(12, 2) NOT NULL +); diff --git a/corpus/include-table-typo/sink/0001_orders.sql b/corpus/include-table-typo/sink/0001_orders.sql new file mode 100644 index 0000000..e755bb2 --- /dev/null +++ b/corpus/include-table-typo/sink/0001_orders.sql @@ -0,0 +1,12 @@ +CREATE TABLE kafka_orders +( + `id` Int64, + `customer_id` Int64, + `total` String +) +ENGINE = Kafka +SETTINGS + kafka_broker_list = '${KAFKA_BROKERS}', + kafka_topic_list = 'shop.public.orders', + kafka_group_name = 'clickhouse-orders', + kafka_format = 'JSONEachRow'; diff --git a/internal/capture/debezium/debezium.go b/internal/capture/debezium/debezium.go index 2ef7f03..5b6dc31 100644 --- a/internal/capture/debezium/debezium.go +++ b/internal/capture/debezium/debezium.go @@ -33,6 +33,7 @@ type Contract struct { columnListKey string tableListKey string columnPatterns []string + tablePatterns []string routers []router Config map[string]string Class string @@ -95,6 +96,11 @@ func Parse(b []byte) (*Contract, error) { if c.tableInclude, c.tableListKey, err = list(cfg, "table.include.list", "table.whitelist"); err != nil { return nil, err } + for _, p := range strings.Split(first(cfg, "table.include.list", "table.whitelist"), ",") { + if p = strings.TrimSpace(p); p != "" { + c.tablePatterns = append(c.tablePatterns, p) + } + } if c.tableExclude, _, err = list(cfg, "table.exclude.list", "table.blacklist"); err != nil { return nil, err } @@ -232,4 +238,9 @@ func (c *Contract) ColumnListSetting() string { return c.columnListKey } // captured-column-missing rule. Exclude lists are not returned: a pattern // there that matches nothing excludes nothing, which is harmless. func (c *Contract) ColumnPatterns() []string { return c.columnPatterns } + +// TablePatterns returns the table include-list patterns as written, for the +// captured-table-missing rule. Exclude lists are left out for the same +// reason as columns. +func (c *Contract) TablePatterns() []string { return c.tablePatterns } func (c *Contract) TableListSetting() string { return c.tableListKey } diff --git a/internal/engine/engine.go b/internal/engine/engine.go index 8370e92..199eb36 100644 --- a/internal/engine/engine.go +++ b/internal/engine/engine.go @@ -6,6 +6,8 @@ package engine import ( "fmt" + "regexp" + "sort" "strings" "github.com/avison9/cdclint/internal/model" @@ -18,7 +20,7 @@ type Input struct { Source *model.Source Contract model.Contract // Patterns, when the contract can list them, lets captured-column-missing - // check each include pattern against the schema. + // and captured-table-missing check each include pattern against the schema. Patterns PatternLister Sinks []*model.Sink Views []clickhouse.View @@ -27,10 +29,11 @@ type Input struct { Base *Base } -// PatternLister is implemented by contracts whose column list is a set of -// patterns that can be checked one by one. +// PatternLister is implemented by contracts whose column and table lists are +// sets of patterns that can be checked one by one. type PatternLister interface { ColumnPatterns() []string + TablePatterns() []string } // Read is one sink column resolved to the source column it expects. @@ -134,6 +137,7 @@ func Run(in *Input) []model.Finding { fs = append(fs, diff...) fs = append(fs, sourceColumnNotCaptured(in, reads, raised)...) fs = append(fs, capturedColumnMissing(in)...) + fs = append(fs, capturedTableMissing(in)...) fs = append(fs, mvColumnMatch(in)...) model.Sort(fs) return fs @@ -291,6 +295,92 @@ func capturedColumnMissing(in *Input) []model.Finding { return fs } +// capturedTableMissing checks each table include-list entry against the +// schema. Debezium logs a warning when an entry matches no table and keeps +// running; debezium/dbz#872 asks it to fail instead, and a maintainer +// answered that a table may be created later. That is why this is a warning: +// when a sink reads the table, sink-table-not-captured raises the error. +// +// Two shapes get a pointed fix because they are how people get it wrong in +// practice (Stack Overflow 74103659 and 51345636): an entry without its schema, +// and a shell glob. Both follow from Debezium's documented matching: each +// entry is a regular expression matched against the whole schema.table name, +// never a substring. +func capturedTableMissing(in *Input) []model.Finding { + if in.Patterns == nil { + return nil + } + var fs []model.Finding + for _, p := range in.Patterns.TablePatterns() { + for _, name := range expand(p) { + m, err := compileAnchored(name) + if err != nil || tablesMatching(in, func(t *model.Table) bool { return m(t.Qualified()) }) != nil { + continue + } + message := fmt.Sprintf("%s in %s matches no table in the source schema", name, in.Contract.TableListSetting()) + fix := "remove it, or check the spelling against the migrations" + if bare := tablesMatching(in, func(t *model.Table) bool { return m(t.Name) }); bare != nil { + message += "\nDebezium matches each entry against the whole schema.table name" + var qualified []string + for _, t := range bare { + qualified = append(qualified, regexp.QuoteMeta(t.Qualified())) + } + fix = "write " + strings.Join(qualified, " or ") + ", or remove it" + } else if re, caught := globReading(in, name); caught != nil { + message += "\nDebezium reads each entry as a regular expression, where * repeats the character before it" + var names []string + for _, t := range caught { + names = append(names, t.Qualified()) + } + fix = "write " + re + " to capture " + joinAnd(names) + ", or remove it" + } + fs = append(fs, model.Finding{ + Rule: "captured-table-missing", Severity: model.Warning, Pos: in.Contract.Pos(), + Message: message, + Fix: fix, + }) + } + } + return fs +} + +// tablesMatching returns the source tables for which match is true, in +// qualified-name order, or nil when there are none. +func tablesMatching(in *Input, match func(*model.Table) bool) []*model.Table { + var out []*model.Table + for _, t := range in.Source.Tables { + if match(t) { + out = append(out, t) + } + } + sort.Slice(out, func(i, j int) bool { return out[i].Qualified() < out[j].Qualified() }) + return out +} + +// globReading reads an entry the way its author most likely meant it, as a +// shell glob where * is any run of characters, and returns the regular +// expression that says so and the tables it would capture. An entry that +// already contains .* was written as a regular expression and is left alone. +func globReading(in *Input, entry string) (string, []*model.Table) { + if !strings.Contains(entry, "*") || strings.Contains(entry, ".*") { + return "", nil + } + re := strings.ReplaceAll(regexp.QuoteMeta(entry), `\*`, ".*") + m, err := compileAnchored(re) + if err != nil { + return "", nil + } + return re, tablesMatching(in, func(t *model.Table) bool { return m(t.Qualified()) }) +} + +// joinAnd lists names as "a", "a and b" or "a, b and c". +func joinAnd(names []string) string { + if len(names) < 2 { + return strings.Join(names, "") + } + return strings.Join(names[:len(names)-1], ", ") + " and " + names[len(names)-1] +} + // expand turns the common `schema.table\.(a|b|c)` shape into one pattern // per alternative so each column is checked on its own. Any other shape is // returned as is. From 6133eb3aadd64ab9d491deb9c96bc173579ffa27 Mon Sep 17 00:00:00 2001 From: avison9 Date: Fri, 25 Sep 2026 22:36:05 +0100 Subject: [PATCH 2/2] Sink columns are matched to the field that actually arrives after the connectors' transforms, and a table that misses a flatten is raised A Flatten transform on the Debezium envelope names every field after., and sinks match fields to columns by name. A ClickHouse table that names its columns as Postgres does therefore receives a row of zeros and empty strings for every change, with no error (ClickHouse/clickhouse-kafka-connect discussion 182). cdclint said "ok" on that exact configuration, because it took every sink column name as a row column name. The same assumption made it raise twelve false sink-column-unknown warnings on ClickHouse's own CDC example, whose table names its columns before. and after., and miss that one of those columns was never captured. A new internal/smt package works out the value's shape from each connector's transforms, in order: the envelope for a Debezium connector (or after.state.only false), the row after ExtractNewRecordState or the Iceberg DebeziumTransform, and the flattened envelope after Flatten, with its delimiter. The sink connector's transforms apply after the source's. When the result is a flattened envelope, after. and before. resolve to source column c, so the capture rules see through them; op, ts_ms and source/transaction fields are envelope metadata; and a column named as the row names it is raised by the new sink-column-flattened rule, once per table. Only documented transforms are modelled. Routers, key transforms and Filter change no field name; any other value transform, or a modelled one under a predicate, makes the shape unknown, and an unknown shape keeps the old behaviour and produces no new finding. Rejected: raising a bare envelope with no transform at all. Whether a sink handles Debezium's envelope itself, or a converter reshapes it, is not in the files, so that would be a guess. Corpus first: flatten-bare-columns (discussion 182's configs), flatten-envelope-columns (ClickHouse/examples cdc/postgresql at ae417ae, with an include list that leaves postcode2 out), flatten-after-unwrap (a guard). RefuseRadar, which unwraps, is unchanged at 0/0/159. --- README.md | 13 +++ corpus/README.md | 3 + corpus/flatten-after-unwrap/connector.json | 13 +++ corpus/flatten-after-unwrap/expected.txt | 1 + .../migrations/0001_orders.sql | 6 ++ .../flatten-after-unwrap/sink-connector.json | 6 ++ .../flatten-after-unwrap/sink/0001_orders.sql | 7 ++ corpus/flatten-bare-columns/connector.json | 21 ++++ corpus/flatten-bare-columns/expected.txt | 5 + .../migrations/0001_testradiomeas.sql | 11 ++ .../flatten-bare-columns/sink-connector.json | 9 ++ .../sink/0001_testradiomeas.sql | 13 +++ .../flatten-envelope-columns/connector.json | 23 ++++ corpus/flatten-envelope-columns/expected.txt | 9 ++ .../migrations/0001_uk_price_paid.sql | 12 +++ .../sink-connector.json | 10 ++ .../sink/0001_uk_price_paid_changes.sql | 17 +++ internal/capture/debezium/debezium.go | 13 +++ internal/engine/diff.go | 4 +- internal/engine/engine.go | 102 ++++++++++++++++-- internal/model/model.go | 27 +++++ internal/sink/connect/connect.go | 10 +- internal/smt/smt.go | 94 ++++++++++++++++ internal/smt/smt_test.go | 85 +++++++++++++++ 24 files changed, 504 insertions(+), 10 deletions(-) create mode 100644 corpus/flatten-after-unwrap/connector.json create mode 100644 corpus/flatten-after-unwrap/expected.txt create mode 100644 corpus/flatten-after-unwrap/migrations/0001_orders.sql create mode 100644 corpus/flatten-after-unwrap/sink-connector.json create mode 100644 corpus/flatten-after-unwrap/sink/0001_orders.sql create mode 100644 corpus/flatten-bare-columns/connector.json create mode 100644 corpus/flatten-bare-columns/expected.txt create mode 100644 corpus/flatten-bare-columns/migrations/0001_testradiomeas.sql create mode 100644 corpus/flatten-bare-columns/sink-connector.json create mode 100644 corpus/flatten-bare-columns/sink/0001_testradiomeas.sql create mode 100644 corpus/flatten-envelope-columns/connector.json create mode 100644 corpus/flatten-envelope-columns/expected.txt create mode 100644 corpus/flatten-envelope-columns/migrations/0001_uk_price_paid.sql create mode 100644 corpus/flatten-envelope-columns/sink-connector.json create mode 100644 corpus/flatten-envelope-columns/sink/0001_uk_price_paid_changes.sql create mode 100644 internal/smt/smt.go create mode 100644 internal/smt/smt_test.go diff --git a/README.md b/README.md index 543e302..1376a63 100644 --- a/README.md +++ b/README.md @@ -57,6 +57,7 @@ your repository, and cdclint reads it: | 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.`, 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 @@ -73,6 +74,7 @@ anything running. | `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.`, 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 | @@ -215,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.`, `before.`, +`source.`, `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, diff --git a/corpus/README.md b/corpus/README.md index 41fdbbe..e171a23 100644 --- a/corpus/README.md +++ b/corpus/README.md @@ -31,6 +31,9 @@ The rest exercise one path each: | `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.`, 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 | diff --git a/corpus/flatten-after-unwrap/connector.json b/corpus/flatten-after-unwrap/connector.json new file mode 100644 index 0000000..9b2c930 --- /dev/null +++ b/corpus/flatten-after-unwrap/connector.json @@ -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": "_" + } +} diff --git a/corpus/flatten-after-unwrap/expected.txt b/corpus/flatten-after-unwrap/expected.txt new file mode 100644 index 0000000..cd1f70b --- /dev/null +++ b/corpus/flatten-after-unwrap/expected.txt @@ -0,0 +1 @@ +ok: source, connector and sink agree diff --git a/corpus/flatten-after-unwrap/migrations/0001_orders.sql b/corpus/flatten-after-unwrap/migrations/0001_orders.sql new file mode 100644 index 0000000..f5d419b --- /dev/null +++ b/corpus/flatten-after-unwrap/migrations/0001_orders.sql @@ -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 +); diff --git a/corpus/flatten-after-unwrap/sink-connector.json b/corpus/flatten-after-unwrap/sink-connector.json new file mode 100644 index 0000000..c950d30 --- /dev/null +++ b/corpus/flatten-after-unwrap/sink-connector.json @@ -0,0 +1,6 @@ +{ + "connector.class": "com.clickhouse.kafka.connect.ClickHouseSinkConnector", + "database": "shop", + "topics": "shop.public.orders", + "topic2TableMap": "shop.public.orders=orders" +} diff --git a/corpus/flatten-after-unwrap/sink/0001_orders.sql b/corpus/flatten-after-unwrap/sink/0001_orders.sql new file mode 100644 index 0000000..6ca16c4 --- /dev/null +++ b/corpus/flatten-after-unwrap/sink/0001_orders.sql @@ -0,0 +1,7 @@ +CREATE TABLE shop.orders +( + `id` Int64, + `total` String +) +ENGINE = MergeTree +ORDER BY id; diff --git a/corpus/flatten-bare-columns/connector.json b/corpus/flatten-bare-columns/connector.json new file mode 100644 index 0000000..8cb27eb --- /dev/null +++ b/corpus/flatten-bare-columns/connector.json @@ -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" + } +} diff --git a/corpus/flatten-bare-columns/expected.txt b/corpus/flatten-bare-columns/expected.txt new file mode 100644 index 0000000..959c48a --- /dev/null +++ b/corpus/flatten-bare-columns/expected.txt @@ -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. + 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 diff --git a/corpus/flatten-bare-columns/migrations/0001_testradiomeas.sql b/corpus/flatten-bare-columns/migrations/0001_testradiomeas.sql new file mode 100644 index 0000000..1d44320 --- /dev/null +++ b/corpus/flatten-bare-columns/migrations/0001_testradiomeas.sql @@ -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 +); diff --git a/corpus/flatten-bare-columns/sink-connector.json b/corpus/flatten-bare-columns/sink-connector.json new file mode 100644 index 0000000..71501d8 --- /dev/null +++ b/corpus/flatten-bare-columns/sink-connector.json @@ -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" +} diff --git a/corpus/flatten-bare-columns/sink/0001_testradiomeas.sql b/corpus/flatten-bare-columns/sink/0001_testradiomeas.sql new file mode 100644 index 0000000..97e4e9c --- /dev/null +++ b/corpus/flatten-bare-columns/sink/0001_testradiomeas.sql @@ -0,0 +1,13 @@ +-- The table as first written: the columns named as they are in Postgres. +-- The flatten transform names every field after., 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`; diff --git a/corpus/flatten-envelope-columns/connector.json b/corpus/flatten-envelope-columns/connector.json new file mode 100644 index 0000000..66a906d --- /dev/null +++ b/corpus/flatten-envelope-columns/connector.json @@ -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" +} diff --git a/corpus/flatten-envelope-columns/expected.txt b/corpus/flatten-envelope-columns/expected.txt new file mode 100644 index 0000000..26ecad7 --- /dev/null +++ b/corpus/flatten-envelope-columns/expected.txt @@ -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 diff --git a/corpus/flatten-envelope-columns/migrations/0001_uk_price_paid.sql b/corpus/flatten-envelope-columns/migrations/0001_uk_price_paid.sql new file mode 100644 index 0000000..96ae38e --- /dev/null +++ b/corpus/flatten-envelope-columns/migrations/0001_uk_price_paid.sql @@ -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. and after., 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 +); diff --git a/corpus/flatten-envelope-columns/sink-connector.json b/corpus/flatten-envelope-columns/sink-connector.json new file mode 100644 index 0000000..5ab233d --- /dev/null +++ b/corpus/flatten-envelope-columns/sink-connector.json @@ -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" +} diff --git a/corpus/flatten-envelope-columns/sink/0001_uk_price_paid_changes.sql b/corpus/flatten-envelope-columns/sink/0001_uk_price_paid_changes.sql new file mode 100644 index 0000000..a3c8ab7 --- /dev/null +++ b/corpus/flatten-envelope-columns/sink/0001_uk_price_paid_changes.sql @@ -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(); diff --git a/internal/capture/debezium/debezium.go b/internal/capture/debezium/debezium.go index 5b6dc31..e18b91e 100644 --- a/internal/capture/debezium/debezium.go +++ b/internal/capture/debezium/debezium.go @@ -19,6 +19,7 @@ import ( "strings" "github.com/avison9/cdclint/internal/model" + "github.com/avison9/cdclint/internal/smt" ) // Contract implements model.Contract for a Postgres, MySQL or SQL Server @@ -35,6 +36,7 @@ type Contract struct { columnPatterns []string tablePatterns []string routers []router + shape model.Shape Config map[string]string Class string Name string @@ -126,6 +128,7 @@ func Parse(b []byte) (*Contract, error) { c.tableListKey = "table.include.list" } c.routers = routers(cfg) + c.shape = smt.Apply(cfg, "", smt.Start(cfg)) return c, nil } @@ -195,6 +198,16 @@ func anyMatch(res []*regexp.Regexp, s string) bool { func (c *Contract) Pos() model.Pos { return c.pos } +// ValueShape is what the change events look like after this connector's own +// transforms, for the engine to match sink columns against. +func (c *Contract) ValueShape() model.Shape { + s := c.shape + if s.Kind == model.ShapeFlattened { + s.File = c.pos.File + } + return s +} + func (c *Contract) CapturesTable(schema, table string) bool { q := schema + "." + table if len(c.tableInclude) > 0 { diff --git a/internal/engine/diff.go b/internal/engine/diff.go index bcc90d6..154bb3a 100644 --- a/internal/engine/diff.go +++ b/internal/engine/diff.go @@ -79,8 +79,8 @@ func schemaBeforeConnector(in *Input, reads []Read) ([]model.Finding, map[string } read := map[string]bool{} for _, r := range reads { - if r.Table != nil { - read[strings.ToLower(r.Table.Qualified()+"."+r.Column.Name)] = true + if r.Table != nil && !r.Unreached { + read[strings.ToLower(r.Table.Qualified()+"."+r.Field)] = true } } var fs []model.Finding diff --git a/internal/engine/engine.go b/internal/engine/engine.go index 199eb36..f312090 100644 --- a/internal/engine/engine.go +++ b/internal/engine/engine.go @@ -44,6 +44,49 @@ type Read struct { Mapping model.Mapping Topic string Mapper *connect.Mapper + // Field is the source column the sink column receives: its own name, or + // the column under after or before in a flattened envelope. It is + // empty for a column no field reaches (see Unreached). + Field string + // Unreached marks a column named after a source column that receives + // nothing, because the value arrives under after instead. + Unreached bool + Shape model.Shape +} + +// ShapeLister is implemented by contracts that know what their events look +// like after their own transforms. +type ShapeLister interface { + ValueShape() model.Shape +} + +// envelopeFields are the flattened envelope's own fields besides before and +// after (Debezium's change event structure; ts_us and ts_ns since 2.6). +var envelopeFields = map[string]bool{"op": true, "ts_ms": true, "ts_us": true, "ts_ns": true} + +// field works out which source column a sink column receives under shape. +// skip is true for fields the connector or sink adds around the row. +func field(shape model.Shape, table *model.Table, name string) (col string, unreached, skip bool) { + if isMetadata(name) { + return "", false, true + } + if shape.Kind != model.ShapeFlattened { + return name, false, false + } + d := shape.Delimiter + lower := strings.ToLower(name) + for _, image := range []string{"after", "before"} { + if strings.HasPrefix(lower, image+d) { + return name[len(image)+len(d):], false, false + } + } + if envelopeFields[lower] || strings.HasPrefix(lower, "source"+d) || strings.HasPrefix(lower, "transaction"+d) { + return "", false, true + } + if table.Column(name) != nil { + return "", true, false + } + return name, false, false } // isMetadata reports whether a sink column is one the connector or the @@ -64,6 +107,10 @@ func isMetadata(name string) bool { // a topic nothing produces. func resolve(in *Input) []Read { var reads []Read + var source model.Shape + if sl, ok := in.Contract.(ShapeLister); ok { + source = sl.ValueShape() + } byTopic := map[string]*model.Table{} for _, t := range in.Source.Tables { byTopic[in.Contract.Topic(t.Schema, t.Name)] = t @@ -112,11 +159,17 @@ func resolve(in *Input) []Read { } continue } + shape := source + if mapper != nil { + shape = mapper.Reshape(source) + } for _, c := range st.Columns { - if isMetadata(c.Name) { + f, unreached, skip := field(shape, table, c.Name) + if skip { continue } - reads = append(reads, Read{Sink: st, Column: c, Table: table, Mapping: mapping, Topic: topic, Mapper: mapper}) + reads = append(reads, Read{Sink: st, Column: c, Table: table, Mapping: mapping, Topic: topic, Mapper: mapper, + Field: f, Unreached: unreached, Shape: shape}) } } } @@ -131,6 +184,7 @@ func Run(in *Input) []model.Finding { fs = append(fs, sinkTableNotCaptured(in, reads)...) fs = append(fs, sinkColumnNotCaptured(in, reads)...) fs = append(fs, sinkColumnUnknown(in, reads)...) + fs = append(fs, sinkColumnFlattened(reads)...) // The diff rule runs before the inventory so a column it raised as a // warning is not listed again as information. diff, raised := schemaBeforeConnector(in, reads) @@ -216,8 +270,8 @@ func sinkColumnNotCaptured(in *Input, reads []Read) []model.Finding { if r.Table == nil || !in.Contract.CapturesTable(r.Table.Schema, r.Table.Name) { continue } - src := r.Table.Column(r.Column.Name) - if src == nil || in.Contract.CapturesColumn(r.Table.Schema, r.Table.Name, src.Name) { + src := r.Table.Column(r.Field) + if r.Unreached || src == nil || in.Contract.CapturesColumn(r.Table.Schema, r.Table.Name, src.Name) { continue } q := fmt.Sprintf("%s.%s.%s", r.Table.Schema, r.Table.Name, src.Name) @@ -230,10 +284,44 @@ func sinkColumnNotCaptured(in *Input, reads []Read) []model.Finding { return fs } +// sinkColumnFlattened raises a sink table whose columns carry the row's +// names while a flatten transform delivers the envelope: every field arrives +// as aftercolumn, the sink matches by name, and the columns are left at +// their defaults (ClickHouse/clickhouse-kafka-connect discussion 182). One +// finding per table, since one transform causes all of them. +func sinkColumnFlattened(reads []Read) []model.Finding { + var fs []model.Finding + var order []*model.SinkTable + names := map[*model.SinkTable][]string{} + first := map[*model.SinkTable]Read{} + for _, r := range reads { + if !r.Unreached { + continue + } + if _, ok := names[r.Sink]; !ok { + order = append(order, r.Sink) + first[r.Sink] = r + } + names[r.Sink] = append(names[r.Sink], r.Column.Name) + } + for _, st := range order { + r := first[st] + cols := names[st] + fs = append(fs, model.Finding{ + Rule: "sink-column-flattened", Severity: model.Error, Pos: st.Pos, + Message: fmt.Sprintf("%s (%s) reads %s by their names in %s.%s, but %s in %s flattens each change event, so every one arrives as after%s\nevery row will carry the columns' defaults, with no error anywhere", + st.Name, how(r), strings.Join(cols, ", "), r.Table.Schema, r.Table.Name, r.Shape.Transform, r.Shape.File, r.Shape.Delimiter), + Fix: fmt.Sprintf("rename them after%s%s and so on, or unwrap the event with io.debezium.transforms.ExtractNewRecordState in place of %s", + r.Shape.Delimiter, cols[0], r.Shape.Transform), + }) + } + return fs +} + func sinkColumnUnknown(in *Input, reads []Read) []model.Finding { var fs []model.Finding for _, r := range reads { - if r.Table == nil || r.Table.Column(r.Column.Name) != nil { + if r.Table == nil || r.Unreached || r.Table.Column(r.Field) != nil { continue } fs = append(fs, model.Finding{ @@ -248,8 +336,8 @@ func sourceColumnNotCaptured(in *Input, reads []Read, raised map[string]bool) [] var fs []model.Finding declared := map[string]bool{} for _, r := range reads { - if r.Table != nil { - declared[strings.ToLower(r.Table.Qualified()+"."+r.Column.Name)] = true + if r.Table != nil && !r.Unreached { + declared[strings.ToLower(r.Table.Qualified()+"."+r.Field)] = true } } for _, t := range in.Source.Tables { diff --git a/internal/model/model.go b/internal/model/model.go index 3eb70c4..df8cbb3 100644 --- a/internal/model/model.go +++ b/internal/model/model.go @@ -197,3 +197,30 @@ func (f Finding) Format() string { } return b.String() } + +// ShapeKind says what a change event's value looks like when a sink reads it. +type ShapeKind int + +const ( + // ShapeUnknown: a transform cdclint does not model changed the value, or + // the connector is not one whose output it knows. Sink columns are taken + // as the row's column names, as they always were. + ShapeUnknown ShapeKind = iota + // ShapeEnvelope: Debezium's change event, before/after/source/op/ts_ms. + ShapeEnvelope + // ShapeRow: the row itself, unwrapped from the envelope. + ShapeRow + // ShapeFlattened: the envelope flattened into aftercolumn, + // beforecolumn, sourcefield, op, ts_ms. + ShapeFlattened +) + +// Shape is the value a sink receives after every transform on the way. +type Shape struct { + Kind ShapeKind + // Delimiter joins envelope and column names when Kind is ShapeFlattened. + Delimiter string + // Transform and File name the flatten transform, for findings. + Transform string + File string +} diff --git a/internal/sink/connect/connect.go b/internal/sink/connect/connect.go index 0ac38c7..2d2907f 100644 --- a/internal/sink/connect/connect.go +++ b/internal/sink/connect/connect.go @@ -13,6 +13,7 @@ import ( "strings" "github.com/avison9/cdclint/internal/model" + "github.com/avison9/cdclint/internal/smt" ) // Kind is which sink connector family a config belongs to. @@ -37,6 +38,7 @@ type Mapper struct { setting string // which setting held the map routes []route // Iceberg: table -> route regex single string // Iceberg: the one table when iceberg.tables lists one + cfg map[string]string // the whole config, for its transforms } type route struct { @@ -44,6 +46,12 @@ type route struct { re *regexp.Regexp } +// Reshape applies this sink connector's transforms to the value the source +// connector produced; a sink connector can unwrap or flatten on its side. +func (m *Mapper) Reshape(in model.Shape) model.Shape { + return smt.Apply(m.cfg, m.Pos.File, in) +} + // ReadFile parses one sink connector JSON file. func ReadFile(path string) (*Mapper, error) { b, err := os.ReadFile(path) @@ -72,7 +80,7 @@ func Parse(b []byte) (*Mapper, error) { for k, v := range src { cfg[k] = fmt.Sprint(v) } - m := &Mapper{Class: cfg["connector.class"], explicit: map[string]string{}} + m := &Mapper{Class: cfg["connector.class"], explicit: map[string]string{}, cfg: cfg} lc := strings.ToLower(m.Class) switch { case strings.Contains(lc, "snowflake"): diff --git a/internal/smt/smt.go b/internal/smt/smt.go new file mode 100644 index 0000000..99b6f8a --- /dev/null +++ b/internal/smt/smt.go @@ -0,0 +1,94 @@ +// Package smt works out what a change event's value looks like after a +// connector's single message transforms, so a sink column can be matched +// to the field that actually arrives. +// +// Only transforms whose effect on field names is documented are modelled: +// Debezium's ExtractNewRecordState (and its old name UnwrapFromEnvelope) and +// the Iceberg DebeziumTransform, which unwrap the envelope into the row, and +// Kafka Connect's Flatten, which joins nested names with a delimiter ("." by +// default). Transforms that only route (RegexRouter, TimestampRouter, +// ByLogicalTableRouter) or drop whole records (Filter) change no field name. +// Anything else that touches the value, or a modelled transform applied +// under a predicate to only some records, makes the shape unknown, and an +// unknown shape produces no finding. +package smt + +import ( + "strings" + + "github.com/avison9/cdclint/internal/model" +) + +// Start is the value a source connector emits before its transforms. It is +// Debezium's envelope for a Debezium connector, and for a connector that +// states it with after.state.only (Confluent's managed Postgres CDC source +// keeps the envelope when it is false and sends only the row when it is +// true). Anything else is unknown. +func Start(cfg map[string]string) model.Shape { + switch strings.ToLower(strings.TrimSpace(cfg["after.state.only"])) { + case "false": + return model.Shape{Kind: model.ShapeEnvelope} + case "true": + return model.Shape{Kind: model.ShapeRow} + } + if strings.HasPrefix(cfg["connector.class"], "io.debezium.") { + return model.Shape{Kind: model.ShapeEnvelope} + } + return model.Shape{} +} + +// Apply walks cfg's transforms in the order transforms= lists them. +func Apply(cfg map[string]string, file string, in model.Shape) model.Shape { + s := in + for _, name := range strings.Split(cfg["transforms"], ",") { + name = strings.TrimSpace(name) + if name == "" || s.Kind == model.ShapeUnknown { + continue + } + key := "transforms." + name + typ := strings.TrimSpace(cfg[key+".type"]) + conditional := cfg[key+".predicate"] != "" + switch { + case isRouter(typ), strings.HasSuffix(typ, "$Key"): + // Routes the record or changes the key; the value is untouched. + case typ == "org.apache.kafka.connect.transforms.Filter": + // Drops whole records; the ones that pass are unchanged. + case isUnwrap(typ) && !conditional: + if s.Kind == model.ShapeEnvelope { + s = model.Shape{Kind: model.ShapeRow} + } else { + s = model.Shape{} + } + case typ == "org.apache.kafka.connect.transforms.Flatten$Value" && !conditional: + if s.Kind == model.ShapeEnvelope { + d := cfg[key+".delimiter"] + if d == "" { + d = "." + } + s = model.Shape{Kind: model.ShapeFlattened, Delimiter: d, Transform: key, File: file} + } + // A row is already flat, and a flattened envelope has nothing + // left to flatten: either way the names do not change. + default: + s = model.Shape{} + } + } + return s +} + +func isRouter(typ string) bool { + return strings.HasSuffix(typ, "RegexRouter") || + typ == "org.apache.kafka.connect.transforms.TimestampRouter" || + typ == "io.debezium.transforms.ByLogicalTableRouter" +} + +func isUnwrap(typ string) bool { + switch typ { + case "io.debezium.transforms.ExtractNewRecordState", + "io.debezium.transforms.UnwrapFromEnvelope", + "io.tabular.iceberg.connect.transforms.DebeziumTransform", + "org.apache.iceberg.connect.transforms.DebeziumTransform": + return true + } + return false +} diff --git a/internal/smt/smt_test.go b/internal/smt/smt_test.go new file mode 100644 index 0000000..792e63c --- /dev/null +++ b/internal/smt/smt_test.go @@ -0,0 +1,85 @@ +package smt + +import ( + "testing" + + "github.com/avison9/cdclint/internal/model" +) + +const ( + debezium = "io.debezium.connector.postgresql.PostgresConnector" + unwrap = "io.debezium.transforms.ExtractNewRecordState" + flatten = "org.apache.kafka.connect.transforms.Flatten$Value" + router = "org.apache.kafka.connect.transforms.RegexRouter" +) + +func TestShapeAfterTransforms(t *testing.T) { + for _, tc := range []struct { + name string + cfg map[string]string + want model.Shape + }{ + {"a Debezium connector with no transforms emits the envelope", + map[string]string{"connector.class": debezium}, + model.Shape{Kind: model.ShapeEnvelope}}, + {"after.state.only false keeps the envelope whatever the class", + map[string]string{"connector.class": "PostgresCdcSource", "after.state.only": "false"}, + model.Shape{Kind: model.ShapeEnvelope}}, + {"after.state.only true sends the row", + map[string]string{"connector.class": "PostgresCdcSource", "after.state.only": "true"}, + model.Shape{Kind: model.ShapeRow}}, + {"a connector cdclint does not know is unknown", + map[string]string{"connector.class": "PostgresCdcSource"}, + model.Shape{}}, + {"unwrap turns the envelope into the row", + map[string]string{"connector.class": debezium, "transforms": "unwrap", "transforms.unwrap.type": unwrap}, + model.Shape{Kind: model.ShapeRow}}, + {"flatten on the envelope uses . by default", + map[string]string{"connector.class": debezium, "transforms": "flat", "transforms.flat.type": flatten}, + model.Shape{Kind: model.ShapeFlattened, Delimiter: ".", Transform: "transforms.flat", File: "c.json"}}, + {"flatten keeps its configured delimiter, and routers around it change nothing", + map[string]string{"connector.class": debezium, "transforms": "route, flat ,route2", + "transforms.route.type": router, "transforms.flat.type": flatten, "transforms.flat.delimiter": "_", + "transforms.route2.type": "io.confluent.connect.cloud.transforms.TopicRegexRouter"}, + model.Shape{Kind: model.ShapeFlattened, Delimiter: "_", Transform: "transforms.flat", File: "c.json"}}, + {"flatten after unwrap leaves the row as it is", + map[string]string{"connector.class": debezium, "transforms": "unwrap,flat", + "transforms.unwrap.type": unwrap, "transforms.flat.type": flatten}, + model.Shape{Kind: model.ShapeRow}}, + {"unwrap after flatten is not modelled", + map[string]string{"connector.class": debezium, "transforms": "flat,unwrap", + "transforms.unwrap.type": unwrap, "transforms.flat.type": flatten}, + model.Shape{}}, + {"a flatten under a predicate applies to some records only, so the shape is unknown", + map[string]string{"connector.class": debezium, "transforms": "flat", + "transforms.flat.type": flatten, "transforms.flat.predicate": "isOrders"}, + model.Shape{}}, + {"a transform that renames value fields makes the shape unknown", + map[string]string{"connector.class": debezium, "transforms": "flat,rename", + "transforms.flat.type": flatten, "transforms.rename.type": "org.apache.kafka.connect.transforms.ReplaceField$Value"}, + model.Shape{}}, + {"key and filter transforms leave the value alone", + map[string]string{"connector.class": debezium, "transforms": "key,drop,flat", + "transforms.key.type": "org.apache.kafka.connect.transforms.ExtractField$Key", + "transforms.drop.type": "org.apache.kafka.connect.transforms.Filter", + "transforms.flat.type": flatten}, + model.Shape{Kind: model.ShapeFlattened, Delimiter: ".", Transform: "transforms.flat", File: "c.json"}}, + } { + t.Run(tc.name, func(t *testing.T) { + if got := Apply(tc.cfg, "c.json", Start(tc.cfg)); got != tc.want { + t.Errorf("got %+v, want %+v", got, tc.want) + } + }) + } +} + +func TestSinkTransformsApplyAfterTheSource(t *testing.T) { + sink := map[string]string{"transforms": "unwrap", "transforms.unwrap.type": unwrap} + if got := Apply(sink, "sink.json", model.Shape{Kind: model.ShapeEnvelope}); got.Kind != model.ShapeRow { + t.Errorf("unwrap on the sink side: got %+v, want the row", got) + } + flat := map[string]string{"transforms": "f", "transforms.f.type": flatten} + if got := Apply(flat, "sink.json", model.Shape{Kind: model.ShapeRow}); got.Kind != model.ShapeRow { + t.Errorf("flatten on the sink side of an unwrapped row: got %+v, want the row", got) + } +}