diff --git a/docs/docs/primary-key-table/changelog-producer.md b/docs/docs/primary-key-table/changelog-producer.md index ecbc2755a363..c1fc24c8f77b 100644 --- a/docs/docs/primary-key-table/changelog-producer.md +++ b/docs/docs/primary-key-table/changelog-producer.md @@ -103,12 +103,51 @@ Lookup uses memory and local disk caches: | `lookup.cache-max-disk-size` | Unlimited | Bound local disk usage | | `lookup.cache-max-memory-size` | `256 mb` | Bound in-memory cache usage | -In Flink, `execution.checkpointing.max-concurrent-checkpoints` can also affect throughput when -checkpoint completion waits for compaction. Tune it with checkpoint duration and resource usage. - `lookup` is incompatible with `full-compaction.delta-commits`. For periodic full compaction with changelog generation, use `full-compaction` instead. +Set `'changelog-producer.event-metadata-fields'` to a comma-separated list of columns whose +post-merge values should be stored as event metadata fields in lookup changelog records. Metadata +fields are named by concatenating the configured prefix and column name (`__internal__` by +default). The same name is used as the Flink metadata key. For retractions (`-U`, `-D`), regular +columns contain the before-image while metadata fields contain the event values; for forward records +(`+I`, `+U`), they mirror the regular values. The post-merge values may differ from the raw incoming +row when the merge engine aggregates values. External sinks that need event timestamps for conflict +resolution can read these metadata fields. + +These fields are intended to be passed through to downstream sinks. Do not use them in filters or +aggregations: those operations can remove or combine changelog records, including retractions, and +leave downstream sinks with an incomplete changelog. + +Paimon readers such as Spark expose these generated fields as regular columns using the configured +names. Flink SQL must declare the field as a metadata column on the Paimon source, for example +`METADATA FROM '__internal__event_ts'` with the default prefix. The Flink column alias is not a +physical Paimon column and is not automatically visible to Spark. + +This option is supported only by the `lookup` changelog producer. Set +`'changelog-producer.metadata-field-prefix'` if the default prefix conflicts with an existing column +name. Changelog files written before this option was enabled expose these metadata fields as `NULL`. + +```sql +-- Source table with event metadata preservation +CREATE TABLE my_table ( + id INT PRIMARY KEY NOT ENFORCED, + data STRING, + event_ts BIGINT, + source_event_ts BIGINT METADATA FROM '__internal__event_ts' VIRTUAL +) WITH ( + 'changelog-producer' = 'lookup', + 'sequence.field' = 'event_ts', + 'changelog-producer.event-metadata-fields' = 'event_ts' +); + +-- external_sink is defined by its datastore connector and has a writable event_ts input. +``` + +The Paimon source key `__internal__event_ts` populates the source alias `source_event_ts`. Map that +alias to the external datastore's writable `event_ts` input. Any writable metadata key on the sink +is separate and must be advertised by that sink connector. + ## Full Compaction ```sql diff --git a/docs/generated/core_configuration.html b/docs/generated/core_configuration.html index a1c67478ac1d..b3d6246dca4e 100644 --- a/docs/generated/core_configuration.html +++ b/docs/generated/core_configuration.html @@ -212,6 +212,12 @@

Enum

Whether to double write to a changelog file. This changelog file keeps the details of data changes, it can be read directly during stream reads. This can be applied to tables with primary keys.

Possible values: + +
changelog-producer.event-metadata-fields
+ (none) + String + A comma-separated list of column names whose post-merge values should be stored as event metadata fields in lookup changelog records. In retraction records (-U, -D), regular value columns retain the correct before-image while the event metadata fields contain the event values. In forward records (+I, +U), the event metadata fields mirror the regular values. The post-merge values may differ from the raw input row when the merge engine aggregates values. The event metadata fields can be read as Flink metadata columns by sinks that need event timestamps for conflict resolution. Only valid when changelog-producer is lookup. +
changelog-producer.ignore-delete
false @@ -224,6 +230,12 @@ Boolean Whether to ignore update-before records in the changelog. When set to true, UPDATE_BEFORE (-U) records will not be written to changelog files. This configuration is only valid for the changelog-producer is lookup or full-compaction. + +
changelog-producer.metadata-field-prefix
+ "__internal__" + String + The prefix used for naming the extra metadata columns created by 'changelog-producer.event-metadata-fields'. For example, with the default prefix '__internal__' and a preserved column 'event_ts', the metadata column is named '__internal__event_ts'. That name is also the Flink readable metadata key. The changelog storage column uses the same prefix with the source field ID, keeping its name stable after a rename. Change this if the default prefix conflicts with existing column names. The prefix cannot be changed or reset after the table has snapshots. +
changelog-producer.row-deduplicate
false @@ -398,12 +410,6 @@ MemorySize When incremental size is bigger than this threshold, force a full compaction. - -
continuous-compaction.initial-scan-mode
- earliest -

Enum

- Initial snapshot mode for dedicated streaming compaction. When set to 'earliest' (the default), compaction starts from the earliest available snapshot if no COMPACT snapshot exists; when a COMPACT snapshot exists, compaction always resumes from the snapshot after it. When set to 'latest', the latest snapshot is read in ALL mode as the initial baseline and subsequent scans start from the next snapshot. The 'latest' mode skips historical snapshot changes and should only be used when historical changelog replay is not required.

Possible values: -
compaction.max-size-amplification-percent
200 @@ -494,6 +500,12 @@

Enum

Specify the consumer consistency mode for table.

Possible values: + +
continuous-compaction.initial-scan-mode
+ earliest +

Enum

+ Initial snapshot mode for dedicated streaming compaction. When set to 'earliest' (the default), compaction starts from the earliest available snapshot if no COMPACT snapshot exists; when a COMPACT snapshot exists, compaction always resumes from the snapshot after it. When set to 'latest', the latest snapshot is read in ALL mode as the initial baseline and subsequent scans start from the next snapshot. The 'latest' mode skips historical snapshot changes and should only be used when historical changelog replay is not required.

Possible values: +
continuous.discovery-interval
10 s diff --git a/paimon-api/src/main/java/org/apache/paimon/CoreOptions.java b/paimon-api/src/main/java/org/apache/paimon/CoreOptions.java index 98f0dfe048e6..b9a81a26e1dd 100644 --- a/paimon-api/src/main/java/org/apache/paimon/CoreOptions.java +++ b/paimon-api/src/main/java/org/apache/paimon/CoreOptions.java @@ -1142,6 +1142,40 @@ public InlineElement getDescription() { .withDescription( "Fields that are ignored for comparison while generating -U, +U changelog for the same record. This configuration is only valid for the changelog-producer.row-deduplicate is true."); + public static final ConfigOption CHANGELOG_PRODUCER_EVENT_METADATA_FIELDS = + key("changelog-producer.event-metadata-fields") + .stringType() + .noDefaultValue() + .withDescription( + "A comma-separated list of column names whose post-merge values should " + + "be stored as event metadata fields in lookup changelog records. " + + "In retraction records (-U, -D), regular value columns retain " + + "the correct before-image while the event metadata fields " + + "contain the event values. In forward records (+I, +U), the " + + "event metadata fields mirror the regular values. The " + + "post-merge values may differ from the raw input row when the " + + "merge engine aggregates values. The event metadata fields " + + "can be read as Flink metadata columns by sinks that need event " + + "timestamps for conflict resolution. " + + "Only valid when changelog-producer is lookup."); + + @Immutable + public static final ConfigOption CHANGELOG_PRODUCER_METADATA_FIELD_PREFIX = + key("changelog-producer.metadata-field-prefix") + .stringType() + .defaultValue("__internal__") + .withDescription( + "The prefix used for naming the extra metadata columns created by " + + "'changelog-producer.event-metadata-fields'. For example, " + + "with the default prefix '__internal__' and a preserved column " + + "'event_ts', the metadata column is named '__internal__event_ts'. " + + "That name is also the Flink readable metadata key. The " + + "changelog storage column uses the same prefix with the " + + "source field ID, keeping its name stable after a rename. " + + "Change this if the default prefix conflicts with existing " + + "column names. The prefix cannot be changed or reset after " + + "the table has snapshots."); + public static final ConfigOption TABLE_READ_SEQUENCE_NUMBER_ENABLED = key("table-read.sequence-number.enabled") .booleanType() @@ -4033,6 +4067,20 @@ public List changelogRowDeduplicateIgnoreFields() { .orElse(Collections.emptyList()); } + public List changelogEventMetadataFields() { + return options.getOptional(CHANGELOG_PRODUCER_EVENT_METADATA_FIELDS) + .map( + s -> + Arrays.stream(s.split(",")) + .map(String::trim) + .collect(Collectors.toList())) + .orElse(Collections.emptyList()); + } + + public String changelogMetadataFieldPrefix() { + return options.get(CHANGELOG_PRODUCER_METADATA_FIELD_PREFIX); + } + public boolean tableReadSequenceNumberEnabled() { return options.get(TABLE_READ_SEQUENCE_NUMBER_ENABLED); } diff --git a/paimon-core/src/main/java/org/apache/paimon/KeyValueFileStore.java b/paimon-core/src/main/java/org/apache/paimon/KeyValueFileStore.java index 8f84f5e504e5..e8671a9473f6 100644 --- a/paimon-core/src/main/java/org/apache/paimon/KeyValueFileStore.java +++ b/paimon-core/src/main/java/org/apache/paimon/KeyValueFileStore.java @@ -40,6 +40,7 @@ import org.apache.paimon.schema.TableSchema; import org.apache.paimon.table.BucketMode; import org.apache.paimon.table.CatalogEnvironment; +import org.apache.paimon.table.system.ChangelogEventMetadata; import org.apache.paimon.types.RowType; import org.apache.paimon.utils.KeyComparatorSupplier; import org.apache.paimon.utils.UserDefinedSeqComparator; @@ -87,6 +88,7 @@ public KeyValueFileStore( options.changelogRowDeduplicate() ? ValueEqualiserSupplier.fromIgnoreFields(valueType, logDedupIgnoreFields) : () -> null; + ChangelogEventMetadata.validate(valueType, options); } @Override @@ -126,16 +128,26 @@ public RawFileSplitRead newBatchRawFileRead() { } public KeyValueFileReaderFactory.Builder newReaderFactoryBuilder() { - return KeyValueFileReaderFactory.builder( - fileIO, - schemaManager, - schema, - keyType, - valueType, - FileFormatDiscover.of(options), - pathFactory(), - keyValueFieldsExtractor, - options); + KeyValueFileReaderFactory.Builder builder = + KeyValueFileReaderFactory.builder( + fileIO, + schemaManager, + schema, + keyType, + valueType, + FileFormatDiscover.of(options), + pathFactory(), + keyValueFieldsExtractor, + options); + if (options.changelogProducer() == CoreOptions.ChangelogProducer.LOOKUP + && !options.changelogEventMetadataFields().isEmpty()) { + List extraFields = + ChangelogEventMetadata.storageValueFields(valueType, options); + if (!extraFields.isEmpty()) { + builder.withChangelogExtraValueFields(extraFields); + } + } + return builder; } @Override diff --git a/paimon-core/src/main/java/org/apache/paimon/io/ChainKeyValueFileReaderFactory.java b/paimon-core/src/main/java/org/apache/paimon/io/ChainKeyValueFileReaderFactory.java index ff07c95fe097..cdd892820e25 100644 --- a/paimon-core/src/main/java/org/apache/paimon/io/ChainKeyValueFileReaderFactory.java +++ b/paimon-core/src/main/java/org/apache/paimon/io/ChainKeyValueFileReaderFactory.java @@ -62,7 +62,8 @@ public ChainKeyValueFileReaderFactory( DeletionVector.Factory dvFactory, ChainReadContext chainReadContext, CoreOptions coreOptions, - @Nullable ReadBatchSizer readBatchSizer) { + @Nullable ReadBatchSizer readBatchSizer, + @Nullable int[] metadataFallbackMapping) { super( fileIO, schemaManager, @@ -74,7 +75,8 @@ public ChainKeyValueFileReaderFactory( partition, dvFactory, coreOptions, - readBatchSizer); + readBatchSizer, + metadataFallbackMapping); this.chainReadContext = chainReadContext; CoreOptions options = new CoreOptions(schema.options()); this.currentBranch = options.branch(); @@ -123,7 +125,9 @@ protected FileRecordReader createRecordReader( valueType, file.level(), overrideSequenceWithSnapshotId, - file.minSequenceNumber()); + file.minSequenceNumber(), + metadataFallbackMapping, + !isChangelogFile(file)); if (deletionVector.isPresent() && !deletionVector.get().isEmpty()) { return new ExposeDeletionKeyValueReader(reader, deletionVector.get()); @@ -165,7 +169,8 @@ public ChainKeyValueFileReaderFactory build( dvFactory, chainReadContext, wrapped.options, - wrapped.readBatchSizer); + wrapped.readBatchSizer, + wrapped.createMetadataFallbackMapping()); } } } diff --git a/paimon-core/src/main/java/org/apache/paimon/io/KeyValueDataFileRecordReader.java b/paimon-core/src/main/java/org/apache/paimon/io/KeyValueDataFileRecordReader.java index 1e4db7046454..613feb0b38c3 100644 --- a/paimon-core/src/main/java/org/apache/paimon/io/KeyValueDataFileRecordReader.java +++ b/paimon-core/src/main/java/org/apache/paimon/io/KeyValueDataFileRecordReader.java @@ -20,6 +20,7 @@ import org.apache.paimon.KeyValue; import org.apache.paimon.KeyValueSerializer; +import org.apache.paimon.casting.FallbackMappingRow; import org.apache.paimon.data.InternalRow; import org.apache.paimon.reader.FileRecordIterator; import org.apache.paimon.reader.FileRecordReader; @@ -38,6 +39,7 @@ public class KeyValueDataFileRecordReader implements FileRecordReader private final int level; private final boolean overrideSequenceWithSnapshotId; private final long snapshotId; + @Nullable private final FallbackMappingRow metadataFallbackRow; public KeyValueDataFileRecordReader( FileRecordReader reader, @@ -45,12 +47,18 @@ public KeyValueDataFileRecordReader( RowType valueType, int level, boolean overrideSequenceWithSnapshotId, - long snapshotId) { + long snapshotId, + @Nullable int[] metadataFallbackMapping, + boolean applyMetadataFallback) { this.reader = reader; this.serializer = new KeyValueSerializer(keyType, valueType); this.level = level; this.overrideSequenceWithSnapshotId = overrideSequenceWithSnapshotId; this.snapshotId = snapshotId; + this.metadataFallbackRow = + applyMetadataFallback && metadataFallbackMapping != null + ? new FallbackMappingRow(metadataFallbackMapping) + : null; } @Nullable @@ -67,6 +75,9 @@ public FileRecordIterator readBatch() throws IOException { return null; } KeyValue kv = serializer.fromRow(internalRow).setLevel(level); + if (metadataFallbackRow != null) { + kv.replaceValue(metadataFallbackRow.replace(kv.value(), kv.value())); + } // In snapshot-ordering mode, an APPEND file's on-disk per-record sequence // numbers are stale; we override them with the commit snapshot id so later // snapshots win during merge. Any read path bypassing this override would diff --git a/paimon-core/src/main/java/org/apache/paimon/io/KeyValueFileReaderFactory.java b/paimon-core/src/main/java/org/apache/paimon/io/KeyValueFileReaderFactory.java index 24a28cb5b2d1..d45f95274191 100644 --- a/paimon-core/src/main/java/org/apache/paimon/io/KeyValueFileReaderFactory.java +++ b/paimon-core/src/main/java/org/apache/paimon/io/KeyValueFileReaderFactory.java @@ -49,6 +49,7 @@ import javax.annotation.Nullable; import java.io.IOException; +import java.util.ArrayList; import java.util.Collections; import java.util.List; import java.util.Map; @@ -76,6 +77,8 @@ public class KeyValueFileReaderFactory implements FileReaderFactory { private final BinaryRow partition; protected final DeletionVector.Factory dvFactory; @Nullable private final ReadBatchSizer readBatchSizer; + @Nullable protected final int[] metadataFallbackMapping; + private final String changelogFilePrefix; protected KeyValueFileReaderFactory( FileIO fileIO, @@ -88,7 +91,8 @@ protected KeyValueFileReaderFactory( BinaryRow partition, DeletionVector.Factory dvFactory, CoreOptions coreOptions, - @Nullable ReadBatchSizer readBatchSizer) { + @Nullable ReadBatchSizer readBatchSizer, + @Nullable int[] metadataFallbackMapping) { this.fileIO = fileIO; this.schemaManager = schemaManager; this.schema = schema; @@ -104,6 +108,8 @@ protected KeyValueFileReaderFactory( this.formatReaderMappings = new ConcurrentHashMap<>(); this.dvFactory = dvFactory; this.readBatchSizer = readBatchSizer; + this.metadataFallbackMapping = metadataFallbackMapping; + this.changelogFilePrefix = coreOptions.changelogFilePrefix(); } public TableSchema schema() { @@ -148,7 +154,13 @@ protected FileRecordReader createRecordReader( valueType, file.level(), overrideSequenceWithSnapshotId, - file.minSequenceNumber()); + file.minSequenceNumber(), + metadataFallbackMapping, + !isChangelogFile(file)); + } + + protected boolean isChangelogFile(DataFileMeta file) { + return file.fileName().startsWith(changelogFilePrefix); } private FileRecordReader createRecordReader( @@ -246,6 +258,7 @@ public static class Builder { protected RowType readKeyType; protected RowType readValueType; @Nullable protected ReadBatchSizer readBatchSizer; + @Nullable protected List changelogExtraValueFields; private Builder( FileIO fileIO, @@ -284,6 +297,7 @@ public Builder copyWithoutProjection() { extractor, options); copy.readBatchSizer = readBatchSizer; + copy.changelogExtraValueFields = changelogExtraValueFields; return copy; } @@ -318,6 +332,12 @@ public Builder withReadBatchSizer(ReadBatchSizer sizer) { return this; } + public Builder withChangelogExtraValueFields( + @Nullable List changelogExtraValueFields) { + this.changelogExtraValueFields = changelogExtraValueFields; + return this; + } + public RowType keyType() { return keyType; } @@ -348,6 +368,7 @@ public KeyValueFileReaderFactory build( boolean projectKeys, @Nullable List filters) { FormatReaderMapping.Builder builder = formatReaderMappingBuilder(projectKeys, filters); + int[] metadataFallbackMapping = createMetadataFallbackMapping(); return new KeyValueFileReaderFactory( fileIO, schemaManager, @@ -359,19 +380,66 @@ public KeyValueFileReaderFactory build( partition, dvFactory, options, - readBatchSizer); + readBatchSizer, + metadataFallbackMapping); + } + + @Nullable + protected int[] createMetadataFallbackMapping() { + if (changelogExtraValueFields == null + || changelogExtraValueFields.isEmpty() + || options.changelogEventMetadataFields().isEmpty()) { + return null; + } + + int[] mapping = new int[readValueType.getFieldCount()]; + java.util.Arrays.fill(mapping, -1); + List readFieldNames = readValueType.getFieldNames(); + List preserveColumns = options.changelogEventMetadataFields(); + for (int i = 0; i < changelogExtraValueFields.size(); i++) { + if (i >= preserveColumns.size()) { + break; + } + int metadataFieldId = changelogExtraValueFields.get(i).id(); + int metadataIndex = + readValueType.containsField(metadataFieldId) + ? readValueType.getFieldIndexByFieldId(metadataFieldId) + : -1; + int valueIndex = readFieldNames.indexOf(preserveColumns.get(i)); + if (metadataIndex >= 0 && valueIndex >= 0) { + mapping[metadataIndex] = valueIndex; + } + } + for (int index : mapping) { + if (index >= 0) { + return mapping; + } + } + return null; } protected FormatReaderMapping.Builder formatReaderMappingBuilder( boolean projectKeys, @Nullable List filters) { RowType finalReadKeyType = projectKeys ? this.readKeyType : keyType; + List readValueFields = new ArrayList<>(readValueType.getFields()); + if (changelogExtraValueFields != null) { + for (DataField extraField : changelogExtraValueFields) { + if (!readValueType.containsField(extraField.id())) { + readValueFields.add(extraField); + } + } + } List readTableFields = - KeyValue.createKeyValueFields( - finalReadKeyType.getFields(), readValueType.getFields()); + KeyValue.createKeyValueFields(finalReadKeyType.getFields(), readValueFields); + List extraFields = changelogExtraValueFields; Function> fieldsExtractor = schema -> { List dataKeyFields = extractor.keyFields(schema); - List dataValueFields = extractor.valueFields(schema); + List dataValueFields = + new ArrayList<>(extractor.valueFields(schema)); + if (extraFields != null) { + dataValueFields.addAll(extraFields); + } return KeyValue.createKeyValueFields(dataKeyFields, dataValueFields); }; return new FormatReaderMapping.Builder( diff --git a/paimon-core/src/main/java/org/apache/paimon/io/KeyValueFileWriterFactory.java b/paimon-core/src/main/java/org/apache/paimon/io/KeyValueFileWriterFactory.java index 7c55360fef4c..9eba0980bb81 100644 --- a/paimon-core/src/main/java/org/apache/paimon/io/KeyValueFileWriterFactory.java +++ b/paimon-core/src/main/java/org/apache/paimon/io/KeyValueFileWriterFactory.java @@ -67,6 +67,7 @@ public class KeyValueFileWriterFactory { private final RowType keyType; private final RowType valueType; private final FileWriterContextFactory formatContext; + @Nullable private final FileWriterContextFactory changelogFormatContext; private final long suggestedFileSize; private final CoreOptions options; private final FileIndexOptions fileIndexOptions; @@ -77,6 +78,7 @@ private KeyValueFileWriterFactory( FileIO fileIO, long schemaId, FileWriterContextFactory formatContext, + @Nullable FileWriterContextFactory changelogFormatContext, long suggestedFileSize, CoreOptions options) { this.fileIO = fileIO; @@ -84,6 +86,7 @@ private KeyValueFileWriterFactory( this.keyType = formatContext.keyType; this.valueType = formatContext.valueType; this.formatContext = formatContext; + this.changelogFormatContext = changelogFormatContext; this.suggestedFileSize = suggestedFileSize; this.options = options; this.fileIndexOptions = options.indexColumnsOptions(); @@ -156,16 +159,19 @@ public RollingFileWriter createRollingMergeTreeFileWrite public RollingFileWriter createRollingChangelogFileWriter(int level) { WriteFormatKey key = new WriteFormatKey(level, true); - FormatWriterFactory writerFactory = formatContext.createWriterFactory(key); + FileWriterContextFactory ctx = + changelogFormatContext != null ? changelogFormatContext : formatContext; + FormatWriterFactory writerFactory = ctx.createWriterFactory(key); return new RollingFileWriterImpl<>( () -> { - DataFilePathFactory pathFactory = formatContext.pathFactory(key); + DataFilePathFactory pathFactory = ctx.pathFactory(key); return createDataFileWriter( pathFactory.newChangelogPath(), key, FileSource.APPEND, pathFactory.isExternalPath(), - writerFactory); + writerFactory, + ctx); }, suggestedFileSize, Long.MAX_VALUE); @@ -212,18 +218,31 @@ private KeyValueDataFileWriter createDataFileWriter( FileSource fileSource, boolean isExternalPath, FormatWriterFactory writerFactory) { + return createDataFileWriter( + path, key, fileSource, isExternalPath, writerFactory, formatContext); + } + + private KeyValueDataFileWriter createDataFileWriter( + Path path, + WriteFormatKey key, + FileSource fileSource, + boolean isExternalPath, + FormatWriterFactory writerFactory, + FileWriterContextFactory ctx) { + RowType writerKeyType = ctx.keyType; + RowType writerValueType = ctx.valueType; // Changelog is sequentially consumed, file index is unnecessary. FileIndexOptions indexOptions = key.isChangelog ? new FileIndexOptions() : fileIndexOptions; Set dataFileManagedBlobFields = key.isChangelog ? Collections.emptySet() : managedBlobFields; - return formatContext.thinModeEnabled + return ctx.thinModeEnabled ? new KeyValueThinDataFileWriterImpl( fileIO, - formatContext.fileWriterContext(key, writerFactory), + ctx.fileWriterContext(key, writerFactory), path, - new KeyValueThinSerializer(keyType, valueType)::toRow, - keyType, - valueType, + new KeyValueThinSerializer(writerKeyType, writerValueType)::toRow, + writerKeyType, + writerValueType, schemaId, key.level, options, @@ -233,11 +252,11 @@ private KeyValueDataFileWriter createDataFileWriter( dataFileManagedBlobFields) : new KeyValueDataFileWriterImpl( fileIO, - formatContext.fileWriterContext(key, writerFactory), + ctx.fileWriterContext(key, writerFactory), path, - new KeyValueSerializer(keyType, valueType)::toRow, - keyType, - valueType, + new KeyValueSerializer(writerKeyType, writerValueType)::toRow, + writerKeyType, + writerValueType, schemaId, key.level, options, @@ -295,6 +314,7 @@ public static class Builder { private final FileFormat fileFormat; private final Function format2PathFactory; private final long suggestedFileSize; + @Nullable private RowType changelogValueType; private Builder( FileIO fileIO, @@ -313,6 +333,11 @@ private Builder( this.suggestedFileSize = suggestedFileSize; } + public Builder withChangelogValueType(@Nullable RowType changelogValueType) { + this.changelogValueType = changelogValueType; + return this; + } + public KeyValueFileWriterFactory build( BinaryRow partition, int bucket, CoreOptions options) { FileWriterContextFactory context = @@ -324,8 +349,20 @@ public KeyValueFileWriterFactory build( fileFormat, format2PathFactory, options); + FileWriterContextFactory changelogContext = null; + if (changelogValueType != null) { + changelogContext = + new FileWriterContextFactory( + partition, + bucket, + keyType, + changelogValueType, + fileFormat, + format2PathFactory, + options); + } return new KeyValueFileWriterFactory( - fileIO, schemaId, context, suggestedFileSize, options); + fileIO, schemaId, context, changelogContext, suggestedFileSize, options); } } diff --git a/paimon-core/src/main/java/org/apache/paimon/mergetree/compact/FirstRowMergeFunctionWrapper.java b/paimon-core/src/main/java/org/apache/paimon/mergetree/compact/FirstRowMergeFunctionWrapper.java index ce595063b425..366c1a03e34b 100644 --- a/paimon-core/src/main/java/org/apache/paimon/mergetree/compact/FirstRowMergeFunctionWrapper.java +++ b/paimon-core/src/main/java/org/apache/paimon/mergetree/compact/FirstRowMergeFunctionWrapper.java @@ -22,6 +22,8 @@ import org.apache.paimon.data.InternalRow; import org.apache.paimon.utils.Filter; +import javax.annotation.Nullable; + import static org.apache.paimon.utils.Preconditions.checkArgument; /** Wrapper for {@link MergeFunction}s to produce changelog by lookup for first row. */ @@ -31,8 +33,20 @@ public class FirstRowMergeFunctionWrapper implements MergeFunctionWrapper mergeFunctionFactory, Filter contains) { + this(mergeFunctionFactory, contains, null); + } + + public FirstRowMergeFunctionWrapper( + MergeFunctionFactory mergeFunctionFactory, + Filter contains, + @Nullable int[] preserveFieldIndices) { this.contains = contains; MergeFunction mergeFunction = mergeFunctionFactory.create(); checkArgument( @@ -40,6 +54,13 @@ public FirstRowMergeFunctionWrapper( "Merge function should be a FirstRowMergeFunction, but is %s, there is a bug.", mergeFunction.getClass().getName()); this.mergeFunction = (FirstRowMergeFunction) mergeFunction; + boolean hasMetadata = preserveFieldIndices != null && preserveFieldIndices.length > 0; + this.eventMetadata = + hasMetadata + ? new LookupChangelogMergeFunctionWrapper.EventMetadataAppendRow( + preserveFieldIndices) + : null; + this.reusedChangelog = hasMetadata ? new KeyValue() : null; } @Override @@ -67,6 +88,17 @@ public ChangelogResult getResult() { } // new record, output changelog - return reusedResult.setResult(result).addChangelog(result); + return reusedResult.setResult(result).addChangelog(withEventMetadata(result)); + } + + private KeyValue withEventMetadata(KeyValue result) { + LookupChangelogMergeFunctionWrapper.EventMetadataAppendRow metadata = eventMetadata; + KeyValue changelog = reusedChangelog; + if (metadata == null || changelog == null) { + return result; + } + metadata.replace(result.value(), result.value()); + return changelog.replace( + result.key(), result.sequenceNumber(), result.valueKind(), metadata); } } diff --git a/paimon-core/src/main/java/org/apache/paimon/mergetree/compact/LookupChangelogMergeFunctionWrapper.java b/paimon-core/src/main/java/org/apache/paimon/mergetree/compact/LookupChangelogMergeFunctionWrapper.java index 7283a3030d01..f0cb8f127253 100644 --- a/paimon-core/src/main/java/org/apache/paimon/mergetree/compact/LookupChangelogMergeFunctionWrapper.java +++ b/paimon-core/src/main/java/org/apache/paimon/mergetree/compact/LookupChangelogMergeFunctionWrapper.java @@ -20,7 +20,15 @@ import org.apache.paimon.KeyValue; import org.apache.paimon.codegen.RecordEqualiser; +import org.apache.paimon.data.BinaryString; +import org.apache.paimon.data.Blob; +import org.apache.paimon.data.Decimal; +import org.apache.paimon.data.InternalArray; +import org.apache.paimon.data.InternalMap; import org.apache.paimon.data.InternalRow; +import org.apache.paimon.data.InternalVector; +import org.apache.paimon.data.Timestamp; +import org.apache.paimon.data.variant.Variant; import org.apache.paimon.deletionvectors.BucketedDvMaintainer; import org.apache.paimon.lookup.LookupStrategy; import org.apache.paimon.mergetree.lookup.FilePosition; @@ -64,6 +72,8 @@ public class LookupChangelogMergeFunctionWrapper private final LookupStrategy lookupStrategy; private final @Nullable BucketedDvMaintainer deletionVectorsMaintainer; private final Comparator comparator; + @Nullable private final EventMetadataAppendRow reusedBeforeAppendRow; + @Nullable private final EventMetadataAppendRow reusedAfterAppendRow; public LookupChangelogMergeFunctionWrapper( MergeFunctionFactory mergeFunctionFactory, @@ -72,6 +82,24 @@ public LookupChangelogMergeFunctionWrapper( LookupStrategy lookupStrategy, @Nullable BucketedDvMaintainer deletionVectorsMaintainer, @Nullable UserDefinedSeqComparator userDefinedSeqComparator) { + this( + mergeFunctionFactory, + lookup, + valueEqualiser, + lookupStrategy, + deletionVectorsMaintainer, + userDefinedSeqComparator, + null); + } + + public LookupChangelogMergeFunctionWrapper( + MergeFunctionFactory mergeFunctionFactory, + Function lookup, + @Nullable RecordEqualiser valueEqualiser, + LookupStrategy lookupStrategy, + @Nullable BucketedDvMaintainer deletionVectorsMaintainer, + @Nullable UserDefinedSeqComparator userDefinedSeqComparator, + @Nullable int[] preserveFieldIndices) { MergeFunction mergeFunction = mergeFunctionFactory.create(); checkArgument( mergeFunction instanceof LookupMergeFunction, @@ -88,6 +116,11 @@ public LookupChangelogMergeFunctionWrapper( this.lookupStrategy = lookupStrategy; this.deletionVectorsMaintainer = deletionVectorsMaintainer; this.comparator = createSequenceComparator(userDefinedSeqComparator); + boolean hasMetadata = preserveFieldIndices != null && preserveFieldIndices.length > 0; + this.reusedBeforeAppendRow = + hasMetadata ? new EventMetadataAppendRow(preserveFieldIndices) : null; + this.reusedAfterAppendRow = + hasMetadata ? new EventMetadataAppendRow(preserveFieldIndices) : null; } @Override @@ -148,25 +181,39 @@ public ChangelogResult getResult() { private void setChangelog(@Nullable KeyValue before, KeyValue after) { if (before == null || !before.isAdd()) { if (after.isAdd()) { - reusedResult.addChangelog(replaceAfter(RowKind.INSERT, after)); + reusedResult.addChangelog(replaceAfterWithEventMetadata(RowKind.INSERT, after)); } } else { if (!after.isAdd()) { - reusedResult.addChangelog(replaceBefore(RowKind.DELETE, before)); + reusedResult.addChangelog( + replaceBeforeWithEventMetadata(RowKind.DELETE, before, after)); } else if (valueEqualiser == null || !valueEqualiser.equals(before.value(), after.value())) { reusedResult - .addChangelog(replaceBefore(RowKind.UPDATE_BEFORE, before)) - .addChangelog(replaceAfter(RowKind.UPDATE_AFTER, after)); + .addChangelog( + replaceBeforeWithEventMetadata( + RowKind.UPDATE_BEFORE, before, after)) + .addChangelog(replaceAfterWithEventMetadata(RowKind.UPDATE_AFTER, after)); } } } - private KeyValue replaceBefore(RowKind valueKind, KeyValue from) { - return replace(reusedBefore, valueKind, from); + private KeyValue replaceBeforeWithEventMetadata( + RowKind valueKind, KeyValue before, KeyValue after) { + if (reusedBeforeAppendRow != null) { + reusedBeforeAppendRow.replace(before.value(), after.value()); + return reusedBefore.replace( + before.key(), before.sequenceNumber(), valueKind, reusedBeforeAppendRow); + } + return replace(reusedBefore, valueKind, before); } - private KeyValue replaceAfter(RowKind valueKind, KeyValue from) { + private KeyValue replaceAfterWithEventMetadata(RowKind valueKind, KeyValue from) { + if (reusedAfterAppendRow != null) { + reusedAfterAppendRow.replace(from.value(), from.value()); + return reusedAfter.replace( + from.key(), from.sequenceNumber(), valueKind, reusedAfterAppendRow); + } return replace(reusedAfter, valueKind, from); } @@ -188,4 +235,141 @@ private Comparator createSequenceComparator( return Long.compare(o1.sequenceNumber(), o2.sequenceNumber()); }; } + + /** + * An {@link InternalRow} that presents the before-image row with event metadata columns + * appended. Positions {@code [0, N-1]} delegate to the primary (before-image) row unchanged. + * Positions {@code [N, N+K-1]} read the preserved field values from the event (after) row, + * remapping through {@code preserveFieldIndices}. + */ + static class EventMetadataAppendRow implements InternalRow { + + private final int[] preserveFieldIndices; + private InternalRow primaryRow; + private InternalRow eventRow; + + EventMetadataAppendRow(int[] preserveFieldIndices) { + this.preserveFieldIndices = preserveFieldIndices; + } + + EventMetadataAppendRow replace(InternalRow primaryRow, InternalRow eventRow) { + this.primaryRow = primaryRow; + this.eventRow = eventRow; + return this; + } + + private InternalRow rowFor(int pos) { + return pos < primaryRow.getFieldCount() ? primaryRow : eventRow; + } + + private int posFor(int pos) { + int primaryCount = primaryRow.getFieldCount(); + return pos < primaryCount ? pos : preserveFieldIndices[pos - primaryCount]; + } + + @Override + public int getFieldCount() { + return primaryRow.getFieldCount() + preserveFieldIndices.length; + } + + @Override + public RowKind getRowKind() { + return primaryRow.getRowKind(); + } + + @Override + public void setRowKind(RowKind kind) { + primaryRow.setRowKind(kind); + } + + @Override + public boolean isNullAt(int pos) { + return rowFor(pos).isNullAt(posFor(pos)); + } + + @Override + public boolean getBoolean(int pos) { + return rowFor(pos).getBoolean(posFor(pos)); + } + + @Override + public byte getByte(int pos) { + return rowFor(pos).getByte(posFor(pos)); + } + + @Override + public short getShort(int pos) { + return rowFor(pos).getShort(posFor(pos)); + } + + @Override + public int getInt(int pos) { + return rowFor(pos).getInt(posFor(pos)); + } + + @Override + public long getLong(int pos) { + return rowFor(pos).getLong(posFor(pos)); + } + + @Override + public float getFloat(int pos) { + return rowFor(pos).getFloat(posFor(pos)); + } + + @Override + public double getDouble(int pos) { + return rowFor(pos).getDouble(posFor(pos)); + } + + @Override + public BinaryString getString(int pos) { + return rowFor(pos).getString(posFor(pos)); + } + + @Override + public Decimal getDecimal(int pos, int precision, int scale) { + return rowFor(pos).getDecimal(posFor(pos), precision, scale); + } + + @Override + public Timestamp getTimestamp(int pos, int precision) { + return rowFor(pos).getTimestamp(posFor(pos), precision); + } + + @Override + public byte[] getBinary(int pos) { + return rowFor(pos).getBinary(posFor(pos)); + } + + @Override + public Variant getVariant(int pos) { + return rowFor(pos).getVariant(posFor(pos)); + } + + @Override + public Blob getBlob(int pos) { + return rowFor(pos).getBlob(posFor(pos)); + } + + @Override + public InternalArray getArray(int pos) { + return rowFor(pos).getArray(posFor(pos)); + } + + @Override + public InternalVector getVector(int pos) { + return rowFor(pos).getVector(posFor(pos)); + } + + @Override + public InternalMap getMap(int pos) { + return rowFor(pos).getMap(posFor(pos)); + } + + @Override + public InternalRow getRow(int pos, int numFields) { + return rowFor(pos).getRow(posFor(pos), numFields); + } + } } diff --git a/paimon-core/src/main/java/org/apache/paimon/mergetree/compact/LookupMergeTreeCompactRewriter.java b/paimon-core/src/main/java/org/apache/paimon/mergetree/compact/LookupMergeTreeCompactRewriter.java index 36cc4cd6a4d5..949e7e532e9d 100644 --- a/paimon-core/src/main/java/org/apache/paimon/mergetree/compact/LookupMergeTreeCompactRewriter.java +++ b/paimon-core/src/main/java/org/apache/paimon/mergetree/compact/LookupMergeTreeCompactRewriter.java @@ -199,14 +199,17 @@ public static class LookupMergeFunctionWrapperFactory @Nullable private final RecordEqualiser valueEqualiser; private final LookupStrategy lookupStrategy; @Nullable private final UserDefinedSeqComparator userDefinedSeqComparator; + @Nullable private final int[] preserveFieldIndices; public LookupMergeFunctionWrapperFactory( @Nullable RecordEqualiser valueEqualiser, LookupStrategy lookupStrategy, - @Nullable UserDefinedSeqComparator userDefinedSeqComparator) { + @Nullable UserDefinedSeqComparator userDefinedSeqComparator, + @Nullable int[] preserveFieldIndices) { this.valueEqualiser = valueEqualiser; this.lookupStrategy = lookupStrategy; this.userDefinedSeqComparator = userDefinedSeqComparator; + this.preserveFieldIndices = preserveFieldIndices; } @Override @@ -227,7 +230,8 @@ public MergeFunctionWrapper create( valueEqualiser, lookupStrategy, deletionVectorsMaintainer, - userDefinedSeqComparator); + userDefinedSeqComparator, + preserveFieldIndices); } } @@ -235,6 +239,16 @@ public MergeFunctionWrapper create( public static class FirstRowMergeFunctionWrapperFactory implements MergeFunctionWrapperFactory { + @Nullable private final int[] preserveFieldIndices; + + public FirstRowMergeFunctionWrapperFactory() { + this(null); + } + + public FirstRowMergeFunctionWrapperFactory(@Nullable int[] preserveFieldIndices) { + this.preserveFieldIndices = preserveFieldIndices; + } + @Override public MergeFunctionWrapper create( MergeFunctionFactory mfFactory, @@ -249,7 +263,8 @@ public MergeFunctionWrapper create( } catch (IOException e) { throw new UncheckedIOException(e); } - }); + }, + preserveFieldIndices); } } } diff --git a/paimon-core/src/main/java/org/apache/paimon/mergetree/compact/MergeTreeCompactManagerFactory.java b/paimon-core/src/main/java/org/apache/paimon/mergetree/compact/MergeTreeCompactManagerFactory.java index 87ed4868accb..f50471b935d5 100644 --- a/paimon-core/src/main/java/org/apache/paimon/mergetree/compact/MergeTreeCompactManagerFactory.java +++ b/paimon-core/src/main/java/org/apache/paimon/mergetree/compact/MergeTreeCompactManagerFactory.java @@ -58,6 +58,7 @@ import org.apache.paimon.options.Options; import org.apache.paimon.schema.SchemaManager; import org.apache.paimon.schema.TableSchema; +import org.apache.paimon.table.system.ChangelogEventMetadata; import org.apache.paimon.types.RowType; import org.apache.paimon.utils.FieldsComparator; import org.apache.paimon.utils.UserDefinedSeqComparator; @@ -134,6 +135,14 @@ public MergeTreeCompactManagerFactory( this.schema = schema; this.recordLevelExpire = recordLevelExpire; this.cacheManager = cacheManager; + + ChangelogEventMetadata.validate(valueType, options); + if (options.changelogProducer() == ChangelogProducer.LOOKUP + && !options.changelogEventMetadataFields().isEmpty()) { + writerFactoryBuilder.withChangelogValueType( + ChangelogEventMetadata.appendStorageMetadataFields( + valueType, valueType, options)); + } } @Override @@ -332,6 +341,10 @@ private MergeTreeCompactRewriter createRewriter( PersistProcessor.Factory processorFactory; LookupMergeTreeCompactRewriter.MergeFunctionWrapperFactory wrapperFactory; FileReaderFactory lookupReaderFactory = readerFactory; + int[] preserveFieldIndices = + options.changelogProducer() == ChangelogProducer.LOOKUP + ? ChangelogEventMetadata.preserveFieldIndices(valueType, options) + : null; if (lookupStrategy.isFirstRow) { if (options.deletionVectorsEnabled()) { throw new UnsupportedOperationException( @@ -343,7 +356,7 @@ private MergeTreeCompactRewriter createRewriter( .withReadValueType(RowType.of()) .build(partition, bucket, dvFactory); processorFactory = PersistEmptyProcessor.factory(); - wrapperFactory = new FirstRowMergeFunctionWrapperFactory(); + wrapperFactory = new FirstRowMergeFunctionWrapperFactory(preserveFieldIndices); } else { if (lookupStrategy.deletionVector) { if (lookupStrategy.produceChangelog @@ -368,7 +381,8 @@ private MergeTreeCompactRewriter createRewriter( new LookupMergeFunctionWrapperFactory<>( logDedupEqualSupplier.get(), lookupStrategy, - UserDefinedSeqComparator.create(valueType, options)); + UserDefinedSeqComparator.create(valueType, options), + preserveFieldIndices); } LookupLevels lookupLevels = createLookupLevels( diff --git a/paimon-core/src/main/java/org/apache/paimon/operation/MergeFileSplitRead.java b/paimon-core/src/main/java/org/apache/paimon/operation/MergeFileSplitRead.java index 9e66aad52fca..ceff944f4c86 100644 --- a/paimon-core/src/main/java/org/apache/paimon/operation/MergeFileSplitRead.java +++ b/paimon-core/src/main/java/org/apache/paimon/operation/MergeFileSplitRead.java @@ -52,6 +52,7 @@ import org.apache.paimon.table.source.DataSplit; import org.apache.paimon.table.source.DeletionFile; import org.apache.paimon.table.source.Split; +import org.apache.paimon.table.system.ChangelogEventMetadata; import org.apache.paimon.types.DataField; import org.apache.paimon.types.RowType; import org.apache.paimon.utils.ProjectedRow; @@ -119,6 +120,7 @@ public MergeFileSplitRead( CoreOptions.fromMap(tableSchema.options()), keyType, valueType, null); this.sequenceFields = options.sequenceField(); this.sequenceOrder = options.sequenceFieldSortOrderIsAscending(); + ChangelogEventMetadata.validate(tableSchema.logicalRowType(), options); } public Comparator keyComparator() { @@ -172,6 +174,31 @@ public RowType adjustReadType(RowType readType) { adjustedReadType = new RowType(allFields); } } + + // Metadata columns are backed by values from the incoming event. Ordinary data files do + // not physically contain the synthetic metadata fields, so retain the corresponding + // physical fields while reading whenever a metadata field was requested. The outer read + // projection removes these internal dependencies after the reader has populated the + // metadata columns. + List preserveColumns = options.changelogEventMetadataFields(); + if (!preserveColumns.isEmpty()) { + List readFieldNames = adjustedReadType.getFieldNames(); + List extraFields = new ArrayList<>(); + RowType logicalRowType = tableSchema.logicalRowType(); + for (String preserveColumn : preserveColumns) { + String metadataName = + ChangelogEventMetadata.metadataFieldName(preserveColumn, options); + if (readFieldNames.contains(metadataName) + && !readFieldNames.contains(preserveColumn)) { + extraFields.add(logicalRowType.getField(preserveColumn)); + } + } + if (!extraFields.isEmpty()) { + List allFields = new ArrayList<>(adjustedReadType.getFields()); + allFields.addAll(extraFields); + adjustedReadType = new RowType(allFields); + } + } adjustedReadType = mfFactory.adjustReadType(adjustedReadType); return adjustedReadType; } diff --git a/paimon-core/src/main/java/org/apache/paimon/operation/RawFileSplitRead.java b/paimon-core/src/main/java/org/apache/paimon/operation/RawFileSplitRead.java index 84b32eaa1d97..5d211afd3d4a 100644 --- a/paimon-core/src/main/java/org/apache/paimon/operation/RawFileSplitRead.java +++ b/paimon-core/src/main/java/org/apache/paimon/operation/RawFileSplitRead.java @@ -19,6 +19,7 @@ package org.apache.paimon.operation; import org.apache.paimon.CoreOptions; +import org.apache.paimon.casting.FallbackMappingRow; import org.apache.paimon.data.BinaryRow; import org.apache.paimon.data.InternalRow; import org.apache.paimon.deletionvectors.ApplyDeletionVectorReader; @@ -40,6 +41,7 @@ import org.apache.paimon.predicate.Predicate; import org.apache.paimon.predicate.TopN; import org.apache.paimon.reader.EmptyFileRecordReader; +import org.apache.paimon.reader.FileRecordIterator; import org.apache.paimon.reader.FileRecordReader; import org.apache.paimon.reader.LimitRecordReader; import org.apache.paimon.reader.ReadBatchSizer; @@ -51,11 +53,14 @@ import org.apache.paimon.table.source.DeletionFile; import org.apache.paimon.table.source.IncrementalSplit; import org.apache.paimon.table.source.Split; +import org.apache.paimon.table.system.ChangelogEventMetadata; +import org.apache.paimon.types.DataField; import org.apache.paimon.types.RowType; import org.apache.paimon.utils.FileStorePathFactory; import org.apache.paimon.utils.FormatReaderMapping; import org.apache.paimon.utils.FormatReaderMapping.Builder; import org.apache.paimon.utils.IOExceptionSupplier; +import org.apache.paimon.utils.ProjectedRow; import org.apache.paimon.utils.RoaringBitmap32; import org.slf4j.Logger; @@ -65,6 +70,7 @@ import java.io.IOException; import java.util.ArrayList; +import java.util.Arrays; import java.util.HashMap; import java.util.List; import java.util.Map; @@ -88,8 +94,13 @@ public class RawFileSplitRead implements SplitRead { private final boolean nestedFieldEnabled; private final boolean ignoreCorruptFiles; private final boolean ignoreLostFiles; + private final List changelogExtraValueFields; + private final List metadataPreserveColumns; + private final String metadataFieldPrefix; private RowType readRowType; + @Nullable private RowType outerReadRowType; + @Nullable private int[] metadataFallbackMapping; @Nullable private List filters; @Nullable private TopN topN; @Nullable private Long limit; @@ -114,7 +125,12 @@ public RawFileSplitRead( this.ignoreLostFiles = coreOptions.scanIgnoreLostFile(); this.rowTrackingEnabled = coreOptions.rowTrackingEnabled(); this.nestedFieldEnabled = coreOptions.dataEvolutionNestedFieldEnabled(); - this.readRowType = rowType; + this.metadataPreserveColumns = coreOptions.changelogEventMetadataFields(); + this.metadataFieldPrefix = coreOptions.changelogMetadataFieldPrefix(); + this.changelogExtraValueFields = createChangelogExtraValueFields(schema, coreOptions); + this.readRowType = readTypeWithMetadataDependencies(rowType); + this.outerReadRowType = this.readRowType.equals(rowType) ? null : rowType; + this.metadataFallbackMapping = createMetadataFallbackMapping(this.readRowType); } @Override @@ -129,10 +145,13 @@ public SplitRead withIOManager(@Nullable IOManager ioManager) { @Override public SplitRead withReadType(RowType readRowType) { - if (!this.readRowType.equals(readRowType)) { + RowType adjustedReadType = readTypeWithMetadataDependencies(readRowType); + if (!this.readRowType.equals(adjustedReadType)) { formatReaderMappings.clear(); } - this.readRowType = readRowType; + this.readRowType = adjustedReadType; + this.outerReadRowType = adjustedReadType.equals(readRowType) ? null : readRowType; + this.metadataFallbackMapping = createMetadataFallbackMapping(adjustedReadType); return this; } @@ -264,12 +283,13 @@ private Builder createFormatReaderMappingBuilder( formatDiscover, outputRowType.getFields(), schema -> { + List fields = new ArrayList<>(schema.fields()); + fields.addAll(changelogExtraValueFields); if (rowTrackingEnabled) { // maybe file has no row id and sequence number, but in manifest entry - return rowTypeWithRowTracking(schema.logicalRowType(), true, true) - .getFields(); + return rowTypeWithRowTracking(new RowType(fields), true, true).getFields(); } - return schema.fields(); + return fields; }, filters, pushDownTopN, @@ -387,8 +407,98 @@ private FileRecordReader createFileReader( } if (deletionVector != null && !deletionVector.isEmpty()) { - return new ApplyDeletionVectorReader(fileRecordReader, deletionVector); + fileRecordReader = new ApplyDeletionVectorReader(fileRecordReader, deletionVector); } - return fileRecordReader; + return applyMetadataFallbackAndOuterProjection(fileRecordReader); + } + + private RowType readTypeWithMetadataDependencies(RowType requestedReadType) { + if (changelogExtraValueFields.isEmpty()) { + return requestedReadType; + } + + List readFieldNames = requestedReadType.getFieldNames(); + List dependencies = new ArrayList<>(); + for (String preserveColumn : metadataPreserveColumns) { + String metadataName = metadataFieldPrefix + preserveColumn; + if (readFieldNames.contains(metadataName) + && !readFieldNames.contains(preserveColumn) + && schema.logicalRowType().containsField(preserveColumn)) { + dependencies.add(schema.logicalRowType().getField(preserveColumn)); + } + } + + if (dependencies.isEmpty()) { + return requestedReadType; + } + List fields = new ArrayList<>(requestedReadType.getFields()); + fields.addAll(dependencies); + return new RowType(fields); + } + + @Nullable + private int[] createMetadataFallbackMapping(RowType rowType) { + if (metadataPreserveColumns.isEmpty()) { + return null; + } + + int[] mapping = new int[rowType.getFieldCount()]; + Arrays.fill(mapping, -1); + boolean hasMapping = false; + List fieldNames = rowType.getFieldNames(); + for (int i = 0; i < metadataPreserveColumns.size(); i++) { + String preserveColumn = metadataPreserveColumns.get(i); + int metadataFieldId = changelogExtraValueFields.get(i).id(); + int metadataIndex = + rowType.containsField(metadataFieldId) + ? rowType.getFieldIndexByFieldId(metadataFieldId) + : -1; + int physicalIndex = fieldNames.indexOf(preserveColumn); + if (metadataIndex >= 0 && physicalIndex >= 0) { + mapping[metadataIndex] = physicalIndex; + hasMapping = true; + } + } + return hasMapping ? mapping : null; + } + + private static List createChangelogExtraValueFields( + TableSchema schema, CoreOptions options) { + return ChangelogEventMetadata.storageValueFields(schema.logicalRowType(), options); + } + + private FileRecordReader applyMetadataFallbackAndOuterProjection( + FileRecordReader reader) { + if (metadataFallbackMapping == null && outerReadRowType == null) { + return reader; + } + + final FallbackMappingRow fallbackRow = + metadataFallbackMapping == null + ? null + : new FallbackMappingRow(metadataFallbackMapping); + final ProjectedRow projectedRow = + outerReadRowType == null ? null : ProjectedRow.from(outerReadRowType, readRowType); + return new FileRecordReader() { + @Nullable + @Override + public FileRecordIterator readBatch() throws IOException { + FileRecordIterator iterator = reader.readBatch(); + if (iterator == null) { + return null; + } + return iterator.transform( + row -> { + InternalRow result = + fallbackRow == null ? row : fallbackRow.replace(row, row); + return projectedRow == null ? result : projectedRow.replaceRow(result); + }); + } + + @Override + public void close() throws IOException { + reader.close(); + } + }; } } diff --git a/paimon-core/src/main/java/org/apache/paimon/schema/SchemaManagerUtils.java b/paimon-core/src/main/java/org/apache/paimon/schema/SchemaManagerUtils.java index a71451c865d6..5dcaafe5fe4c 100644 --- a/paimon-core/src/main/java/org/apache/paimon/schema/SchemaManagerUtils.java +++ b/paimon-core/src/main/java/org/apache/paimon/schema/SchemaManagerUtils.java @@ -72,6 +72,7 @@ import static org.apache.paimon.CoreOptions.AGG_FUNCTION; import static org.apache.paimon.CoreOptions.BUCKET_KEY; +import static org.apache.paimon.CoreOptions.CHANGELOG_PRODUCER_EVENT_METADATA_FIELDS; import static org.apache.paimon.CoreOptions.CLUSTERING_COLUMNS; import static org.apache.paimon.CoreOptions.DELETION_VECTORS_ENABLED; import static org.apache.paimon.CoreOptions.DELETION_VECTORS_MODIFIABLE; @@ -636,6 +637,20 @@ static Map applyRenameColumnsToOptions( newOptions.put(SEQUENCE_FIELD.key(), String.join(",", newSequenceFields)); } + // changelog metadata source fields rename + String metadataFieldsStr = options.get(CHANGELOG_PRODUCER_EVENT_METADATA_FIELDS.key()); + if (!StringUtils.isNullOrWhitespaceOnly(metadataFieldsStr)) { + List metadataFields = + Arrays.stream(metadataFieldsStr.split(",")) + .map(String::trim) + .collect(Collectors.toList()); + List newMetadataFields = + applyNotNestedColumnRename(metadataFields, renameMappings); + newOptions.put( + CHANGELOG_PRODUCER_EVENT_METADATA_FIELDS.key(), + String.join(",", newMetadataFields)); + } + // case 2: the option key is composed of certain fixed prefixes, suffixes, and the field // name, while the option value doesn't contain field names. List> fieldNameToOptionKeys = diff --git a/paimon-core/src/main/java/org/apache/paimon/schema/SchemaValidation.java b/paimon-core/src/main/java/org/apache/paimon/schema/SchemaValidation.java index ed316a72238d..a41928f6e163 100644 --- a/paimon-core/src/main/java/org/apache/paimon/schema/SchemaValidation.java +++ b/paimon-core/src/main/java/org/apache/paimon/schema/SchemaValidation.java @@ -41,6 +41,7 @@ import org.apache.paimon.options.ConfigOption; import org.apache.paimon.options.Options; import org.apache.paimon.table.BucketMode; +import org.apache.paimon.table.system.ChangelogEventMetadata; import org.apache.paimon.types.ArrayType; import org.apache.paimon.types.BigIntType; import org.apache.paimon.types.DataField; @@ -199,6 +200,7 @@ public static void validateTableSchema(TableSchema schema, Set dynamicOp "Can not set %s on table without primary keys, please define primary keys.", CHANGELOG_PRODUCER.key())); } + ChangelogEventMetadata.validate(new RowType(schema.fields()), options); if (options.streamingReadOverwrite() && (changelogProducer == ChangelogProducer.FULL_COMPACTION || changelogProducer == ChangelogProducer.LOOKUP)) { diff --git a/paimon-core/src/main/java/org/apache/paimon/table/system/AuditLogTable.java b/paimon-core/src/main/java/org/apache/paimon/table/system/AuditLogTable.java index 5ccc84886ee7..baa49fa80e51 100644 --- a/paimon-core/src/main/java/org/apache/paimon/table/system/AuditLogTable.java +++ b/paimon-core/src/main/java/org/apache/paimon/table/system/AuditLogTable.java @@ -169,7 +169,8 @@ public String name() { public RowType rowType() { List fields = new ArrayList<>(specialFields); fields.addAll(wrapped.rowType().getFields()); - return new RowType(fields); + RowType baseType = new RowType(fields); + return ChangelogEventMetadataTable.computeExtendedRowType(wrapped, baseType); } @Override diff --git a/paimon-core/src/main/java/org/apache/paimon/table/system/ChangelogEventMetadata.java b/paimon-core/src/main/java/org/apache/paimon/table/system/ChangelogEventMetadata.java new file mode 100644 index 000000000000..4bf5ed60962d --- /dev/null +++ b/paimon-core/src/main/java/org/apache/paimon/table/system/ChangelogEventMetadata.java @@ -0,0 +1,265 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.paimon.table.system; + +import org.apache.paimon.CoreOptions; +import org.apache.paimon.CoreOptions.ChangelogProducer; +import org.apache.paimon.table.SpecialFields; +import org.apache.paimon.types.DataField; +import org.apache.paimon.types.RowType; + +import javax.annotation.Nullable; + +import java.util.ArrayList; +import java.util.Collections; +import java.util.HashMap; +import java.util.HashSet; +import java.util.List; +import java.util.Map; +import java.util.Set; +import java.util.function.Function; + +/** Shared validation and row-type construction for changelog event metadata fields. */ +public final class ChangelogEventMetadata { + + private ChangelogEventMetadata() {} + + /** + * Validates the event metadata configuration against the physical value type. + * + *

The fields are appended to records produced by the lookup changelog wrapper. They must not + * be enabled for another changelog producer because those producers share the ordinary writer + * and do not emit the appended values. + */ + public static void validate(RowType valueType, CoreOptions options) { + List preserveColumns = options.changelogEventMetadataFields(); + if (preserveColumns.isEmpty()) { + return; + } + + if (options.changelogProducer() != ChangelogProducer.LOOKUP) { + throw new IllegalArgumentException( + String.format( + "Option '%s' can only be used when '%s' is '%s', but it is '%s'.", + CoreOptions.CHANGELOG_PRODUCER_EVENT_METADATA_FIELDS.key(), + CoreOptions.CHANGELOG_PRODUCER.key(), + ChangelogProducer.LOOKUP, + options.changelogProducer())); + } + + Set valueFieldNames = new HashSet<>(valueType.getFieldNames()); + Set preservedColumns = new HashSet<>(); + Set metadataFieldNames = new HashSet<>(); + Map metadataFieldSources = new HashMap<>(); + Map storageFieldSources = new HashMap<>(); + for (String preserveColumn : preserveColumns) { + if (!preservedColumns.add(preserveColumn)) { + throw new IllegalArgumentException( + String.format( + "Column '%s' is specified more than once in '%s'.", + preserveColumn, + CoreOptions.CHANGELOG_PRODUCER_EVENT_METADATA_FIELDS.key())); + } + + if (!valueFieldNames.contains(preserveColumn)) { + throw new IllegalArgumentException( + String.format( + "Column '%s' specified in '%s' not found in value type. Available columns: %s", + preserveColumn, + CoreOptions.CHANGELOG_PRODUCER_EVENT_METADATA_FIELDS.key(), + valueType.getFieldNames())); + } + + DataField physicalField = valueType.getField(preserveColumn); + String metadataFieldName = metadataFieldName(preserveColumn, options); + if (valueFieldNames.contains(metadataFieldName)) { + throw new IllegalArgumentException( + String.format( + "Metadata field '%s' created by '%s' conflicts with an existing value column.", + metadataFieldName, + CoreOptions.CHANGELOG_PRODUCER_EVENT_METADATA_FIELDS.key())); + } + if (SpecialFields.isSystemField(metadataFieldName)) { + throw new IllegalArgumentException( + String.format( + "Metadata field '%s' created by '%s' conflicts with a system field.", + metadataFieldName, + CoreOptions.CHANGELOG_PRODUCER_EVENT_METADATA_FIELDS.key())); + } + if (!metadataFieldNames.add(metadataFieldName)) { + throw new IllegalArgumentException( + String.format( + "Metadata field '%s' is created more than once by '%s'.", + metadataFieldName, + CoreOptions.CHANGELOG_PRODUCER_EVENT_METADATA_FIELDS.key())); + } + + String storageFieldName = storageMetadataFieldName(physicalField, options); + if (valueFieldNames.contains(storageFieldName)) { + throw new IllegalArgumentException( + String.format( + "Storage metadata field '%s' created by '%s' " + + "conflicts with an existing value column.", + storageFieldName, + CoreOptions.CHANGELOG_PRODUCER_EVENT_METADATA_FIELDS.key())); + } + if (SpecialFields.isSystemField(storageFieldName)) { + throw new IllegalArgumentException( + String.format( + "Storage metadata field '%s' created by '%s' " + + "conflicts with a system field.", + storageFieldName, + CoreOptions.CHANGELOG_PRODUCER_EVENT_METADATA_FIELDS.key())); + } + if (storageFieldSources.put(storageFieldName, preserveColumn) != null) { + throw new IllegalArgumentException( + String.format( + "Storage metadata field '%s' is created more than once by '%s'.", + storageFieldName, + CoreOptions.CHANGELOG_PRODUCER_EVENT_METADATA_FIELDS.key())); + } + metadataFieldSources.put(metadataFieldName, preserveColumn); + } + + for (Map.Entry metadataField : metadataFieldSources.entrySet()) { + String storageSource = storageFieldSources.get(metadataField.getKey()); + if (storageSource != null && !storageSource.equals(metadataField.getValue())) { + throw new IllegalArgumentException( + String.format( + "Metadata field '%s' created by '%s' conflicts with the storage " + + "metadata field for column '%s'.", + metadataField.getKey(), + CoreOptions.CHANGELOG_PRODUCER_EVENT_METADATA_FIELDS.key(), + storageSource)); + } + } + } + + /** Returns the nullable public metadata fields appended to a table row. */ + public static List extraValueFields(RowType valueType, CoreOptions options) { + validate(valueType, options); + return metadataValueFields( + valueType, + options, + physicalField -> metadataFieldName(physicalField.name(), options)); + } + + /** + * Returns the nullable fields used to store event metadata in changelog files. + * + *

Storage names use the configured prefix and source field ID, so they remain stable when + * the source column is renamed. + */ + public static List storageValueFields(RowType valueType, CoreOptions options) { + validate(valueType, options); + return metadataValueFields( + valueType, + options, + physicalField -> storageMetadataFieldName(physicalField, options)); + } + + /** Returns the physical value-field positions copied into event metadata columns. */ + @Nullable + public static int[] preserveFieldIndices(RowType valueType, CoreOptions options) { + List preserveColumns = options.changelogEventMetadataFields(); + if (preserveColumns.isEmpty()) { + return null; + } + validate(valueType, options); + int[] indices = new int[preserveColumns.size()]; + for (int i = 0; i < preserveColumns.size(); i++) { + indices[i] = valueType.getFieldIndex(preserveColumns.get(i)); + } + return indices; + } + + /** + * Appends event metadata fields to a row type while resolving their source fields from the + * physical value type. + * + *

The two row types intentionally differ for system tables such as {@code audit_log}, whose + * base row prepends system fields to the physical value fields. + */ + public static RowType appendMetadataFields( + RowType baseRowType, RowType valueType, CoreOptions options) { + return appendFields(baseRowType, extraValueFields(valueType, options)); + } + + /** Appends the internal storage metadata fields to a row type. */ + public static RowType appendStorageMetadataFields( + RowType baseRowType, RowType valueType, CoreOptions options) { + return appendFields(baseRowType, storageValueFields(valueType, options)); + } + + private static RowType appendFields(RowType baseRowType, List extraFields) { + if (extraFields.isEmpty()) { + return baseRowType; + } + + boolean hasExistingMetadata = false; + for (DataField extraField : extraFields) { + if (baseRowType.containsField(extraField.name())) { + hasExistingMetadata = true; + } + } + if (hasExistingMetadata) { + boolean allExisting = + extraFields.stream().allMatch(field -> baseRowType.containsField(field.name())); + if (allExisting) { + return baseRowType; + } + throw new IllegalArgumentException( + "Changelog event metadata fields conflict with the requested row type."); + } + + List fields = new ArrayList<>(baseRowType.getFields()); + fields.addAll(extraFields); + return new RowType(fields); + } + + /** Returns the public metadata field name for a preserved physical field. */ + public static String metadataFieldName(String preserveColumn, CoreOptions options) { + return options.changelogMetadataFieldPrefix() + preserveColumn; + } + + /** Returns the internal changelog storage name for a preserved physical field. */ + private static String storageMetadataFieldName(DataField physicalField, CoreOptions options) { + return options.changelogMetadataFieldPrefix() + "field_id_" + physicalField.id(); + } + + private static List metadataValueFields( + RowType valueType, CoreOptions options, Function metadataName) { + List preserveColumns = options.changelogEventMetadataFields(); + if (preserveColumns.isEmpty()) { + return Collections.emptyList(); + } + + int nextId = RowType.currentHighestFieldId(valueType.getFields()) + 1; + List extraFields = new ArrayList<>(preserveColumns.size()); + for (String preserveColumn : preserveColumns) { + DataField physicalField = valueType.getField(preserveColumn); + extraFields.add( + new DataField( + nextId++, + metadataName.apply(physicalField), + physicalField.type().copy(true))); + } + return extraFields; + } +} diff --git a/paimon-core/src/main/java/org/apache/paimon/table/system/ChangelogEventMetadataTable.java b/paimon-core/src/main/java/org/apache/paimon/table/system/ChangelogEventMetadataTable.java new file mode 100644 index 000000000000..77468239fbfe --- /dev/null +++ b/paimon-core/src/main/java/org/apache/paimon/table/system/ChangelogEventMetadataTable.java @@ -0,0 +1,193 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.paimon.table.system; + +import org.apache.paimon.CoreOptions; +import org.apache.paimon.Snapshot; +import org.apache.paimon.consumer.ConsumerManager; +import org.apache.paimon.fs.FileIO; +import org.apache.paimon.fs.Path; +import org.apache.paimon.manifest.IndexManifestEntry; +import org.apache.paimon.manifest.ManifestEntry; +import org.apache.paimon.manifest.ManifestFileMeta; +import org.apache.paimon.schema.SchemaManager; +import org.apache.paimon.table.DataTable; +import org.apache.paimon.table.FileStoreTable; +import org.apache.paimon.table.ReadonlyTable; +import org.apache.paimon.table.Table; +import org.apache.paimon.table.source.DataTableScan; +import org.apache.paimon.table.source.InnerTableRead; +import org.apache.paimon.table.source.StreamDataTableScan; +import org.apache.paimon.table.source.snapshot.SnapshotReader; +import org.apache.paimon.types.RowType; +import org.apache.paimon.utils.BranchManager; +import org.apache.paimon.utils.ChangelogManager; +import org.apache.paimon.utils.SimpleFileReader; +import org.apache.paimon.utils.SnapshotManager; +import org.apache.paimon.utils.TagManager; + +import java.util.List; +import java.util.Map; +import java.util.Optional; + +/** + * A {@link Table} wrapper that extends the row type with configured-prefix metadata columns for + * changelog retraction records. These columns carry the incoming event's preserved field values + * alongside the correct before-image. + */ +public class ChangelogEventMetadataTable implements DataTable, ReadonlyTable { + + private final FileStoreTable wrapped; + private final RowType extendedRowType; + + public ChangelogEventMetadataTable(FileStoreTable wrapped) { + this.wrapped = wrapped; + this.extendedRowType = computeExtendedRowType(wrapped, wrapped.rowType()); + } + + public static RowType computeExtendedRowType(FileStoreTable table, RowType baseRowType) { + return ChangelogEventMetadata.appendMetadataFields( + baseRowType, table.schema().logicalRowType(), CoreOptions.fromMap(table.options())); + } + + @Override + public RowType rowType() { + return extendedRowType; + } + + @Override + public InnerTableRead newRead() { + return wrapped.newRead().withReadType(extendedRowType); + } + + @Override + public String name() { + return wrapped.name(); + } + + @Override + public List partitionKeys() { + return wrapped.partitionKeys(); + } + + @Override + public Map options() { + return wrapped.options(); + } + + @Override + public List primaryKeys() { + return wrapped.primaryKeys(); + } + + @Override + public Optional latestSnapshot() { + return wrapped.latestSnapshot(); + } + + @Override + public Snapshot snapshot(long snapshotId) { + return wrapped.snapshot(snapshotId); + } + + @Override + public SimpleFileReader manifestListReader() { + return wrapped.manifestListReader(); + } + + @Override + public SimpleFileReader manifestFileReader() { + return wrapped.manifestFileReader(); + } + + @Override + public SimpleFileReader indexManifestFileReader() { + return wrapped.indexManifestFileReader(); + } + + @Override + public SnapshotReader newSnapshotReader() { + return wrapped.newSnapshotReader(); + } + + @Override + public DataTableScan newScan() { + return wrapped.newScan(); + } + + @Override + public StreamDataTableScan newStreamScan() { + return wrapped.newStreamScan(); + } + + @Override + public CoreOptions coreOptions() { + return wrapped.coreOptions(); + } + + @Override + public Path location() { + return wrapped.location(); + } + + @Override + public SnapshotManager snapshotManager() { + return wrapped.snapshotManager(); + } + + @Override + public ChangelogManager changelogManager() { + return wrapped.changelogManager(); + } + + @Override + public ConsumerManager consumerManager() { + return wrapped.consumerManager(); + } + + @Override + public SchemaManager schemaManager() { + return wrapped.schemaManager(); + } + + @Override + public TagManager tagManager() { + return wrapped.tagManager(); + } + + @Override + public BranchManager branchManager() { + return wrapped.branchManager(); + } + + @Override + public DataTable switchToBranch(String branchName) { + return new ChangelogEventMetadataTable(wrapped.switchToBranch(branchName)); + } + + @Override + public Table copy(Map dynamicOptions) { + return new ChangelogEventMetadataTable(wrapped.copy(dynamicOptions)); + } + + @Override + public FileIO fileIO() { + return wrapped.fileIO(); + } +} diff --git a/paimon-core/src/test/java/org/apache/paimon/mergetree/compact/LookupChangelogMergeFunctionWrapperTest.java b/paimon-core/src/test/java/org/apache/paimon/mergetree/compact/LookupChangelogMergeFunctionWrapperTest.java index db063462fb09..dc2e35de08d5 100644 --- a/paimon-core/src/test/java/org/apache/paimon/mergetree/compact/LookupChangelogMergeFunctionWrapperTest.java +++ b/paimon-core/src/test/java/org/apache/paimon/mergetree/compact/LookupChangelogMergeFunctionWrapperTest.java @@ -555,4 +555,248 @@ public void testKeepLowestHighLevel() { kv = result.result(); assertThat(kv.value().getInt(0)).isEqualTo(3); } + + @Test + public void testPreserveFieldOnRetractDelete() { + // Schema: value has two fields: f0 (data), f1 (event_ts to preserve) + Map highLevel = new HashMap<>(); + RowType valueType = + RowType.builder() + .fields( + new DataType[] {DataTypes.INT(), DataTypes.INT()}, + new String[] {"f0", "f1"}) + .build(); + UserDefinedSeqComparator userDefinedSeqComparator = + UserDefinedSeqComparator.create( + valueType, CoreOptions.fromMap(ImmutableMap.of("sequence.field", "f1"))); + assert userDefinedSeqComparator != null; + + // preserve f1 (index 1) — event metadata appended as extra columns + LookupChangelogMergeFunctionWrapper function = + new LookupChangelogMergeFunctionWrapper( + LookupMergeFunction.wrap( + DeduplicateMergeFunction.factory(), null, null, null), + highLevel::get, + null, + LookupStrategy.from(false, true, false, false), + null, + userDefinedSeqComparator, + new int[] {1}); + + // Delete: -D before-image should retain old values, event metadata appended + highLevel.put(row(1), new KeyValue().replace(row(1), 1, INSERT, row(10, 50)).setLevel(2)); + function.reset(); + function.add(new KeyValue().replace(row(1), 2, DELETE, row(10, 100)).setLevel(0)); + ChangelogResult result = function.getResult(); + assertThat(result).isNotNull(); + List changelogs = result.changelogs(); + assertThat(changelogs).hasSize(1); + assertThat(changelogs.get(0).valueKind()).isEqualTo(DELETE); + InternalRow deleteValue = changelogs.get(0).value(); + // before-image is correct: both fields from old row + assertThat(deleteValue.getInt(0)).isEqualTo(10); + assertThat(deleteValue.getInt(1)).isEqualTo(50); + // event metadata appended at position 2 (field count = 2 + 1 preserved) + assertThat(deleteValue.getFieldCount()).isEqualTo(3); + assertThat(deleteValue.getInt(2)).isEqualTo(100); + // sequence number from the before record + assertThat(changelogs.get(0).sequenceNumber()).isEqualTo(1); + } + + @Test + public void testPreserveFieldOnRetractUpdate() { + // Schema: value has two fields: f0 (data), f1 (event_ts to preserve) + Map highLevel = new HashMap<>(); + RowType valueType = + RowType.builder() + .fields( + new DataType[] {DataTypes.INT(), DataTypes.INT()}, + new String[] {"f0", "f1"}) + .build(); + UserDefinedSeqComparator userDefinedSeqComparator = + UserDefinedSeqComparator.create( + valueType, CoreOptions.fromMap(ImmutableMap.of("sequence.field", "f1"))); + assert userDefinedSeqComparator != null; + + // preserve f1 (index 1) — event metadata appended as extra columns + LookupChangelogMergeFunctionWrapper function = + new LookupChangelogMergeFunctionWrapper( + LookupMergeFunction.wrap( + DeduplicateMergeFunction.factory(), null, null, null), + highLevel::get, + null, + LookupStrategy.from(false, true, false, false), + null, + userDefinedSeqComparator, + new int[] {1}); + + // Update: -U before-image should keep old values, event metadata appended + function.reset(); + function.add(new KeyValue().replace(row(1), 1, INSERT, row(10, 50)).setLevel(1)); + function.add(new KeyValue().replace(row(1), 2, INSERT, row(20, 100)).setLevel(0)); + ChangelogResult result = function.getResult(); + assertThat(result).isNotNull(); + List changelogs = result.changelogs(); + assertThat(changelogs).hasSize(2); + + // -U (UPDATE_BEFORE): correct before-image, event metadata appended + assertThat(changelogs.get(0).valueKind()).isEqualTo(UPDATE_BEFORE); + InternalRow ubValue = changelogs.get(0).value(); + assertThat(ubValue.getInt(0)).isEqualTo(10); + assertThat(ubValue.getInt(1)).isEqualTo(50); + assertThat(ubValue.getFieldCount()).isEqualTo(3); + assertThat(ubValue.getInt(2)).isEqualTo(100); + assertThat(changelogs.get(0).sequenceNumber()).isEqualTo(1); + + // +U (UPDATE_AFTER): event values + metadata (mirrors regular values for schema + // consistency) + assertThat(changelogs.get(1).valueKind()).isEqualTo(UPDATE_AFTER); + InternalRow uaValue = changelogs.get(1).value(); + assertThat(uaValue.getInt(0)).isEqualTo(20); + assertThat(uaValue.getInt(1)).isEqualTo(100); + assertThat(uaValue.getFieldCount()).isEqualTo(3); + assertThat(uaValue.getInt(2)).isEqualTo(100); + } + + @Test + public void testPreserveFieldOnRetractNotConfigured() { + // Verify that the old behavior is preserved when no columns are specified + Map highLevel = new HashMap<>(); + RowType valueType = + RowType.builder() + .fields( + new DataType[] {DataTypes.INT(), DataTypes.INT()}, + new String[] {"f0", "f1"}) + .build(); + UserDefinedSeqComparator userDefinedSeqComparator = + UserDefinedSeqComparator.create( + valueType, CoreOptions.fromMap(ImmutableMap.of("sequence.field", "f1"))); + assert userDefinedSeqComparator != null; + + // no preserve columns (null) + LookupChangelogMergeFunctionWrapper function = + new LookupChangelogMergeFunctionWrapper( + LookupMergeFunction.wrap( + DeduplicateMergeFunction.factory(), null, null, null), + highLevel::get, + null, + LookupStrategy.from(false, true, false, false), + null, + userDefinedSeqComparator, + null); + + // Delete: changelog -D should use old row's values (original behavior, no extra columns) + highLevel.put(row(1), new KeyValue().replace(row(1), 1, INSERT, row(10, 50)).setLevel(2)); + function.reset(); + function.add(new KeyValue().replace(row(1), 2, DELETE, row(10, 100)).setLevel(0)); + ChangelogResult result = function.getResult(); + assertThat(result).isNotNull(); + List changelogs = result.changelogs(); + assertThat(changelogs).hasSize(1); + assertThat(changelogs.get(0).valueKind()).isEqualTo(DELETE); + assertThat(changelogs.get(0).value().getInt(0)).isEqualTo(10); + // f1 should be from the OLD row (original behavior) + assertThat(changelogs.get(0).value().getInt(1)).isEqualTo(50); + // field count should be 2 (no extra columns) + assertThat(changelogs.get(0).value().getFieldCount()).isEqualTo(2); + assertThat(changelogs.get(0).sequenceNumber()).isEqualTo(1); + } + + @Test + public void testPreserveFieldFilterCorrectness() { + // Verify that a downstream WHERE event_ts < 75 filter works correctly: + // After +I(id=1, event_ts=50), update to event_ts=100 should produce + // -U with before-image event_ts=50 (matches filter) so the old row is retracted. + Map highLevel = new HashMap<>(); + RowType valueType = + RowType.builder() + .fields( + new DataType[] {DataTypes.INT(), DataTypes.INT()}, + new String[] {"data", "event_ts"}) + .build(); + UserDefinedSeqComparator userDefinedSeqComparator = + UserDefinedSeqComparator.create( + valueType, + CoreOptions.fromMap(ImmutableMap.of("sequence.field", "event_ts"))); + assert userDefinedSeqComparator != null; + + LookupChangelogMergeFunctionWrapper function = + new LookupChangelogMergeFunctionWrapper( + LookupMergeFunction.wrap( + DeduplicateMergeFunction.factory(), null, null, null), + highLevel::get, + null, + LookupStrategy.from(false, true, false, false), + null, + userDefinedSeqComparator, + new int[] {1}); + + // Simulate: old row has event_ts=50, update event has event_ts=100 + function.reset(); + function.add(new KeyValue().replace(row(1), 1, INSERT, row(10, 50)).setLevel(1)); + function.add(new KeyValue().replace(row(1), 2, INSERT, row(20, 100)).setLevel(0)); + ChangelogResult result = function.getResult(); + List changelogs = result.changelogs(); + assertThat(changelogs).hasSize(2); + + // -U: before-image event_ts=50 — a WHERE event_ts < 75 filter WILL see this retraction + InternalRow ubValue = changelogs.get(0).value(); + int beforeImageEventTs = ubValue.getInt(1); + assertThat(beforeImageEventTs).isEqualTo(50); + assertThat(beforeImageEventTs < 75).isTrue(); + + // Event metadata at position 2: event_ts=100 from the incoming event + assertThat(ubValue.getInt(2)).isEqualTo(100); + } + + @Test + public void testPreserveFieldAggregationCorrectness() { + // Verify that a downstream GROUP BY event_ts aggregation works correctly: + // Retraction must target the OLD group (event_ts=50), not the new group (event_ts=100). + Map highLevel = new HashMap<>(); + RowType valueType = + RowType.builder() + .fields( + new DataType[] {DataTypes.INT(), DataTypes.INT()}, + new String[] {"amount", "event_ts"}) + .build(); + UserDefinedSeqComparator userDefinedSeqComparator = + UserDefinedSeqComparator.create( + valueType, + CoreOptions.fromMap(ImmutableMap.of("sequence.field", "event_ts"))); + assert userDefinedSeqComparator != null; + + LookupChangelogMergeFunctionWrapper function = + new LookupChangelogMergeFunctionWrapper( + LookupMergeFunction.wrap( + DeduplicateMergeFunction.factory(), null, null, null), + highLevel::get, + null, + LookupStrategy.from(false, true, false, false), + null, + userDefinedSeqComparator, + new int[] {1}); + + // Old row: amount=10, event_ts=50 → belongs to group event_ts=50 + // Update: amount=20, event_ts=100 → moves to group event_ts=100 + function.reset(); + function.add(new KeyValue().replace(row(1), 1, INSERT, row(10, 50)).setLevel(1)); + function.add(new KeyValue().replace(row(1), 2, INSERT, row(20, 100)).setLevel(0)); + ChangelogResult result = function.getResult(); + List changelogs = result.changelogs(); + assertThat(changelogs).hasSize(2); + + // -U retraction targets old group: event_ts=50 in before-image + InternalRow retractValue = changelogs.get(0).value(); + int retractGroup = retractValue.getInt(1); + assertThat(retractGroup).isEqualTo(50); + + // +U targets new group: event_ts=100 + InternalRow insertValue = changelogs.get(1).value(); + int insertGroup = insertValue.getInt(1); + assertThat(insertGroup).isEqualTo(100); + + // Groups are different — the aggregation correctly decrements old group and increments new + assertThat(retractGroup).isNotEqualTo(insertGroup); + } } diff --git a/paimon-core/src/test/java/org/apache/paimon/schema/SchemaManagerTest.java b/paimon-core/src/test/java/org/apache/paimon/schema/SchemaManagerTest.java index ea2d4d8c9db1..a64c0081a2a8 100644 --- a/paimon-core/src/test/java/org/apache/paimon/schema/SchemaManagerTest.java +++ b/paimon-core/src/test/java/org/apache/paimon/schema/SchemaManagerTest.java @@ -1783,4 +1783,43 @@ public void testResetCannotWeakenAnOptionThatSetCannotWeaken() { .doesNotThrowAnyException(); } } + + @Test + public void testChangelogMetadataFieldPrefixIsImmutable() { + String key = CoreOptions.CHANGELOG_PRODUCER_METADATA_FIELD_PREFIX.key(); + Map options = new HashMap<>(); + options.put(key, "__internal__"); + + assertThatThrownBy( + () -> + SchemaManager.checkAlterTableOption( + options, key, "__internal__", "__event__")) + .isInstanceOf(UnsupportedOperationException.class) + .hasMessageContaining(key); + assertThatThrownBy(() -> SchemaManager.checkResetTableOption(options, key)) + .isInstanceOf(UnsupportedOperationException.class) + .hasMessageContaining(key); + } + + @Test + public void testRenameChangelogMetadataSourceColumn() throws Exception { + Map metadataOptions = new HashMap<>(); + metadataOptions.put( + CoreOptions.CHANGELOG_PRODUCER.key(), + CoreOptions.ChangelogProducer.LOOKUP.toString()); + metadataOptions.put(CoreOptions.CHANGELOG_PRODUCER_EVENT_METADATA_FIELDS.key(), "f2, f1"); + Schema metadataSchema = + new Schema(rowType.getFields(), partitionKeys, primaryKeys, metadataOptions, ""); + retryArtificialException(() -> manager.createTable(metadataSchema)); + + retryArtificialException( + () -> manager.commitChanges(SchemaChange.renameColumn("f2", "renamed_f2"))); + + TableSchema latest = retryArtificialException(() -> manager.latest()).get(); + assertThat(latest.options()) + .containsEntry( + CoreOptions.CHANGELOG_PRODUCER_EVENT_METADATA_FIELDS.key(), + "renamed_f2,f1"); + assertThat(latest.logicalRowType().getField("renamed_f2").id()).isEqualTo(2); + } } diff --git a/paimon-core/src/test/java/org/apache/paimon/schema/SchemaValidationTest.java b/paimon-core/src/test/java/org/apache/paimon/schema/SchemaValidationTest.java index 681edbf1dac0..b98ad12018df 100644 --- a/paimon-core/src/test/java/org/apache/paimon/schema/SchemaValidationTest.java +++ b/paimon-core/src/test/java/org/apache/paimon/schema/SchemaValidationTest.java @@ -2279,6 +2279,48 @@ public void testFullCompactionDeltaCommitsWithLookupChangelogProducer() { assertThatCode(() -> validateTableSchemaExec(options)).doesNotThrowAnyException(); } + @Test + public void testExposeFieldAsMetadataOnlySupportsLookupChangelogProducer() { + Map options = new HashMap<>(); + options.put(CoreOptions.CHANGELOG_PRODUCER_EVENT_METADATA_FIELDS.key(), "f1"); + + for (String producer : Arrays.asList("none", "input", "full-compaction")) { + options.put(CoreOptions.CHANGELOG_PRODUCER.key(), producer); + assertThatThrownBy(() -> validateTableSchemaExec(options)) + .isInstanceOf(IllegalArgumentException.class) + .hasMessageContaining( + CoreOptions.CHANGELOG_PRODUCER_EVENT_METADATA_FIELDS.key()) + .hasMessageContaining(CoreOptions.CHANGELOG_PRODUCER.key()) + .hasMessageContaining("lookup"); + } + + options.put(CoreOptions.CHANGELOG_PRODUCER.key(), "lookup"); + assertThatCode(() -> validateTableSchemaExec(options)).doesNotThrowAnyException(); + } + + @Test + public void testExposeFieldAsMetadataRejectsInvalidColumns() { + Map options = new HashMap<>(); + options.put(CoreOptions.CHANGELOG_PRODUCER.key(), "lookup"); + + options.put(CoreOptions.CHANGELOG_PRODUCER_EVENT_METADATA_FIELDS.key(), "unknown"); + assertThatThrownBy(() -> validateTableSchemaExec(options)) + .isInstanceOf(IllegalArgumentException.class) + .hasMessageContaining("unknown") + .hasMessageContaining("not found"); + + options.put(CoreOptions.CHANGELOG_PRODUCER_EVENT_METADATA_FIELDS.key(), "f1,f1"); + assertThatThrownBy(() -> validateTableSchemaExec(options)) + .isInstanceOf(IllegalArgumentException.class) + .hasMessageContaining("more than once"); + + options.put(CoreOptions.CHANGELOG_PRODUCER_EVENT_METADATA_FIELDS.key(), "f1"); + options.put(CoreOptions.CHANGELOG_PRODUCER_METADATA_FIELD_PREFIX.key(), ""); + assertThatThrownBy(() -> validateTableSchemaExec(options)) + .isInstanceOf(IllegalArgumentException.class) + .hasMessageContaining("conflicts"); + } + private Map mapSharedShreddingOptions() { Map options = new HashMap<>(); options.put(BUCKET.key(), "-1"); diff --git a/paimon-core/src/test/java/org/apache/paimon/table/system/ChangelogEventMetadataTest.java b/paimon-core/src/test/java/org/apache/paimon/table/system/ChangelogEventMetadataTest.java new file mode 100644 index 000000000000..bf2bf6ab599a --- /dev/null +++ b/paimon-core/src/test/java/org/apache/paimon/table/system/ChangelogEventMetadataTest.java @@ -0,0 +1,132 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.paimon.table.system; + +import org.apache.paimon.CoreOptions; +import org.apache.paimon.options.Options; +import org.apache.paimon.types.DataField; +import org.apache.paimon.types.DataTypes; +import org.apache.paimon.types.RowType; + +import org.junit.jupiter.api.Test; + +import java.util.Arrays; + +import static org.assertj.core.api.Assertions.assertThat; + +/** Tests for changelog event metadata row-type construction. */ +class ChangelogEventMetadataTest { + + @Test + void testMetadataUsesPhysicalFieldsWhenBaseRowPrependsSystemFields() { + RowType valueType = + new RowType( + Arrays.asList( + new DataField(0, "id", DataTypes.INT()), + new DataField(1, "event_ts", DataTypes.BIGINT()))); + RowType baseRowType = + new RowType( + Arrays.asList( + new DataField(100, "rowkind", DataTypes.STRING()), + valueType.getField("id"), + valueType.getField("event_ts"))); + Options options = new Options(); + options.set(CoreOptions.CHANGELOG_PRODUCER, CoreOptions.ChangelogProducer.LOOKUP); + options.set(CoreOptions.CHANGELOG_PRODUCER_EVENT_METADATA_FIELDS, "event_ts"); + CoreOptions coreOptions = new CoreOptions(options); + + RowType extended = + ChangelogEventMetadata.appendMetadataFields(baseRowType, valueType, coreOptions); + + assertThat(extended.getFieldNames()) + .containsExactly("rowkind", "id", "event_ts", "__internal__event_ts"); + assertThat(extended.getField("__internal__event_ts").type().isNullable()).isTrue(); + assertThat(extended.getField("__internal__event_ts").id()).isEqualTo(2); + } + + @Test + void testMetadataFieldIdsIncludeNestedFields() { + RowType valueType = + new RowType( + Arrays.asList( + new DataField(0, "id", DataTypes.INT()), + new DataField( + 1, + "payload", + DataTypes.ROW(new DataField(3, "nested", DataTypes.INT()))), + new DataField(2, "event_ts", DataTypes.BIGINT()))); + Options options = new Options(); + options.set(CoreOptions.CHANGELOG_PRODUCER, CoreOptions.ChangelogProducer.LOOKUP); + options.set(CoreOptions.CHANGELOG_PRODUCER_EVENT_METADATA_FIELDS, "event_ts"); + CoreOptions coreOptions = new CoreOptions(options); + + RowType extended = + ChangelogEventMetadata.appendMetadataFields(valueType, valueType, coreOptions); + + assertThat(extended.getField("__internal__event_ts").id()).isEqualTo(4); + assertThat(RowType.currentHighestFieldId(extended.getFields())).isEqualTo(4); + } + + @Test + void testStorageFieldIdentityIsStableAcrossColumnRename() { + RowType originalValueType = + new RowType( + Arrays.asList( + new DataField(0, "id", DataTypes.INT()), + new DataField(1, "event_ts", DataTypes.BIGINT()))); + Options originalOptions = new Options(); + originalOptions.set(CoreOptions.CHANGELOG_PRODUCER, CoreOptions.ChangelogProducer.LOOKUP); + originalOptions.set(CoreOptions.CHANGELOG_PRODUCER_EVENT_METADATA_FIELDS, "event_ts"); + originalOptions.set(CoreOptions.CHANGELOG_PRODUCER_METADATA_FIELD_PREFIX, "__event__"); + + RowType renamedValueType = + new RowType( + Arrays.asList( + new DataField(0, "id", DataTypes.INT()), + new DataField(1, "event_time", DataTypes.BIGINT()))); + Options renamedOptions = new Options(); + renamedOptions.set(CoreOptions.CHANGELOG_PRODUCER, CoreOptions.ChangelogProducer.LOOKUP); + renamedOptions.set(CoreOptions.CHANGELOG_PRODUCER_EVENT_METADATA_FIELDS, "event_time"); + renamedOptions.set(CoreOptions.CHANGELOG_PRODUCER_METADATA_FIELD_PREFIX, "__event__"); + + DataField originalPublicField = + ChangelogEventMetadata.extraValueFields( + originalValueType, new CoreOptions(originalOptions)) + .get(0); + DataField renamedPublicField = + ChangelogEventMetadata.extraValueFields( + renamedValueType, new CoreOptions(renamedOptions)) + .get(0); + DataField originalStorageField = + ChangelogEventMetadata.storageValueFields( + originalValueType, new CoreOptions(originalOptions)) + .get(0); + DataField renamedStorageField = + ChangelogEventMetadata.storageValueFields( + renamedValueType, new CoreOptions(renamedOptions)) + .get(0); + + assertThat(originalPublicField.name()).isEqualTo("__event__event_ts"); + assertThat(renamedPublicField.name()).isEqualTo("__event__event_time"); + assertThat(originalPublicField.id()).isEqualTo(renamedPublicField.id()); + assertThat(originalStorageField.name()).isEqualTo("__event__field_id_1"); + assertThat(renamedStorageField.name()).isEqualTo(originalStorageField.name()); + assertThat(renamedStorageField.id()).isEqualTo(originalStorageField.id()); + } +} diff --git a/paimon-e2e-tests/src/test/java/org/apache/paimon/tests/SparkE2eTest.java b/paimon-e2e-tests/src/test/java/org/apache/paimon/tests/SparkE2eTest.java index ad2e577041e6..04afb95aa1b6 100644 --- a/paimon-e2e-tests/src/test/java/org/apache/paimon/tests/SparkE2eTest.java +++ b/paimon-e2e-tests/src/test/java/org/apache/paimon/tests/SparkE2eTest.java @@ -64,35 +64,70 @@ public void testFlinkWriteAndSparkRead() throws Exception { createCatalogSql("my_spark", warehousePath), createTableSql(table), createInsertSql(table))); - checkQueryResults( - sparkTable, - sql -> { - Container.ExecResult execResult = - getSpark() - .execInContainer( - "/spark/bin/spark-sql", - "--master", - "spark://spark-master:7077", - "--conf", - "spark.sql.extensions=org.apache.paimon.spark.extensions.PaimonSparkSessionExtensions", - "--conf", - "spark.sql.catalog.paimon=org.apache.paimon.spark.SparkCatalog", - "--conf", - "spark.sql.catalog.paimon.warehouse=file:" - + warehousePath, - "-f", - TEST_DATA_DIR + "/" + sql); - if (execResult.getExitCode() != 0) { - LOG.info(execResult.getStdout()); - LOG.info(execResult.getStderr()); - throw new AssertionError("Failed when running spark sql."); - } - String stdout = stripTrailingSparkErrorLogs(execResult.getStdout()); - return Arrays.stream(stdout.split("\n")) - .filter(s -> !s.contains("WARN")) - .collect(Collectors.joining("\n")) - + "\n"; - }); + checkQueryResults(sparkTable, sql -> executeSparkSql(warehousePath, sql)); + } + + @Test + public void testFlinkCreateAndSparkReadChangelogEventMetadata() throws Exception { + String warehousePath = TEST_DATA_DIR + "/" + UUID.randomUUID() + "_warehouse"; + final String table = "event_metadata"; + final String sparkTable = String.format("paimon.default.%s", table); + + runBatchSql( + String.join( + "\n", + createCatalogSql("my_flink", warehousePath), + "CREATE TABLE " + + table + + " (" + + " id INT," + + " data INT," + + " event_ts BIGINT," + + " PRIMARY KEY (id) NOT ENFORCED" + + ") WITH (" + + " 'bucket' = '1'," + + " 'changelog-producer' = 'lookup'," + + " 'sequence.field' = 'event_ts'," + + " 'changelog-producer.event-metadata-fields' = 'event_ts'" + + ");", + "INSERT INTO " + table + " VALUES (1, 10, 50);", + "INSERT INTO " + table + " VALUES (1, 20, 100);")); + + // Flink created the table without a METADATA FROM alias. Spark reads the generated field + // using the physical metadata name stored in the table properties. + checkQueryResult( + sql -> executeSparkSql(warehousePath, sql), + "SELECT id, data, event_ts, __internal__event_ts FROM " + + sparkTable + + " ORDER BY id", + "1\t20\t100\t100\n"); + } + + private String executeSparkSql(String warehousePath, String sqlFile) throws Exception { + Container.ExecResult execResult = + getSpark() + .execInContainer( + "/spark/bin/spark-sql", + "--master", + "spark://spark-master:7077", + "--conf", + "spark.sql.extensions=org.apache.paimon.spark.extensions.PaimonSparkSessionExtensions", + "--conf", + "spark.sql.catalog.paimon=org.apache.paimon.spark.SparkCatalog", + "--conf", + "spark.sql.catalog.paimon.warehouse=file:" + warehousePath, + "-f", + TEST_DATA_DIR + "/" + sqlFile); + if (execResult.getExitCode() != 0) { + LOG.info(execResult.getStdout()); + LOG.info(execResult.getStderr()); + throw new AssertionError("Failed when running spark sql."); + } + String stdout = stripTrailingSparkErrorLogs(execResult.getStdout()); + return Arrays.stream(stdout.split("\n")) + .filter(s -> !s.contains("WARN")) + .collect(Collectors.joining("\n")) + + "\n"; } /** diff --git a/paimon-flink/paimon-flink-common/pom.xml b/paimon-flink/paimon-flink-common/pom.xml index f0a38d841743..f975962e4650 100644 --- a/paimon-flink/paimon-flink-common/pom.xml +++ b/paimon-flink/paimon-flink-common/pom.xml @@ -190,6 +190,13 @@ under the License. test + + org.apache.avro + avro + ${avro.version} + test + + org.apache.iceberg iceberg-data diff --git a/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/AbstractFlinkTableFactory.java b/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/AbstractFlinkTableFactory.java index 137a7d753431..59dbc44ed3e5 100644 --- a/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/AbstractFlinkTableFactory.java +++ b/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/AbstractFlinkTableFactory.java @@ -99,7 +99,8 @@ && parseBoolean(options.get(SCAN_BOUNDED.key()))) { } if (origin instanceof SystemCatalogTable) { return new SystemTableSource(table, unbounded, context.getObjectIdentifier()); - } else if (CoreOptions.fromMap(table.options()).dataEvolutionEnabled()) { + } else if (CoreOptions.fromMap(table.options()).dataEvolutionEnabled() + || !CoreOptions.fromMap(table.options()).changelogEventMetadataFields().isEmpty()) { return new DataEvolutionDataTableSource( context.getObjectIdentifier(), table, unbounded, context); } else { diff --git a/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/source/DataTableSource.java b/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/source/DataTableSource.java index 7a18546a1fd6..1e199ec2d18e 100644 --- a/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/source/DataTableSource.java +++ b/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/source/DataTableSource.java @@ -19,6 +19,7 @@ package org.apache.paimon.flink.source; import org.apache.paimon.CoreOptions; +import org.apache.paimon.flink.LogicalTypeConversion; import org.apache.paimon.flink.PaimonDataStreamScanProvider; import org.apache.paimon.flink.Projection; import org.apache.paimon.flink.dataevolution.DataEvolutionRowLevelModificationScanContext; @@ -31,7 +32,10 @@ import org.apache.paimon.table.SpecialFields; import org.apache.paimon.table.Table; import org.apache.paimon.table.source.snapshot.TimeTravelUtil; +import org.apache.paimon.table.system.ChangelogEventMetadata; +import org.apache.paimon.table.system.ChangelogEventMetadataTable; import org.apache.paimon.table.system.RowTrackingTable; +import org.apache.paimon.types.DataField; import org.apache.flink.api.common.eventtime.WatermarkStrategy; import org.apache.flink.table.api.DataTypes; @@ -47,6 +51,7 @@ import org.apache.flink.table.plan.stats.ColumnStats; import org.apache.flink.table.plan.stats.TableStats; import org.apache.flink.table.types.DataType; +import org.apache.flink.table.types.utils.TypeConversions; import javax.annotation.Nullable; @@ -173,25 +178,60 @@ public RowLevelModificationScanContext applyRowLevelModificationScan( } public Map listReadableMetadata() { - // Flink calls this after applyRowLevelModificationScan for row-level operations. - if (rowLevelModificationSnapshotId == null || !isDataEvolutionTable()) { - return Collections.emptyMap(); - } Map metadata = new LinkedHashMap<>(); - metadata.put(SpecialFields.ROW_ID.name(), DataTypes.BIGINT().notNull()); + + // Row-level modification metadata + if (rowLevelModificationSnapshotId != null && isDataEvolutionTable()) { + metadata.put(SpecialFields.ROW_ID.name(), DataTypes.BIGINT().notNull()); + } + + // Event metadata from changelog-producer.event-metadata-fields + List preserveColumns = eventPreserveColumns(); + if (!preserveColumns.isEmpty() && table instanceof FileStoreTable) { + org.apache.paimon.types.RowType valueType = + ((FileStoreTable) table).schema().logicalRowType(); + CoreOptions coreOptions = CoreOptions.fromMap(table.options()); + for (DataField field : + ChangelogEventMetadata.extraValueFields(valueType, coreOptions)) { + DataType flinkType = + TypeConversions.fromLogicalToDataType( + LogicalTypeConversion.toLogicalType(field.type())); + metadata.put(field.name(), flinkType.nullable()); + } + } + return metadata; } public void applyReadableMetadata(List metadataKeys, DataType producedDataType) { for (String metadataKey : metadataKeys) { - if (!SpecialFields.ROW_ID.name().equals(metadataKey)) { - throw new UnsupportedOperationException( - "Unsupported Paimon metadata column: " + metadataKey); + if (SpecialFields.ROW_ID.name().equals(metadataKey) + || eventMetadataFieldNames().contains(metadataKey)) { + continue; } + throw new UnsupportedOperationException( + "Unsupported Paimon metadata column: " + metadataKey); } this.metadataKeys = metadataKeys; } + private List eventPreserveColumns() { + return CoreOptions.fromMap(table.options()).changelogEventMetadataFields(); + } + + private List eventMetadataFieldNames() { + if (!(table instanceof FileStoreTable)) { + return Collections.emptyList(); + } + org.apache.paimon.types.RowType valueType = + ((FileStoreTable) table).schema().logicalRowType(); + return ChangelogEventMetadata.extraValueFields( + valueType, CoreOptions.fromMap(table.options())) + .stream() + .map(DataField::name) + .collect(Collectors.toList()); + } + @Override public ScanRuntimeProvider getScanRuntimeProvider(ScanContext scanContext) { if (rowLevelModificationSnapshotId == null @@ -223,7 +263,7 @@ public ScanRuntimeProvider getScanRuntimeProvider(ScanContext scanContext) { @Override protected Table tableForScan() { if (rowLevelModificationSnapshotId == null) { - return table; + return wrapForEventMetadata(table); } FileStoreTable fileStoreTable = (FileStoreTable) table; @@ -237,7 +277,19 @@ protected Table tableForScan() { fileStoreTable = (FileStoreTable) fileStoreTable.copy(options); } - return metadataKeys.isEmpty() ? fileStoreTable : new RowTrackingTable(fileStoreTable); + if (metadataKeys.isEmpty()) { + return fileStoreTable; + } + return new RowTrackingTable(fileStoreTable); + } + + private Table wrapForEventMetadata(Table scanTable) { + boolean hasEventMetadata = + metadataKeys.stream().anyMatch(k -> eventMetadataFieldNames().contains(k)); + if (hasEventMetadata && scanTable instanceof FileStoreTable) { + return new ChangelogEventMetadataTable((FileStoreTable) scanTable); + } + return scanTable; } @Override @@ -255,10 +307,22 @@ protected int[][] projectFieldsForScan() { } } + List metadataFieldNames = eventMetadataFieldNames(); int[][] projection = Arrays.copyOf(physicalProjection, physicalProjection.length + metadataKeys.size()); for (int i = 0; i < metadataKeys.size(); i++) { - projection[physicalProjection.length + i] = new int[] {physicalFieldCount}; + String key = metadataKeys.get(i); + if (metadataFieldNames.contains(key)) { + int preserveIdx = metadataFieldNames.indexOf(key); + if (preserveIdx < 0) { + throw new UnsupportedOperationException( + "Unknown event metadata column: " + key); + } + projection[physicalProjection.length + i] = + new int[] {physicalFieldCount + preserveIdx}; + } else { + projection[physicalProjection.length + i] = new int[] {physicalFieldCount}; + } } return projection; } diff --git a/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/LookupChangelogEventMetadataITCase.java b/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/LookupChangelogEventMetadataITCase.java new file mode 100644 index 000000000000..8371d7fe1fe1 --- /dev/null +++ b/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/LookupChangelogEventMetadataITCase.java @@ -0,0 +1,196 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.paimon.flink; + +import org.apache.paimon.utils.BlockingIterator; + +import org.apache.flink.types.Row; +import org.apache.flink.types.RowKind; +import org.junit.jupiter.api.Test; + +import java.util.stream.Collectors; + +import static org.assertj.core.api.Assertions.assertThat; + +/** End-to-end tests for lookup changelog event metadata. */ +public class LookupChangelogEventMetadataITCase extends CatalogITCaseBase { + + @Test + public void testEventMetadataCanBeReadAsWriteTime() throws Exception { + sql( + "CREATE TABLE source_table (" + + "id INT PRIMARY KEY NOT ENFORCED, " + + "data INT, " + + "event_ts BIGINT, " + + "writetime BIGINT METADATA FROM '__internal__event_ts' VIRTUAL" + + ") WITH (" + + "'bucket'='1', " + + "'changelog-producer'='lookup', " + + "'sequence.field'='event_ts', " + + "'changelog-producer.event-metadata-fields'='event_ts')"); + + // The metadata column models an external sink populated from an event timestamp. The + // physical event_ts column remains available to normal Flink operators, while writetime + // carries the incoming event value even on an UPDATE_BEFORE retraction. + BlockingIterator iterator = + streamSqlBlockIter("SELECT id, data, event_ts, writetime FROM source_table"); + + sql("INSERT INTO source_table VALUES (1, 10, 50)"); + assertThat(iterator.collect(1)).containsExactly(Row.of(1, 10, 50L, 50L)); + + sql("INSERT INTO source_table VALUES (1, 20, 100)"); + assertThat(iterator.collect(2)) + .containsExactly( + Row.ofKind(RowKind.UPDATE_BEFORE, 1, 10, 50L, 100L), + Row.ofKind(RowKind.UPDATE_AFTER, 1, 20, 100L, 100L)); + + iterator.close(); + } + + @Test + public void testPhysicalFilterStillSeesBeforeImage() throws Exception { + sql( + "CREATE TABLE filtered_source (" + + "id INT PRIMARY KEY NOT ENFORCED, " + + "data INT, " + + "event_ts BIGINT, " + + "writetime BIGINT METADATA FROM '__event__event_ts' VIRTUAL" + + ") WITH (" + + "'bucket'='1', " + + "'changelog-producer'='lookup', " + + "'sequence.field'='event_ts', " + + "'changelog-producer.metadata-field-prefix'='__event__', " + + "'changelog-producer.event-metadata-fields'='event_ts')"); + + BlockingIterator iterator = + streamSqlBlockIter( + "SELECT id, data, event_ts, writetime " + + "FROM filtered_source WHERE event_ts < 75"); + + sql("INSERT INTO filtered_source VALUES (1, 10, 50)"); + assertThat(iterator.collect(1)).containsExactly(Row.of(1, 10, 50L, 50L)); + + // The update-after value is filtered out, but the update-before must retain the old + // physical event_ts so that the downstream filter can retract the old row. Its metadata + // value remains the incoming event timestamp, which is the value an external sink needs. + sql("INSERT INTO filtered_source VALUES (1, 20, 100)"); + assertThat(iterator.collect(1)) + .containsExactly(Row.ofKind(RowKind.UPDATE_BEFORE, 1, 10, 50L, 100L)); + + iterator.close(); + } + + @Test + public void testNestedValueCanBeReadWithFlinkMetadataAlias() throws Exception { + sql( + "CREATE TABLE nested_source (" + + "id INT PRIMARY KEY NOT ENFORCED, " + + "event_ts BIGINT, " + + "payload ROW, " + + "writetime BIGINT METADATA FROM '__internal__event_ts' VIRTUAL" + + ") WITH (" + + "'bucket'='1', " + + "'changelog-producer'='lookup', " + + "'sequence.field'='event_ts', " + + "'changelog-producer.event-metadata-fields'='event_ts')"); + + BlockingIterator iterator = + streamSqlBlockIter("SELECT id, event_ts, payload, writetime FROM nested_source"); + + sql("INSERT INTO nested_source VALUES " + "(1, 50, CAST(ROW(10) AS ROW))"); + assertThat(iterator.collect(1)).containsExactly(Row.of(1, 50L, Row.of(10), 50L)); + + sql("INSERT INTO nested_source VALUES " + "(1, 100, CAST(ROW(20) AS ROW))"); + assertThat(iterator.collect(2)) + .containsExactly( + Row.ofKind(RowKind.UPDATE_BEFORE, 1, 50L, Row.of(10), 100L), + Row.ofKind(RowKind.UPDATE_AFTER, 1, 100L, Row.of(20), 100L)); + + iterator.close(); + } + + @Test + public void testPhysicalTableCanBeReadWithFlinkMetadataAlias() throws Exception { + String tableName = "physical_table"; + sql( + "CREATE TABLE " + + tableName + + " (" + + "id INT PRIMARY KEY NOT ENFORCED, " + + "data INT, " + + "event_ts BIGINT" + + ") WITH (" + + "'bucket'='1', " + + "'changelog-producer'='lookup', " + + "'sequence.field'='event_ts', " + + "'changelog-producer.event-metadata-fields'='event_ts')"); + + // A table created without a Flink metadata alias stores only its physical columns. Verify + // that Flink can read that schema before registering a metadata alias. + assertThat( + table(tableName).getUnresolvedSchema().getColumns().stream() + .map(column -> column.getName()) + .collect(Collectors.toList())) + .containsExactly("id", "data", "event_ts"); + + // Without a metadata column in the Flink schema, a wildcard sees only the three physical + // columns and the changelog contains the same three-column shape. + BlockingIterator physicalIterator = + streamSqlBlockIter("SELECT * FROM " + tableName); + + sql("INSERT INTO " + tableName + " VALUES (1, 10, 50)"); + assertThat(physicalIterator.collect(1)).containsExactly(Row.of(1, 10, 50L)); + assertThat(sql("SELECT * FROM " + tableName)).containsExactly(Row.of(1, 10, 50L)); + + // Read the same physical table through a separate Flink connector definition. This is + // equivalent to registering a metadata alias when consuming a table created by another + // engine. + sEnv.executeSql( + String.format( + "CREATE TEMPORARY TABLE flink_reader (" + + "id INT PRIMARY KEY NOT ENFORCED, " + + "data INT, " + + "event_ts BIGINT, " + + "writetime BIGINT METADATA FROM '__internal__event_ts' VIRTUAL" + + ") WITH (" + + "'connector'='paimon', " + + "'path'='%s')", + getTableDirectory(tableName))); + + // A wildcard projection should include the physical columns followed by the declared + // metadata alias, so verify the same shape while reading the changelog. + BlockingIterator iterator = streamSqlBlockIter("SELECT * FROM flink_reader"); + + assertThat(iterator.collect(1)).containsExactly(Row.of(1, 10, 50L, 50L)); + + sql("INSERT INTO " + tableName + " VALUES (1, 20, 100)"); + assertThat(iterator.collect(2)) + .containsExactly( + Row.ofKind(RowKind.UPDATE_BEFORE, 1, 10, 50L, 100L), + Row.ofKind(RowKind.UPDATE_AFTER, 1, 20, 100L, 100L)); + + assertThat(physicalIterator.collect(2)) + .containsExactly( + Row.ofKind(RowKind.UPDATE_BEFORE, 1, 10, 50L), + Row.ofKind(RowKind.UPDATE_AFTER, 1, 20, 100L)); + + physicalIterator.close(); + iterator.close(); + } +} diff --git a/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/PaimonSparkTableBase.scala b/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/PaimonSparkTableBase.scala index d42d129925dd..339537967d33 100644 --- a/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/PaimonSparkTableBase.scala +++ b/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/PaimonSparkTableBase.scala @@ -27,8 +27,9 @@ import org.apache.paimon.spark.read.PaimonSplitScanBuilder import org.apache.paimon.spark.schema.PaimonMetadataColumn import org.apache.paimon.spark.util.OptionUtils import org.apache.paimon.spark.write.{PaimonV2WriteBuilder, PaimonWriteBuilder} -import org.apache.paimon.table.{CatalogTableType, Table, _} +import org.apache.paimon.table.{CatalogTableType, FileStoreTable, Table, _} import org.apache.paimon.table.BucketMode.{BUCKET_UNAWARE, HASH_FIXED, POSTPONE_MODE} +import org.apache.paimon.table.system.ChangelogEventMetadataTable import org.apache.spark.sql.connector.catalog._ import org.apache.spark.sql.connector.read.ScanBuilder @@ -48,6 +49,16 @@ abstract class PaimonSparkTableBase(val table: Table) lazy val coreOptions = new CoreOptions(table.options()) + override lazy val schema: org.apache.spark.sql.types.StructType = { + val baseRowType = table.rowType() + val extendedRowType = table match { + case fst: FileStoreTable if !coreOptions.changelogEventMetadataFields().isEmpty => + ChangelogEventMetadataTable.computeExtendedRowType(fst, baseRowType) + case _ => baseRowType + } + SparkTypeUtils.fromPaimonRowType(extendedRowType) + } + lazy val useV2Write: Boolean = { val v2WriteConfigured = OptionUtils.useV2Write() v2WriteConfigured && supportsV2Write diff --git a/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/catalyst/analysis/PaimonAnalysis.scala b/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/catalyst/analysis/PaimonAnalysis.scala index 527ca1faad64..43875bb0d8ec 100644 --- a/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/catalyst/analysis/PaimonAnalysis.scala +++ b/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/catalyst/analysis/PaimonAnalysis.scala @@ -25,6 +25,7 @@ import org.apache.paimon.spark.catalyst.analysis.PaimonRelation.isPaimonTable import org.apache.paimon.spark.catalyst.plans.logical.{PaimonDropPartitions, PaimonHiveDynamicPartitionQuery} import org.apache.paimon.spark.commands.{PaimonAnalyzeFormatTablePartitionsCommand, PaimonAnalyzeTableColumnCommand, PaimonDynamicPartitionOverwriteCommand, PaimonShowColumnsCommand, SchemaEvolutionHelper} import org.apache.paimon.spark.format.PaimonFormatTable +import org.apache.paimon.spark.schema.SparkSystemColumns import org.apache.paimon.spark.util.OptionUtils import org.apache.paimon.table.FileStoreTable @@ -119,7 +120,8 @@ class PaimonAnalysis(session: SparkSession) extends Rule[LogicalPlan] { table: DataSourceV2Relation, options: Options, mergeSchemaEnabled: Boolean): LogicalPlan = { - val query = stripHiveDynamicPartitionMarker(v2WriteCommand.query) + val query = + stripChangelogMetadataColumns(stripHiveDynamicPartitionMarker(v2WriteCommand.query), table) val hiveStyleDynamicPartitionEnabled = OptionUtils.hiveStyleDynamicPartitionEnabled() hiveDynamicPartitionColumns(v2WriteCommand.query) match { case Some(dynamicPartitionColumns) @@ -163,16 +165,30 @@ class PaimonAnalysis(session: SparkSession) extends Rule[LogicalPlan] { query.transformDown { case PaimonHiveDynamicPartitionQuery(_, child) => child } } + /** Synthetic event metadata is readable but is not part of a physical write row. */ + private def stripChangelogMetadataColumns( + query: LogicalPlan, + table: DataSourceV2Relation): LogicalPlan = { + table.table.asInstanceOf[SparkTable].getTable match { + case fileStoreTable: FileStoreTable => + val physicalOutput = + SparkSystemColumns.filterChangelogMetadataColumns(query.output, fileStoreTable) + if (physicalOutput.size == query.output.size) query else Project(physicalOutput, query) + case _ => query + } + } + private def resolveDynamicPartitionWrite( query: LogicalPlan, table: DataSourceV2Relation, hiveStyleOutput: Option[Seq[Attribute]], options: Options, mergeSchemaEnabled: Boolean): LogicalPlan = { + val physicalTableOutput = physicalOutput(table) hiveStyleOutput match { case Some(hiveStyleOutput) - if !sameOutputNames(query.output, table.output) && - !sameOutputNames(hiveStyleOutput, table.output) => + if !sameOutputNames(query.output, physicalTableOutput) && + !sameOutputNames(hiveStyleOutput, physicalTableOutput) => val hiveStyleQuery = resolveWriteOutput(query, table.name, hiveStyleOutput, byName = false, mergeSchemaEnabled) resolveWriteOutput( @@ -235,6 +251,7 @@ class PaimonAnalysis(session: SparkSession) extends Rule[LogicalPlan] { private def hiveStyleDynamicPartitionOutput( table: DataSourceV2Relation, dynamicPartitionColumns: Seq[String]): Option[Seq[Attribute]] = { + val physicalTableOutput = physicalOutput(table) val partitionKeys = table.table.asInstanceOf[SparkTable].getTable.partitionKeys().asScala.toSeq if (partitionKeys.isEmpty || dynamicPartitionColumns.isEmpty) { None @@ -244,9 +261,10 @@ class PaimonAnalysis(session: SparkSession) extends Rule[LogicalPlan] { partition => dynamicPartitionColumns.exists(dynamic => conf.resolver(dynamic, partition)) } .flatMap { - dynamicPartition => table.output.find(attr => conf.resolver(attr.name, dynamicPartition)) + dynamicPartition => + physicalTableOutput.find(attr => conf.resolver(attr.name, dynamicPartition)) } - val dataAttrs = table.output.filterNot { + val dataAttrs = physicalTableOutput.filterNot { attr => dynamicPartitionColumns.exists(partition => conf.resolver(attr.name, partition)) } val hiveStyleOutput = dataAttrs ++ dynamicPartitionAttrs @@ -258,6 +276,14 @@ class PaimonAnalysis(session: SparkSession) extends Rule[LogicalPlan] { } } + private def physicalOutput(table: DataSourceV2Relation): Seq[Attribute] = { + table.table.asInstanceOf[SparkTable].getTable match { + case fileStoreTable: FileStoreTable => + SparkSystemColumns.filterChangelogMetadataColumns(table.output, fileStoreTable) + case _ => table.output + } + } + private def sameOutputNames(left: Seq[Attribute], right: Seq[Attribute]): Boolean = { left.length == right.length && left.zip(right).forall { case (l, r) => conf.resolver(l.name, r.name) } diff --git a/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/commands/SchemaEvolutionHelper.scala b/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/commands/SchemaEvolutionHelper.scala index 7bfb7ee90501..c34d8494d3d0 100644 --- a/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/commands/SchemaEvolutionHelper.scala +++ b/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/commands/SchemaEvolutionHelper.scala @@ -99,7 +99,7 @@ private[spark] object SchemaEvolutionHelper { table: FileStoreTable, dataSchema: StructType, flags: SchemaEvolutionFlags): Option[TableSchema] = { - val filtered = SparkSystemColumns.filterSparkSystemColumns(dataSchema) + val filtered = SparkSystemColumns.filterSparkSystemColumns(dataSchema, table) val dataRowType = SparkTypeUtils.toPaimonType(filtered).asInstanceOf[RowType] val current = table.schema() val merged = @@ -122,7 +122,7 @@ private[spark] object SchemaEvolutionHelper { dataSchema: StructType, sparkSession: SparkSession, options: Options = new Options()): Boolean = { - val filtered = SparkSystemColumns.filterSparkSystemColumns(dataSchema) + val filtered = SparkSystemColumns.filterSparkSystemColumns(dataSchema, table) val flags = readFlags(sparkSession, options) val dataRowType = SparkTypeUtils.toPaimonType(filtered).asInstanceOf[RowType] table @@ -177,14 +177,20 @@ private[spark] object SchemaEvolutionHelper { isByName: Boolean, sparkSession: SparkSession): Seq[Attribute] = { val flags = readFlags(sparkSession, options) - if (!isByName || !mergeSchemaEnabled(options) || !flags.typeWidening) return table.output + val physicalOutput = + table.table.asInstanceOf[SparkTable].getTable match { + case fileStoreTable: FileStoreTable => + SparkSystemColumns.filterChangelogMetadataColumns(table.output, fileStoreTable) + case _ => table.output + } + if (!isByName || !mergeSchemaEnabled(options) || !flags.typeWidening) return physicalOutput table.table.asInstanceOf[SparkTable].getTable match { case fst: FileStoreTable => computeMergedSchema(fst, querySchema, flags) .map(s => PaimonUtils.toAttributes(SparkTypeUtils.fromPaimonRowType(s.logicalRowType()))) - .getOrElse(table.output) - case _ => table.output + .getOrElse(physicalOutput) + case _ => physicalOutput } } diff --git a/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/commands/WriteIntoPaimonTable.scala b/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/commands/WriteIntoPaimonTable.scala index 937a47c526e9..06455864e0a1 100644 --- a/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/commands/WriteIntoPaimonTable.scala +++ b/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/commands/WriteIntoPaimonTable.scala @@ -24,6 +24,7 @@ import org.apache.paimon.options.Options import org.apache.paimon.spark._ import org.apache.paimon.spark.catalyst.analysis.ReplacePaimonFunctions import org.apache.paimon.spark.catalyst.analysis.expressions.ExpressionHelper +import org.apache.paimon.spark.schema.SparkSystemColumns import org.apache.paimon.spark.write.PaimonWriteOptions import org.apache.paimon.table.FileStoreTable @@ -47,10 +48,14 @@ case class WriteIntoPaimonTable( with Logging { override def run(sparkSession: SparkSession): Seq[Row] = { - val replacedData = + val replacedDataWithMetadata = PaimonUtils.createDataset( sparkSession, ReplacePaimonFunctions(sparkSession)(_data.queryExecution.analyzed)) + val replacedData = + SparkSystemColumns + .changelogMetadataFieldNames(table) + .foldLeft(replacedDataWithMetadata)((data, fieldName) => data.drop(fieldName)) mergeSchema(sparkSession, replacedData, options) val (dynamicPartitionOverwriteMode, overwritePartition) = parseSaveMode() diff --git a/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/read/BaseScan.scala b/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/read/BaseScan.scala index b7854ea4e297..14be066d4a9c 100644 --- a/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/read/BaseScan.scala +++ b/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/read/BaseScan.scala @@ -25,9 +25,10 @@ import org.apache.paimon.spark.{PaimonBatch, PaimonInputPartition, PaimonNumSpli import org.apache.paimon.spark.schema.PaimonMetadataColumn import org.apache.paimon.spark.schema.PaimonMetadataColumn._ import org.apache.paimon.spark.util.{OptionUtils, SplitUtils} -import org.apache.paimon.table.{SpecialFields, Table} +import org.apache.paimon.table.{FileStoreTable, SpecialFields, Table} import org.apache.paimon.table.BlobDescriptorReadUtils import org.apache.paimon.table.source.{ReadBuilder, Split} +import org.apache.paimon.table.system.ChangelogEventMetadataTable import org.apache.paimon.types.RowType import org.apache.spark.internal.Logging @@ -72,14 +73,19 @@ trait BaseScan extends Scan with SupportsReportStatistics with Logging { val coreOptions: CoreOptions = CoreOptions.fromMap(table.options()) lazy val tableRowType: RowType = { + var rowType = table.rowType() if ( coreOptions - .rowTrackingEnabled() && !table.rowType().containsField(SpecialFields.ROW_ID.name()) + .rowTrackingEnabled() && !rowType.containsField(SpecialFields.ROW_ID.name()) ) { - SpecialFields.rowTypeWithRowTracking(table.rowType()) - } else { - table.rowType() + rowType = SpecialFields.rowTypeWithRowTracking(rowType) + } + if (!coreOptions.changelogEventMetadataFields().isEmpty && table.isInstanceOf[FileStoreTable]) { + rowType = ChangelogEventMetadataTable.computeExtendedRowType( + table.asInstanceOf[FileStoreTable], + rowType) } + rowType } /** diff --git a/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/schema/SparkSystemColumns.scala b/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/schema/SparkSystemColumns.scala index 4e6d181f6e38..5abcde2ebdd6 100644 --- a/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/schema/SparkSystemColumns.scala +++ b/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/schema/SparkSystemColumns.scala @@ -18,8 +18,15 @@ package org.apache.paimon.spark.schema +import org.apache.paimon.CoreOptions +import org.apache.paimon.table.FileStoreTable +import org.apache.paimon.table.system.ChangelogEventMetadata + +import org.apache.spark.sql.catalyst.expressions.Attribute import org.apache.spark.sql.types.StructType +import scala.collection.JavaConverters._ + /** System columns for paimon spark. */ object SparkSystemColumns { @@ -34,4 +41,29 @@ object SparkSystemColumns { def filterSparkSystemColumns(schema: StructType): StructType = { StructType(schema.fields.filterNot(field => SPARK_SYSTEM_COLUMNS_NAME.contains(field.name))) } + + /** Names exposed by a table only for changelog reads, never as physical write columns. */ + def changelogMetadataFieldNames(table: FileStoreTable): Set[String] = { + val options = CoreOptions.fromMap(table.options()) + ChangelogEventMetadata + .extraValueFields(table.schema().logicalRowType(), options) + .asScala + .map(_.name()) + .toSet + } + + def filterSparkSystemColumns(schema: StructType, table: FileStoreTable): StructType = { + val metadataNames = changelogMetadataFieldNames(table) + StructType( + schema.fields.filterNot( + field => + SPARK_SYSTEM_COLUMNS_NAME.contains(field.name) || metadataNames.contains(field.name))) + } + + def filterChangelogMetadataColumns( + attributes: Seq[Attribute], + table: FileStoreTable): Seq[Attribute] = { + val metadataNames = changelogMetadataFieldNames(table) + attributes.filterNot(attribute => metadataNames.contains(attribute.name)) + } } diff --git a/paimon-spark/paimon-spark-ut/src/test/scala/org/apache/paimon/spark/PaimonCDCSourceTest.scala b/paimon-spark/paimon-spark-ut/src/test/scala/org/apache/paimon/spark/PaimonCDCSourceTest.scala index dd1aecfe4175..5f5c745dfd96 100644 --- a/paimon-spark/paimon-spark-ut/src/test/scala/org/apache/paimon/spark/PaimonCDCSourceTest.scala +++ b/paimon-spark/paimon-spark-ut/src/test/scala/org/apache/paimon/spark/PaimonCDCSourceTest.scala @@ -18,10 +18,15 @@ package org.apache.paimon.spark +import org.apache.paimon.data.GenericRow +import org.apache.paimon.disk.IOManagerImpl + import org.apache.spark.sql.{Dataset, Row} import org.apache.spark.sql.paimon.shims.memstream.MemoryStream import org.apache.spark.sql.streaming.StreamTest +import scala.collection.JavaConverters._ + class PaimonCDCSourceTest extends PaimonSparkTestBase with StreamTest { import testImplicits._ @@ -214,6 +219,111 @@ class PaimonCDCSourceTest extends PaimonSparkTestBase with StreamTest { } } + test("Paimon CDC Source: Spark reads exposed event metadata as a regular column") { + withTempDirs { + (checkpointDir, ioManagerDir) => + val tableName = "T" + spark.sql(s""" + |CREATE TABLE $tableName (id INT, data INT, event_ts BIGINT) + |TBLPROPERTIES ( + | 'primary-key'='id', + | 'bucket'='1', + | 'changelog-producer' = 'lookup', + | 'sequence.field' = 'event_ts', + | 'changelog-producer.event-metadata-fields' = 'event_ts', + | 'changelog-producer.metadata-field-prefix' = '__event__') + |""".stripMargin) + + val table = loadTable(tableName) + val ioManager = new IOManagerImpl(ioManagerDir.getCanonicalPath) + val write = table.newWrite(commitUser).withIOManager(ioManager) + val commit = table.newCommit(commitUser) + try { + // Write only the physical columns, as a Flink writer would. Spark exposes the + // generated metadata field when reading the table, but it is not an input column. + write.write(GenericRow.of(1, 10, 50L)) + commit.commit(0, write.prepareCommit(true, 0)) + + // Spark exposes the generated field as a regular column. No Flink metadata alias is + // needed to read it. + checkAnswer( + spark.sql(s"SELECT id, data, event_ts, __event__event_ts FROM $tableName"), + Row(1, 10, 50L, 50L) :: Nil) + + val location = table.location().toString + val readStream = spark.readStream + .format("paimon") + .option("read.changelog", "true") + .load(location) + .writeStream + .format("memory") + .option("checkpointLocation", checkpointDir.getCanonicalPath) + .queryName("mem_table") + .outputMode("append") + .start() + + val currentResult = () => spark.sql("SELECT * FROM mem_table") + try { + readStream.processAllAvailable() + checkAnswer(currentResult(), Row("+I", 1, 10, 50L, 50L) :: Nil) + + write.write(GenericRow.of(1, 20, 100L)) + commit.commit(1, write.prepareCommit(true, 1)) + readStream.processAllAvailable() + checkAnswer( + currentResult(), + Row("+I", 1, 10, 50L, 50L) :: + Row("-U", 1, 10, 50L, 100L) :: + Row("+U", 1, 20, 100L, 100L) :: Nil) + } finally { + readStream.stop() + } + } finally { + write.close() + commit.close() + ioManager.close() + } + } + } + + test("Paimon CDC Source: exposed event metadata is read-only for Spark writes") { + withTable("T") { + withSparkSQLConf("spark.paimon.write.merge-schema" -> "true") { + spark.sql(""" + |CREATE TABLE T (id INT, data INT, event_ts BIGINT) + |TBLPROPERTIES ( + | 'primary-key' = 'id', + | 'bucket' = '1', + | 'changelog-producer' = 'lookup', + | 'sequence.field' = 'event_ts', + | 'changelog-producer.event-metadata-fields' = 'event_ts') + |""".stripMargin) + + // Positional writes must not require the generated metadata field as an input column. + spark.sql("INSERT INTO T VALUES (1, 10, 50)") + + // A self-insert reads the generated field through SELECT *, so it also verifies that a + // readable metadata field is removed before write resolution and schema evolution. + spark.sql("INSERT INTO T SELECT * FROM T") + + if (gteqSpark3_5) { + // The physical extra column should still participate in merge-schema evolution, while + // the generated metadata field must not be committed as a physical column. + spark.sql( + "INSERT INTO T BY NAME " + + "SELECT 2 AS id, 20 AS data, 200L AS event_ts, 'extra' AS extra") + } + + val physicalFieldNames = loadTable("T").copyWithLatestSchema().schema().fieldNames().asScala + assert(!physicalFieldNames.contains("__internal__event_ts")) + if (gteqSpark3_5) { + assert(physicalFieldNames.contains("extra")) + } + assert(spark.table("T").schema.fieldNames.contains("__internal__event_ts")) + } + } + } + test("Paimon CDC Source: streaming read change-log with audit_log system table") { withTable("T") { withTempDir {