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
44 changes: 38 additions & 6 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -12,7 +12,7 @@ GitHub Marketplace.
A change-data-capture pipeline has three schemas that must agree and nothing
that makes them agree:

1. **The source schema.** Postgres tables, evolved by migrations, changed
1. **The source schema.** Postgres or MySQL tables, evolved by migrations, changed
weekly by the application team.
2. **The capture contract.** The Debezium connector JSON: `table.include.list`,
`column.include.list`, replica identity, topic naming. Written once by
Expand Down Expand Up @@ -51,14 +51,16 @@ your repository, and cdclint reads it:
| what you see | what is usually wrong | rule |
|---|---|---|
| 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` |
| A column added in Postgres or MySQL 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
It reads files only: Postgres and MySQL (or MariaDB) sources, no type checks, no connection to
anything running.

## What it checks
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 @@ -198,6 +202,23 @@ line; `fail-on`, `version` and `working-directory` are the other inputs. The
action downloads the release binary for the runner and verifies its
checksum before running it.

## Sources

The migrations are read in filename order, and the connector's class picks
the dialect: `MySqlConnector` and `MariaDbConnector` read MySQL, every other
Debezium source reads Postgres. `--migrations mysql:DIR` or `postgres:DIR`
says it outright.

MySQL has no schemas: Debezium names a table `database.table`, so cdclint
needs the database the migrations run in. It takes it from a `USE db;` in the
migrations, or from the connector when `database.include.list` names one
database (or every `table.include.list` entry starts with the same one), and
stops with that advice when neither says. `database.include.list` and
`database.exclude.list` are applied before the table lists, as Debezium does,
and a finding names whichever list keeps a table out. MySQL's re-runnable
migrations wrap `ALTER TABLE` in a string run through `PREPARE`, because MySQL
has no `ADD COLUMN IF NOT EXISTS`; cdclint reads the DDL in those strings too.

## Sinks, out of the box in v1

| sink | how tables are found | how a table maps to a topic |
Expand All @@ -210,14 +231,25 @@ checksum before running it.

Without a sink connector config, a sink table maps to the source table of the
same name, and the finding says so. Debezium's `RegexRouter` transform is
applied when computing topic names. Other sources (MySQL, SQL Server) and
applied when computing topic names. Other sources (SQL Server, Oracle) 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,
not monitoring.
- **Pluggable ends.** Postgres and Debezium as the source and capture
- **Pluggable ends.** Postgres or MySQL, and Debezium, as the source and capture
readers; ClickHouse, BigQuery, Snowflake and Iceberg as sink readers, each
behind a small interface so the next one is a package, not a rewrite.
- **The corpus is the spec.** [`corpus/`](corpus/) holds one directory per
Expand Down
6 changes: 3 additions & 3 deletions cmd/cdclint/corpus_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -9,7 +9,7 @@ import (
"testing"

"github.com/avison9/cdclint/internal/engine"
"github.com/avison9/cdclint/internal/source/postgres"
"github.com/avison9/cdclint/internal/source"
)

// update rewrites every expected.txt from the current output. Run it after
Expand Down Expand Up @@ -110,13 +110,13 @@ func baseFromDir(t *testing.T, dir, headConnector string) *engine.Base {
if err != nil {
t.Fatal(err)
}
var files []postgres.NamedFile
var files []source.NamedFile
for _, e := range entries {
b, err := os.ReadFile(filepath.Join(dir, "migrations", e.Name()))
if err != nil {
t.Fatal(err)
}
files = append(files, postgres.NamedFile{Path: filepath.Join(dir, "migrations", e.Name()), Text: string(b)})
files = append(files, source.NamedFile{Path: filepath.Join(dir, "migrations", e.Name()), Text: string(b)})
}
sort.Slice(files, func(i, j int) bool { return files[i].Path < files[j].Path })
// No connector.json under base/ means the connector is new in the
Expand Down
52 changes: 37 additions & 15 deletions cmd/cdclint/main.go
Original file line number Diff line number Diff line change
Expand Up @@ -24,7 +24,7 @@ import (
"github.com/avison9/cdclint/internal/sink/clickhouse"
"github.com/avison9/cdclint/internal/sink/connect"
"github.com/avison9/cdclint/internal/sink/sqlddl"
"github.com/avison9/cdclint/internal/source/postgres"
"github.com/avison9/cdclint/internal/source"
)

// version is set by the release build with -ldflags "-X main.version=...".
Expand All @@ -42,7 +42,7 @@ func main() {
func run(args []string) int {
fs := flag.NewFlagSet("cdclint", flag.ContinueOnError)
var (
migrations = fs.String("migrations", "", "directory of source migrations, applied in name order (Postgres)")
migrations = fs.String("migrations", "", "directory of source migrations, applied in name order; Postgres or MySQL, from the connector class, or say it with postgres:DIR or mysql:DIR")
connector = fs.String("connector", "", "Debezium source connector JSON")
sinks multi
sinkConns multi
Expand Down Expand Up @@ -144,11 +144,15 @@ func load(migrations, connector string, sinks, sinkConns []string) (*engine.Inpu
// Load is exported for the corpus test, which runs the tool exactly as the
// command line does.
func Load(migrations, connector string, sinks, sinkConns []string) (*engine.Input, error) {
src, _, err := readSource(migrations)
c, err := debezium.ReadFile(connector)
if err != nil {
return nil, err
}
c, err := debezium.ReadFile(connector)
reader, dir, err := pickSource(migrations, c)
if err != nil {
return nil, err
}
src, _, err := reader.readDir(dir)
if err != nil {
return nil, err
}
Expand Down Expand Up @@ -198,55 +202,73 @@ func LoadBase(ref, migrations, connector string) (*engine.Base, error) {
if err != nil {
return nil, fmt.Errorf("--base: %w", err)
}
// The dialect and the database come from the connector as it is now:
// the base's may not exist, and the source did not change dialect.
head, err := debezium.ReadFile(connector)
if err != nil {
return nil, err
}
reader, migrations, err := pickSource(migrations, head)
if err != nil {
return nil, err
}
files, err := gitread.Dir(ref, migrations)
if err != nil {
return nil, fmt.Errorf("--base: %w", err)
}
var named []postgres.NamedFile
var named []source.NamedFile
for _, f := range files {
named = append(named, postgres.NamedFile{Path: f.Path, Text: f.Text})
named = append(named, source.NamedFile{Path: f.Path, Text: f.Text})
}
exists, err := gitread.Exists(ref, connector)
if err != nil {
return nil, fmt.Errorf("--base: %w", err)
}
if !exists {
return baseFrom(id, named, nil, true)
return baseFrom(id, reader, named, nil, true)
}
changed, err := gitread.Changed(ref, connector)
if err != nil {
return nil, fmt.Errorf("--base: %w", err)
}
if !changed {
return baseFrom(id, named, nil, false)
return baseFrom(id, reader, named, nil, false)
}
text, err := gitread.Show(ref, connector)
if err != nil {
return nil, fmt.Errorf("--base: %w", err)
}
return baseFrom(id, named, []byte(text), true)
return baseFrom(id, reader, named, []byte(text), true)
}

// BaseFromFiles builds the base from migrations already in memory and the
// connector's text at the base and now; the corpus test feeds it from a
// base/ directory, so the diff rule is tested without a repository. A nil
// baseConnector is a connector that did not exist at the base.
func BaseFromFiles(ref string, migrations []postgres.NamedFile, baseConnector, headConnector []byte) (*engine.Base, error) {
func BaseFromFiles(ref string, migrations []source.NamedFile, baseConnector, headConnector []byte) (*engine.Base, error) {
head, err := debezium.Parse(headConnector)
if err != nil {
return nil, err
}
reader, _, err := pickSource("", head)
if err != nil {
return nil, err
}
if baseConnector == nil {
return baseFrom(ref, migrations, nil, true)
return baseFrom(ref, reader, migrations, nil, true)
}
if string(baseConnector) == string(headConnector) {
return baseFrom(ref, migrations, nil, false)
return baseFrom(ref, reader, migrations, nil, false)
}
return baseFrom(ref, migrations, baseConnector, true)
return baseFrom(ref, reader, migrations, baseConnector, true)
}

// baseFrom parses the base's migrations and, when given, its connector. A
// parse failure there is an error, because the base once ran and its
// config once parsed, so the failure is in the reader and hiding it would
// hide the rule.
func baseFrom(ref string, migrations []postgres.NamedFile, connector []byte, changed bool) (*engine.Base, error) {
src, err := postgres.ReadFiles(migrations)
func baseFrom(ref string, reader sourceReader, migrations []source.NamedFile, connector []byte, changed bool) (*engine.Base, error) {
src, err := reader.readFiles(migrations)
if err != nil {
return nil, fmt.Errorf("--base %s: %w", ref, err)
}
Expand Down
49 changes: 46 additions & 3 deletions cmd/cdclint/source.go
Original file line number Diff line number Diff line change
@@ -1,12 +1,55 @@
package main

import (
"fmt"
"strings"

"github.com/avison9/cdclint/internal/capture/debezium"
"github.com/avison9/cdclint/internal/model"
"github.com/avison9/cdclint/internal/source"
"github.com/avison9/cdclint/internal/source/mysql"
"github.com/avison9/cdclint/internal/source/postgres"
)

// readSource is the one place the source dialect is chosen. Postgres is the
// only one in v1; a MySQL reader slots in here.
func readSource(dir string) (*model.Source, []string, error) {
// sourceReader is the one place the source dialect is chosen. A
// "postgres:" or "mysql:" prefix on --migrations says it outright;
// otherwise the connector's class decides, since a MySqlConnector (or
// MariaDbConnector) reads a MySQL database and every other Debezium source
// cdclint knows reads Postgres.
type sourceReader struct {
mysql bool
database string // MySQL: the database unqualified migrations run in
}

func pickSource(migrations string, c *debezium.Contract) (sourceReader, string, error) {
r := sourceReader{mysql: c.MySQL()}
dir := migrations
if i := strings.Index(migrations, ":"); i > 0 && !strings.Contains(migrations[:i], "/") {
switch strings.ToLower(migrations[:i]) {
case "mysql", "mariadb":
r.mysql, dir = true, migrations[i+1:]
case "postgres", "postgresql":
r.mysql, dir = false, migrations[i+1:]
default:
return r, "", fmt.Errorf("unknown source dialect %q in --migrations %q", migrations[:i], migrations)
}
}
if r.mysql {
r.database = c.DefaultDatabase()
}
return r, dir, nil
}

func (r sourceReader) readDir(dir string) (*model.Source, []string, error) {
if r.mysql {
return mysql.ReadDir(dir, r.database)
}
return postgres.ReadDir(dir)
}

func (r sourceReader) readFiles(files []source.NamedFile) (*model.Source, error) {
if r.mysql {
return mysql.ReadFiles(files, r.database)
}
return postgres.ReadFiles(files)
}
10 changes: 10 additions & 0 deletions corpus/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -28,6 +28,16 @@ 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 |
| `mysql-column-never-captured` | the column-never-captured incident in MySQL syntax: backticks, `KEY` lines, `ADD COLUMN ... AFTER`, a `database.include.list` naming the database the migrations run in |
| `mysql-change-column` | MySQL renames a column with `CHANGE`; the include list still names the old column and the new one stops arriving |
| `mysql-database-not-included` | a table in a second database (`USE billing`) that `database.include.list` leaves out; the fix names the database list and the database |
| `mysql-diff-adds-column` | the diff rule through the MySQL reader: `base/` without the new column, the change adds it and leaves the connector alone |
| `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;
Loading