Skip to content

[core][flink][spark] Expose changelog event metadata as separate columns - #9822

Open
junmuz wants to merge 11 commits into
apache:masterfrom
junmuz:feature/expose-field-as-metadata
Open

junmuz wants to merge 11 commits into
apache:masterfrom
junmuz:feature/expose-field-as-metadata

Conversation

@junmuz

@junmuz junmuz commented Sep 14, 2026 •

Copy link
Copy Markdown
Contributor

Preserve correct before-images in changelog retraction records and append event values as configurable metadata columns, accessible via Flink SupportsReadingMetadata and Spark table schema.

Purpose

  • Added support for reading the metadata through Flink using <prefix><field> metadata columns
    such as Cassandra WRITETIME.
  • Updated Spark and core read paths to include and populate the synthetic metadata fields across raw and merged file reads.
  • Added validation for metadata projections, file reads, filters, and lookup changelog generation.

Tests

  • Added unit coverage for insert, update, delete, before-image, filter, aggregation, and unconfigured cases.
  • Added Flink integration coverage for reading event metadata and preserving correct physical before-images.
  • paimon-docs verify passes on Java 8 and Java 11.
  • Clean Java 11 Maven installation completed successfully across all 76 modules.

@junmuz
junmuz force-pushed the feature/expose-field-as-metadata branch from 7681a1c to 2bd8780 Compare September 17, 2026 10:22
Comment thread docs/docs/primary-key-table/changelog-producer.md Outdated
Comment thread docs/docs/primary-key-table/changelog-producer.md Outdated
Comment thread docs/docs/primary-key-table/changelog-producer.md Outdated
@junmuz
junmuz marked this pull request as ready for review September 17, 2026 16:53
@junmuz
junmuz marked this pull request as draft September 17, 2026 17:36
@junmuz
junmuz force-pushed the feature/expose-field-as-metadata branch from 1f3861a to c7a4fe2 Compare September 17, 2026 21:55
@junmuz
junmuz marked this pull request as ready for review September 18, 2026 05:30
@junmuz

junmuz commented Sep 18, 2026

Copy link
Copy Markdown
Contributor Author

@JingsongLi Can I get a review on this?

@JingsongLi

Copy link
Copy Markdown
Contributor

Thanks for the feature. I reviewed the write and read paths end to end and found two blocking issues plus a consistency concern.

1. [P1] The option is applied for every changelog producer, but only the lookup path fills the extra fields

MergeTreeCompactManagerFactory's constructor widens the changelog value type whenever the option is set:

RowType changelogValueType = computeChangelogValueType(valueType, options);
if (changelogValueType != null) {
    writerFactoryBuilder.withChangelogValueType(changelogValueType);
}

There is no changelog-producer == lookup check here, and the option is not validated anywhere else (docs/description only say "Only valid when changelog-producer is lookup"). The extra fields are actually populated only by the lookup merge-function wrapper (preserveFieldIndices is only passed to LookupMergeFunctionWrapperFactory inside the lookup branch of createRewriter).

Because the same KeyValueFileWriterFactory.Builder instance is shared with KeyValueFileStoreWrite (the compact-manager factory is created in KeyValueFileStoreWrite's constructor, and createWriter later calls writerFactoryBuilder.build(...)), the merge-tree writer's changelog format context is widened to N+K for all producers:

  • changelog-producer='input': MergeTreeWriter flushes changelog records built from the incoming values (MergeTreeWriter.java:219-220 -> createRollingChangelogFileWriter) whose value rows have N fields, but the writer now expects N+K value fields.
  • changelog-producer='full-compaction': same for ChangelogMergeTreeRewriter.java:147.

Result: the flush/compaction fails while serializing the shorter row (e.g. AIOOBE in the row writer), or worse, corrupt changelog files are written. Please gate the widening on options.changelogProducer() == ChangelogProducer.LOOKUP (and/or reject the option at schema validation for other producers).

2. [P1] The widened Spark DSv2 schema breaks ordinary Spark writes to such tables

PaimonSparkTableBase.schema (and BaseScan.tableRowType) add the K metadata columns to the DSv2 table schema. Paimon tables advertise ACCEPT_ANY_SCHEMA, so Spark's own ResolveOutputRelation is skipped and writes are resolved by Paimon's PaimonAnalysis -> PaimonOutputResolver. For a by-position write, SchemaEvolutionHelper.expectedAttrsForCatalogWrite returns table.output (now N+K columns), and PaimonOutputResolver.resolveColumnsByPosition throws on the size mismatch:

Cannot write to `t`, the number of data columns (3) doesn't match the table schema's (4).

So on a table with changelog-producer.expose-field-as-metadata set, plain Spark SQL writes such as INSERT INTO t SELECT id, data, event_ts FROM src or INSERT INTO t VALUES (...) fail, even though the metadata columns are generated rather than user input - the new test itself says "it is not an input column". By-name writes happen to survive (top-level NULL fill), which makes the failure mode even more confusing.

Additionally, INSERT INTO t SELECT * FROM t produces N+K columns, which resolves positionally; with spark.paimon.write.merge-schema enabled, SchemaEvolutionHelper.commitSchemaEvolution will merge __internal__event_ts into the persisted table schema, because SparkSystemColumns only filters _bucket_/_row_kind_.

Suggested fix: keep the generated metadata columns out of write-side resolution, e.g. filter them in expectedAttrsForCatalogWrite/SparkSystemColumns, or expose them only on the scan side while keeping PaimonSparkTableBase.schema at the physical N columns.

3. [P2] Inconsistent unknown-column handling and duplicated extension logic

  • MergeTreeCompactManagerFactory.computeChangelogValueType throws IllegalArgumentException for a configured column that does not exist, while KeyValueFileStore.newReaderFactoryBuilder and ChangelogEventMetadataTable.computeExtendedRowType silently skip it (if (idx >= 0)). A typo in the option therefore fails at write/compaction time but is ignored on the read path (metadata column silently NULL) - please pick one behavior, ideally validating at schema/option level.
  • The "append prefix + name fields with max(id) + 1 and nullable type copy" logic is duplicated in MergeTreeCompactManagerFactory (twice), KeyValueFileStore, ChangelogEventMetadataTable.computeExtendedRowType and again in a different variant in MergeFileSplitRead/RawFileSplitRead. These must stay in sync for the reader mapping (createMetadataFallbackMapping) to line up; a single helper would be safer.

Minor: for changelog files written before the option was enabled, the metadata columns read as NULL (the fallback mapping is only applied to data files, via !isChangelogFile(file)). Worth one sentence in the docs if that is intended.

@junmuz
junmuz force-pushed the feature/expose-field-as-metadata branch from c7a4fe2 to 5ec5070 Compare September 21, 2026 15:54

@JingsongLi JingsongLi left a comment •

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Given the significant nature of the changes, I need specific use cases to justify their validity; please provide a detailed explanation of how this would be implemented in a real-world production environment.

@junmuz

junmuz commented Sep 22, 2026 •

Copy link
Copy Markdown
Contributor Author

Given the significant nature of the changes, I need specific use cases to justify their validity; please provide a detailed explanation of how this would be implemented in a real-world production environment.

@JingsongLi Let me explain this in detail below. Yeah I agree the changes are more complex that's why I initially went with minimal change approach.

Use case and motivation

Paimon’s lookup changelog producer reconstructs UPDATE_BEFORE and DELETE records by looking up the previous table state. As a result, regular columns in those records describe the old row, not necessarily the event that caused the change. For example:

UPDATE_BEFORE: event_ts = 50, __internal__event_ts = 100
UPDATE_AFTER: event_ts = 100, __internal__event_ts = 100

Downstream CDC consumers often need the event timestamp or source sequence to apply changes deterministically.

Concurrent Sinks to a Single Cassandra table

Multiple sources can fan into the same Cassandra table while being processed by independent jobs. A common example is migrating from Kafka to Paimon: Kafka retains the original event payload for retract events, whereas Paimon’s lookup-generated changelog may represent UPDATE_BEFORE and DELETE records using the previous row image. This makes it difficult to replace a Kafka-based job with a Paimon-based job without changing downstream conflict-resolution behavior. Exposing selected fields as changelog metadata preserves the event-level values and gives the sink a consistent, globally comparable event version for conditional writes and tombstones, preventing stale updates from overwriting newer changes or resurrecting deleted records. This behavior allows creating Sinks from multiple Paimon tables merging into a single Cassandra column.

@JingsongLi JingsongLi left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

The Cassandra WRITETIME use case gives this feature concrete end-to-end value. The earlier lookup-only validation, physical before-image preservation, and Spark write-side filtering are present on this head. I found one remaining schema issue before this should ship:

P2 — Allocate synthetic metadata field IDs above nested field IDs (ChangelogEventMetadata.java:117). extraValueFields takes the maximum top-level field ID. A valid table can have id (0), payload ROW<nested INT> (top-level 1, nested field ID 3), and event_ts (2). Exposing event_ts assigns the new __internal__event_ts field ID 3, duplicating payload.nested. I added this case to ChangelogEventMetadataTest in an isolated checkout: appendMetadataFields(...).collectFieldIds(...) fails with Broken schema, field id 3 is duplicated. The generated row type is used by the changelog writer and both read integrations, so its IDs need to remain valid for nested schemas and field-ID based mapping. Please derive the next ID from the full recursive field-ID set (for example, RowType.currentHighestFieldId(valueType.getFields()) + 1) and test a nested table through a changelog read.

Validation: the added core tests passed 16/16 and Spark 3 PaimonCDCSourceTest passed 6/6 on 5ec5070; the strengthened nested-ID test fails as above. I attempted the 3 Flink integration tests locally, but all stopped on the local Avro classpath mismatch DataFileWriter.setEncoder before exercising this change. The PR's Flink and other CI jobs are green on this head.

@junmuz
junmuz force-pushed the feature/expose-field-as-metadata branch from 5ec5070 to b0fd1ea Compare September 24, 2026 10:48
@junmuz
junmuz requested a review from JingsongLi September 24, 2026 11:49
@junmuz

junmuz commented Sep 24, 2026

Copy link
Copy Markdown
Contributor Author

The Cassandra WRITETIME use case gives this feature concrete end-to-end value. The earlier lookup-only validation, physical before-image preservation, and Spark write-side filtering are present on this head. I found one remaining schema issue before this should ship:

Thanks a lot for pointing it out. I have made that fix.

@JingsongLi

Copy link
Copy Markdown
Contributor

Re-reviewed head b0fd1ea for production. The Cassandra WRITETIME use case supplies clear end-to-end value. The latest fix allocates synthetic changelog metadata IDs above the full recursive physical field-ID set, addressing my prior nested-field collision; it also adds a Flink SQL read/write regression with a nested ROW and metadata alias.

Local JDK 8 core suites passed: ChangelogEventMetadataTest 2/2, SchemaValidationTest 72/72, and LookupChangelogMergeFunctionWrapperTest 13/13. I attempted the four Flink integration cases locally, but each failed at its initial manifest write before the changed logic because this isolated checkout loads an incompatible Avro jar (DataFileWriter.setEncoder NoSuchMethodError), the same local classpath issue observed in the earlier review. Exact-head Flink 1/2, Spark, JDK 8/11, E2E, docs and licensing CI are all green. git diff --check passed. I found no remaining code blocker; keep this PR open for merge.

@JingsongLi

Copy link
Copy Markdown
Contributor

I took another look at head b0fd1ea, focusing on the public semantics and schema evolution. The Cassandra WRITETIME use case is valid, and preserving the physical before-image is the right direction. I need to revise my earlier no-blocker assessment: I found two correctness issues with the current metadata surface.

1. [P1] Event metadata breaks retract semantics when used by normal SQL operators.

LookupChangelogMergeFunctionWrapper writes the new event's metadata into the -U record while its physical columns correctly contain the old row. For example, an initial +I(event_ts=50, metadata=50) followed by an update to 100 produces -U(event_ts=50, metadata=100). In Flink, WHERE metadata < 75 forwards the insert but drops its retraction, leaving a stale row. Grouping or aggregating by the metadata column has the same problem. The tests cover filtering on the physical event_ts, but not on the exposed metadata column. This is intrinsic to treating an event attribute as an ordinary column in a retracting dynamic table; please define an event-oriented consumption surface (or another mechanism that prevents relational operators from interpreting it as before-image data) and add a regression test for this sequence.

2. [P1] Changing changelog-producer.metadata-field-prefix loses historical metadata.

The prefix is not marked immutable, so it can be changed after changelog files exist. The write path stored the old prefixed field name, but ChangelogEventMetadata.extraValueFields and the reader construct only the current name for every file schema. After changing __internal__ to __event__, older changelog files have no __event__event_ts field, so historical event metadata reads as NULL; the data-file fallback is intentionally disabled for changelog files. Please keep the on-disk identity stable across option changes, or reject prefix changes once files exist. A write → ALTER prefix → read-old-changelog test would make the contract clear.

There is also a source/sink mismatch in the documentation example: METADATA FROM '__internal__event_ts' needs to be declared on the Paimon source. An external sink's METADATA FROM key belongs to that sink and works only if the sink advertises that writable key. Please show the source alias and an explicit mapping to the sink's actual timestamp input.

Finally, the option names obscure the contract. expose-field-as-metadata persists a list of post-merge event values; it does more than expose a read-time field and does not always preserve the raw incoming value. A name such as changelog-producer.event-metadata-fields would be clearer. The configurable prefix currently couples the on-disk field name, Flink metadata key, and Spark-visible column name. I suggest separating the stable storage identity from engine-specific presentation. Given that the concrete use case is a Flink-to-Cassandra flow, I would also consider landing core/Flink support first and adding the Spark schema/write-path changes with a separate demonstrated consumer.

@junmuz
junmuz force-pushed the feature/expose-field-as-metadata branch from e201136 to c287adb Compare September 27, 2026 22:01
@junmuz

junmuz commented Sep 27, 2026 •

Copy link
Copy Markdown
Contributor Author

@JingsongLi Thanks for the review.

  1. I see your point, and I am of the opinion if someone uses it inside for filtering and aggregation then it is an incorrect configuration. I have also updated the documentation to not use the metadata field with aggregation or filtering. Also, it is quite difficult unless someone chooses to use it that way, because now they have to setup additional table property as well as exposing the metadata column. We have similar behavior with sequence fields as well if someone configures ts as sequence field, and uses `SELECT MAX(ts) as ts from source_table; then retraction doesn't work.

  2. I’ve made changelog-producer.metadata-field-prefix immutable. The stored changelog field uses the prefix with the source field ID, so its physical identity stays stable across source field renames. Spark and Flink can use the same logical metadata name through their respective APIs (Spark as a column and Flink as a metadata key) without engine-specific names.

@junmuz
junmuz force-pushed the feature/expose-field-as-metadata branch from c287adb to 2b5f23f Compare September 27, 2026 22:28
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.

2 participants