Skip to content

Sink columns are matched to the field that actually arrives after the connectors' transforms, and a table that misses a flatten is raised - #14

Merged
avison9 merged 2 commits into
mainfrom
feat/sink-flatten-shape
Sep 25, 2026
Merged

avison9 merged 2 commits into
mainfrom
feat/sink-flatten-shape

Conversation

@avison9

@avison9 avison9 commented Sep 25, 2026 •

Copy link
Copy Markdown
Owner

What it changes

cdclint took every sink column name as a row column name. That is right when the events are unwrapped (RefuseRadar's case) and wrong both ways when a Flatten transform keeps Debezium's envelope and names every field after.<column>:

  • Missed failure. ClickHouse/clickhouse-kafka-connect discussion 182 (and issue 183): a ClickHouse table with the Postgres column names received a row of zeros and empty strings for every change. On that exact configuration cdclint said ok: source, connector and sink agree. Now:

    error sink-column-flattened flatten-bare-columns/sink/0001_testradiomeas.sql:4
      default.testradiomeas (is fed topic testradiomeas by flatten-bare-columns/sink-connector.json) reads measDate, measLongitude, measLatitude, measNetNoCoverage, measNetType by their names in public.testradiomeas, but transforms.flatten in flatten-bare-columns/connector.json flattens each change event, so every one arrives as after.<column>
      every row will carry the columns' defaults, with no error anywhere
      fix: rename them after.measDate and so on, or unwrap the event with io.debezium.transforms.ExtractNewRecordState in place of transforms.flatten
    
  • False alarms, and a miss behind them. ClickHouse's own CDC example names its columns before.<column> and after.<column>. cdclint raised twelve sink-column-unknown warnings on it, and with an include list that leaves postcode2 out it only said postcode2 was not captured and "nothing reads it yet". Now after.postcode2 and before.postcode2 resolve to postcode2, and the capture rule raises the two real errors.

How

A new internal/smt package follows each connector's transforms in order:

start or transform value after it
a Debezium connector (io.debezium.*), or after.state.only: false envelope
after.state.only: true row
ExtractNewRecordState, UnwrapFromEnvelope, the Iceberg DebeziumTransform, on the envelope row
Flatten$Value on the envelope flattened, with its delimiter (. by default)
Flatten$Value on a row row (already flat)
routers (RegexRouter, TopicRegexRouter, TimestampRouter, ByLogicalTableRouter), *$Key, Filter unchanged
anything else touching the value, or a modelled transform under a predicate unknown

The sink connector's transforms apply after the source's. Only a flattened envelope changes anything: after<d>c and before<d>c resolve to source column c, op, ts_ms, ts_us, ts_ns, source<d>* and transaction<d>* are envelope metadata, and a column named as the row names it goes to the new sink-column-flattened rule, once per table. An unknown shape keeps the old behaviour exactly.

What was rejected

  • Raising a bare envelope with no transform at all. Whether the sink handles Debezium's envelope itself, or a converter reshapes it, is not in the files; raising it would be a guess (AGENTS.md rule 4).
  • Modelling ReplaceField, HoistField, Cast and friends. Possible later; for now they make the shape unknown rather than half-modelled.
  • One finding per column for the new rule: one transform causes all of them, so one per table with the columns listed.

Verified

  • Corpus first, expectations written by hand and failing before the code: flatten-bare-columns (discussion 182's configs), flatten-envelope-columns (ClickHouse/examples cdc/postgresql at ae417ae; its sink config has no connector.class, added so the file names its sink), flatten-after-unwrap (a guard, stays ok). Every existing entry is unchanged.
  • internal/smt table tests for each row above, and for sink-side transforms applying after the source's.
  • gofmt -l . clean, go vet, go build, go test ./... pass.
  • RefuseRadar (unwraps, with add.fields): 0 error(s), 0 warning(s), 159 info, unchanged.

Merging

Stacked on #13 (rebased onto main after #15 and #17 landed): until #13 merges, this PR's diff also shows #13's commit. Merge #13 first, then this; each goes in without conflicts.

Next for you

Merge when happy. Ships in the next release.

…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.
… connectors' transforms, and a table that misses a flatten is raised

A Flatten transform on the Debezium envelope names every field
after.<column>, 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.<column> and after.<column>, 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.<c> and
before.<c> 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.
@avison9
avison9 force-pushed the feat/sink-flatten-shape branch from 399b266 to 6133eb3 Compare September 25, 2026 22:23
@avison9
avison9 merged commit 2f243c6 into main Sep 25, 2026
2 checks passed
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant