Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
45 changes: 42 additions & 3 deletions docs/docs/primary-key-table/changelog-producer.md
Original file line number Diff line number Diff line change
Expand Up @@ -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__<column>` 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
Expand Down
24 changes: 18 additions & 6 deletions docs/generated/core_configuration.html
Original file line number Diff line number Diff line change
Expand Up @@ -212,6 +212,12 @@
<td><p>Enum</p></td>
<td>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. <br /><br />Possible values:<ul><li>"none": No changelog file.</li><li>"input": Double write to a changelog file when flushing memory table, the changelog is from input.</li><li>"full-compaction": Generate changelog files with each full compaction.</li><li>"lookup": Generate changelog files through 'lookup' compaction.</li></ul></td>
</tr>
<tr>
<td><h5>changelog-producer.event-metadata-fields</h5></td>
<td style="word-wrap: break-word;">(none)</td>
<td>String</td>
<td>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.</td>
</tr>
<tr>
<td><h5>changelog-producer.ignore-delete</h5></td>
<td style="word-wrap: break-word;">false</td>
Expand All @@ -224,6 +230,12 @@
<td>Boolean</td>
<td>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.</td>
</tr>
<tr>
<td><h5>changelog-producer.metadata-field-prefix</h5></td>
<td style="word-wrap: break-word;">"__internal__"</td>
<td>String</td>
<td>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.</td>
</tr>
<tr>
<td><h5>changelog-producer.row-deduplicate</h5></td>
<td style="word-wrap: break-word;">false</td>
Expand Down Expand Up @@ -398,12 +410,6 @@
<td>MemorySize</td>
<td>When incremental size is bigger than this threshold, force a full compaction.</td>
</tr>
<tr>
<td><h5>continuous-compaction.initial-scan-mode</h5></td>
<td style="word-wrap: break-word;">earliest</td>
<td><p>Enum</p></td>
<td>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.<br /><br />Possible values:<ul><li>"earliest": Read snapshots from the earliest available snapshot.</li><li>"latest": Read the latest snapshot as the initial full baseline.</li></ul></td>
</tr>
<tr>
<td><h5>compaction.max-size-amplification-percent</h5></td>
<td style="word-wrap: break-word;">200</td>
Expand Down Expand Up @@ -494,6 +500,12 @@
<td><p>Enum</p></td>
<td>Specify the consumer consistency mode for table.<br /><br />Possible values:<ul><li>"exactly-once": Readers consume data at snapshot granularity, and strictly ensure that the snapshot-id recorded in the consumer is the snapshot-id + 1 that all readers have exactly consumed.</li><li>"at-least-once": Each reader consumes snapshots at a different rate, and the snapshot with the slowest consumption progress among all readers will be recorded in the consumer.</li></ul></td>
</tr>
<tr>
<td><h5>continuous-compaction.initial-scan-mode</h5></td>
<td style="word-wrap: break-word;">earliest</td>
<td><p>Enum</p></td>
<td>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.<br /><br />Possible values:<ul><li>"earliest": Read snapshots from the earliest available snapshot.</li><li>"latest": Read the latest snapshot as the initial full baseline.</li></ul></td>
</tr>
<tr>
<td><h5>continuous.discovery-interval</h5></td>
<td style="word-wrap: break-word;">10 s</td>
Expand Down
48 changes: 48 additions & 0 deletions paimon-api/src/main/java/org/apache/paimon/CoreOptions.java
Original file line number Diff line number Diff line change
Expand Up @@ -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<String> 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<String> 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<Boolean> TABLE_READ_SEQUENCE_NUMBER_ENABLED =
key("table-read.sequence-number.enabled")
.booleanType()
Expand Down Expand Up @@ -4033,6 +4067,20 @@ public List<String> changelogRowDeduplicateIgnoreFields() {
.orElse(Collections.emptyList());
}

public List<String> 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);
}
Expand Down
32 changes: 22 additions & 10 deletions paimon-core/src/main/java/org/apache/paimon/KeyValueFileStore.java
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -87,6 +88,7 @@ public KeyValueFileStore(
options.changelogRowDeduplicate()
? ValueEqualiserSupplier.fromIgnoreFields(valueType, logDedupIgnoreFields)
: () -> null;
ChangelogEventMetadata.validate(valueType, options);
}

@Override
Expand Down Expand Up @@ -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<org.apache.paimon.types.DataField> extraFields =
ChangelogEventMetadata.storageValueFields(valueType, options);
if (!extraFields.isEmpty()) {
builder.withChangelogExtraValueFields(extraFields);
}
}
return builder;
}

@Override
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand All @@ -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();
Expand Down Expand Up @@ -123,7 +125,9 @@ protected FileRecordReader<KeyValue> createRecordReader(
valueType,
file.level(),
overrideSequenceWithSnapshotId,
file.minSequenceNumber());
file.minSequenceNumber(),
metadataFallbackMapping,
!isChangelogFile(file));

if (deletionVector.isPresent() && !deletionVector.get().isEmpty()) {
return new ExposeDeletionKeyValueReader(reader, deletionVector.get());
Expand Down Expand Up @@ -165,7 +169,8 @@ public ChainKeyValueFileReaderFactory build(
dvFactory,
chainReadContext,
wrapped.options,
wrapped.readBatchSizer);
wrapped.readBatchSizer,
wrapped.createMetadataFallbackMapping());
}
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand All @@ -38,19 +39,26 @@ public class KeyValueDataFileRecordReader implements FileRecordReader<KeyValue>
private final int level;
private final boolean overrideSequenceWithSnapshotId;
private final long snapshotId;
@Nullable private final FallbackMappingRow metadataFallbackRow;

public KeyValueDataFileRecordReader(
FileRecordReader<InternalRow> reader,
RowType keyType,
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
Expand All @@ -67,6 +75,9 @@ public FileRecordIterator<KeyValue> 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
Expand Down
Loading
Loading