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
4 changes: 3 additions & 1 deletion README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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 |
Expand Down
3 changes: 3 additions & 0 deletions corpus/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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 |
Expand Down
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
14 changes: 14 additions & 0 deletions corpus/include-table-typo/migrations/0001_tables.sql
Original file line number Diff line number Diff line change
@@ -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
);
12 changes: 12 additions & 0 deletions corpus/include-table-typo/sink/0001_orders.sql
Original file line number Diff line number Diff line change
@@ -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';
11 changes: 11 additions & 0 deletions internal/capture/debezium/debezium.go
Original file line number Diff line number Diff line change
Expand Up @@ -33,6 +33,7 @@ type Contract struct {
columnListKey string
tableListKey string
columnPatterns []string
tablePatterns []string
routers []router
Config map[string]string
Class string
Expand Down Expand Up @@ -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
}
Expand Down Expand Up @@ -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 }
96 changes: 93 additions & 3 deletions internal/engine/engine.go
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,8 @@ package engine

import (
"fmt"
"regexp"
"sort"
"strings"

"github.com/avison9/cdclint/internal/model"
Expand All @@ -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
Expand All @@ -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.
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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.
Expand Down