This topic explains how to use the Paimon connector for a streaming data lakehouse. For best results, we recommend using the connector with the Paimon Catalog.
Background information
Apache Paimon is a unified streaming and batch data lake storage format that supports high-throughput writes and low-latency queries. Paimon integrates well with popular compute engines available in Alibaba Cloud E-MapReduce, such as Flink, Spark, Hive, and Trino. You can use Apache Paimon to quickly build a data lake on HDFS or OSS and connect to these compute engines for data lake analytics. For more information, see Apache Paimon.
|
Category |
Description |
|
Supported types |
Source table, dimension table, result table, and target for data ingestion |
|
Running mode |
Streaming mode and batch mode |
|
Data format |
Not supported |
|
Monitoring metrics |
None |
|
API types |
SQL and YAML for data ingestion |
|
Updates and deletes on result tables |
Yes |
Key features
Apache Paimon provides the following core capabilities:
-
Build a lightweight, low-cost data lake on HDFS or object storage.
-
Read and write large-scale datasets in both streaming and batch modes.
-
Run batch and OLAP queries with data freshness from minutes down to seconds.
-
Ingest and produce incremental data, serving as a storage layer for both traditional offline and modern streaming data warehouses.
-
Pre-aggregate data to reduce storage costs and downstream compute load.
-
Access historical versions of data.
-
Filter data efficiently.
-
Enable schema evolution.
Limitations and recommendations
-
The Paimon connector requires Flink compute engine VVR 6.0.6 or later.
-
The following table lists the version compatibility between Paimon and VVR.
Apache Paimon version
VVR
1.3.1
11.5, 11.6, 11.7, 11.8
1.3
11.4
1.2
11.2, 11.3
1.1
11.1
1.0
8.0.11
-
Storage recommendations for concurrent writes
When multiple jobs perform concurrent writes to the same Paimon table, using standard OSS storage (oss://) can occasionally cause commit conflicts or job failures due to limitations on atomic file operations.
For stable and consistent writes, use a metadata or storage service that provides strong atomic guarantees. The preferred option is Data Lake Formation (DLF), which offers unified management of Paimon metadata and storage. Alternatively, you can use OSS-HDFS or HDFS.
-
How configuration changes take effect
Changes to the configuration parameters of a Paimon table take effect only after you restart the related jobs. Running jobs do not dynamically load these changes.
-
Delayed physical reclamation of dropped partitions
When you run a DROP PARTITION operation, the system does not immediately delete the underlying physical data files.
This operation performs a logical deletion. Paimon removes the metadata of the target partition only from the latest snapshot. Because Paimon supports the time travel feature, historical snapshots still reference that partition's data files. The physical data files are permanently deleted only after all historical snapshots that reference the partition reach their retention limit and are cleaned up by the snapshot expiration mechanism.
SQL
Use the Paimon connector in SQL jobs as a source table or a result table.
Syntax
-
If you create a Paimon table in a Paimon Catalog, you do not need to specify the
connectorparameter. The syntax is as follows:CREATE TABLE `<YOUR-PAIMON-CATALOG>`.`<YOUR-DB>`.paimon_table ( id BIGINT, data STRING, PRIMARY KEY (id) NOT ENFORCED ) WITH ( ... );NoteIf you have already created a Paimon table in a Paimon Catalog, you can use it directly.
-
If you create a Paimon temporary table in another Catalog, you must specify the 'connector' and 'path' parameters. The syntax is as follows:
CREATE TEMPORARY TABLE paimon_table ( id BIGINT, data STRING, PRIMARY KEY (id) NOT ENFORCED ) WITH ( 'connector' = 'paimon', 'path' = '<path-to-paimon-table-files>', 'auto-create' = 'true', -- If Paimon table data files do not exist at the specified path, they are automatically created. ... );Note-
Path example:
'path' = 'oss://<bucket>/test/order.db/orders'. Do not omit the.dbsuffix. Paimon relies on this suffix to identify the database. -
Multiple jobs that write to the same table must share the same path configuration.
-
If two path configurations are different, Paimon does not recognize them as the same table. Even if the physical path is the same, inconsistent Catalog configurations can lead to concurrent write conflicts, failed compact operations, and data loss. For example, Paimon considers
oss://b/testandoss://b/test/to be different tables due to the trailing slash, even though they may point to the same physical location.
-
WITH parameters
|
Parameter |
Description |
Type |
Required |
Default |
Remarks |
|
connector |
Specifies the connector for the table. |
String |
No |
None |
|
|
path |
The table's storage path. |
String |
No |
None |
|
|
auto-create |
Specifies whether to automatically create table files if they do not exist at the specified path. |
Boolean |
No |
false |
Valid values:
|
|
file.format |
The format of the data files. |
String |
No |
parquet |
Valid values:
|
|
bucket |
The number of buckets per partition. |
Integer |
No |
1 |
Paimon distributes data to buckets based on the Note
We recommend that each bucket contain less than 5 GB of data. |
|
bucket-key |
The column(s) used as the bucket key. |
String |
No |
None |
Specifies the columns used to distribute data into buckets. Separate multiple column names with a comma (,). For example, Note
|
|
changelog-producer |
The changelog production mechanism. |
String |
No |
none |
Paimon can produce a complete changelog for any input stream, which means that every
For more information about how to choose a changelog producer, see Changelog Production. |
|
full-compaction.delta-commits |
The maximum number of commits between two consecutive full compactions. |
Integer |
No |
None |
Specifies the maximum number of snapshot commits to allow before triggering a full compaction. |
|
lookup.cache-max-memory-size |
The memory cache size for a Paimon dimension table. |
String |
No |
256 MB |
This parameter controls the cache size for both dimension table lookups and the |
|
merge-engine |
The mechanism for merging records with the same primary key. |
String |
No |
deduplicate |
Valid values:
For a detailed analysis of merge engines, see Merge Engine. |
|
partial-update.ignore-delete |
Specifies whether to ignore delete (-D) messages. |
Boolean |
No |
false |
Valid values:
Note
|
|
ignore-delete |
Specifies whether to ignore delete (-D) messages. |
Boolean |
No |
false |
The valid values are the same as those for partial-update.ignore-delete. Note
|
|
partition.default-name |
The default partition name. |
String |
No |
__DEFAULT_PARTITION__ |
The partition name to use when a partition column's value is null or an empty string. |
|
partition.expiration-check-interval |
How often the system checks for expired partitions. |
String |
No |
1h |
For details, see How to configure automatic partition expiration. |
|
partition.expiration-time |
The duration after which partitions expire. |
String |
No |
None |
A partition expires when its age exceeds this value. By default, partitions never expire. The system calculates a partition's age from its partition value. For details, see How to configure automatic partition expiration. |
|
partition.timestamp-formatter |
The format string for converting a time string to a timestamp. |
String |
No |
None |
Specifies the format for extracting the partition age from the partition value. For details, see How to configure automatic partition expiration. |
|
partition.timestamp-pattern |
The format string for converting a partition value to a time string. |
String |
No |
None |
Specifies the pattern for extracting a time string from the partition value. For details, see How to configure automatic partition expiration. |
|
scan.bounded.watermark |
The watermark value that signals the end of a scan. The source table stops producing data when its watermark exceeds this value. |
Long |
No |
None |
N/A |
|
scan.mode |
Specifies the consumption position for the Paimon source table. |
String |
No |
default |
For details, see How to set the consumption position for a Paimon source table. |
|
scan.snapshot-id |
Specifies the snapshot from which the Paimon source table starts consumption. |
Integer |
No |
None |
For details, see How to set the consumption position for a Paimon source table. |
|
scan.timestamp-millis |
Specifies the point in time from which the Paimon source table starts consumption. |
Integer |
No |
None |
For details, see How to set the consumption position for a Paimon source table. |
|
snapshot.num-retained.max |
The maximum number of recent snapshots to retain. |
Integer |
No |
2147483647 |
Snapshot expiration is triggered if either this condition or the |
|
snapshot.num-retained.min |
The minimum number of recent snapshots to retain. |
Integer |
No |
10 |
N/A |
|
snapshot.time-retained |
The retention period for snapshots. |
String |
No |
1h |
Snapshot expiration is triggered if either this condition or the |
|
write-mode |
The write mode for the Paimon table. |
String |
No |
change-log |
Valid values:
For more information about write modes, see Write Mode. |
|
scan.infer-parallelism |
Specifies whether to automatically infer the parallelism for a Paimon source table. |
Boolean |
No |
true |
Valid values:
|
|
scan.parallelism |
The parallelism for the Paimon source table. |
Integer |
No |
None |
Note
This parameter is ignored if Resource Mode is set to Mode on the job's tab. |
|
sink.parallelism |
The parallelism for the Paimon sink table. |
Integer |
No |
None |
Note
This parameter is ignored if Resource Mode is set to Mode on the job's tab. |
|
sink.clustering.by-columns |
Specifies the clustering columns for writing to the Paimon sink table. |
String |
No |
None |
For Paimon append-only tables (tables without a primary key), this parameter enables clustered writing in batch jobs. This process improves query performance by grouping data on the specified columns. Separate multiple column names with a comma (,), for example, For more information about clustering, see the Apache Paimon Official Documentation. |
|
sink.delete-strategy |
Specifies a validation strategy to ensure that the system correctly handles retraction messages (-D/-U). |
Enum |
No |
NONE |
Valid values and the expected behavior of the sink operator when handling retraction messages:
Note
|
|
blob-as-descriptor |
Specifies whether to output Blob descriptor bytes when a Blob column is read. |
Boolean |
No |
false |
When this parameter is set to true, queries return the serialized bytes of the BlobDescriptor instead of the actual Blob content. Use this parameter together with the Blob presigned URL function. You do not need to set this parameter at table creation time. You can dynamically set it at read time by using an SQL hint. This parameter is supported in VVR 11.9-preview1 and later. Note
Supported only in VVR 11.9-preview1 and later. |
For more information about configuration options, see the Apache Paimon Official Documentation.
Vector table parameters
The following parameters are supported only in VVR 11.8 and later.
|
Parameter |
Description |
Data type |
Required |
Default value |
Remarks |
|
|
The number of IVF clusters to probe during search. |
Integer |
No |
16 |
Higher values generally improve recall but increase latency. |
|
|
Retrieves top_k × refine_factor IVF candidates and re-ranks them using the original vectors stored in the Paimon table. |
Float |
No |
Disabled |
Disabled by default for all IVF variants. Best suited for compressed indexes (such as ivf-pq and ivf-hnsw-sq) in scenarios where recall is prioritized over latency. |
|
|
The HNSW search width during search. |
Integer |
No |
0 |
Higher values generally improve recall but increase latency. 0 indicates that the native library default is used. |
|
|
The Lumina DiskANN search list size. |
Integer |
No |
max(1.5 × top_k, 16) |
Higher values generally improve recall but increase latency. |
|
|
The Lumina DiskANN search beam width. |
Integer |
No |
4 |
— |
|
|
The number of parallel Lumina searches. |
Integer |
No |
5 |
— |
Feature details
Data freshness and consistency
The Paimon sink table uses the two-phase commit protocol to commit data during each Flink job checkpoint. Therefore, data freshness is determined by the Flink job's checkpoint interval. Each commit produces up to two snapshots.
When two Flink jobs write to the same Paimon table concurrently, if the jobs write to different buckets, they achieve serializable consistency. If the jobs write to the same bucket, they only achieve snapshot isolation. This means the table's data might be a mix of results from both jobs, but no data loss occurs.
Merge engine
When a Paimon sink table receives multiple records with the same primary key, it merges them into a single record to maintain uniqueness. You can control this behavior by setting the merge-engine parameter. The following table describes the available merge engines.
|
Merge engine |
Description |
|
Deduplicate |
The deduplication engine is the default. For multiple records with the same primary key, the Paimon sink table keeps only the latest record and discards the others. Note
If the latest record is a delete message, all records with that primary key are discarded. |
|
Partial Update |
The partial update engine allows you to build a complete record by incrementally updating it with multiple messages. When a new record with the same primary key arrives, its non-null values overwrite the corresponding fields in the existing record. The engine ignores fields that are null in the new record and retains the existing values. For example, assume a Paimon sink table receives the following three records in order:
If the first column is the primary key, the final merged record is <1, 25.2, 10, 'This is a book'>. Note
|
|
Aggregation |
In some use cases, you might only need the aggregated value of records. The aggregation engine combines records that have the same primary key by using the aggregate functions you specify. For each non-primary key column, you must specify an aggregate function using the
The
Note
|
Changelog producer
Set the changelog-producer parameter to configure Paimon to generate a complete changelog (where every update_after record has a corresponding update_before record) for any input stream. The following table describes the available changelog producers. For more details, see the Apache Paimon official documentation.
|
Producer |
Description |
|
None |
When you set For example, if a downstream consumer needs to calculate a column's sum and only sees the latest value of 5, it cannot determine how to update the total. If the previous value was 4, the sum should increase by 1; if the previous value was 6, the sum should decrease by 1. Consumers that are sensitive to Note
If your downstream consumer, such as a database, is not sensitive to |
|
Input |
When you set Use this producer only when the input stream itself is already a complete changelog, such as data from Change Data Capture (CDC). |
|
Lookup |
When you set Compared to the Use this option for use cases that require high data freshness (for example, minute-level). |
|
Full Compaction |
When you set Compared to the Use this option for use cases with low data freshness requirements (for example, hourly). |
Write mode
Paimon tables support the following write modes.
|
Mode |
Description |
|
Change-log |
The |
|
Append-only |
The For a detailed description of the
|
Target for Flink CDC data ingestion
Paimon tables support real-time synchronization of data from single tables or entire databases. Upstream schema changes are synchronized to Paimon tables in real time. For more information, see the Data ingestion section.
Variant read pruning
Reads only the Variant fields referenced by the query, reducing I/O and memory overhead.
How to enable
|
Parameter |
Default value |
Description |
|
|
false |
Set to |
Prerequisites
-
Variant data can only be written by VVR 11.6 or later. Variant data written by earlier versions or other engines does not support read pruning.
-
For upstream data to take effect, you must enable the
variant.inferShreddingSchemaoption to enable automatic Variant shredding inference, or explicitly configurevariant.shreddingSchemaand related options to specify which Variant members are suitable for shredding. For more information, see Coreoptions.
Supported scenarios
Read pruning takes effect when SQL accesses Variant fields by string key, for example:
SELECT v['a'] FROM t; -- Single level
SELECT v['a']['b'] FROM t; -- Nested
SELECT v['a'], v['b'] FROM t; -- Multiple fields
SELECT id, v['a'] FROM t; -- Mixed with regular columns
SELECT v['a'] + 1 FROM t; -- Referenced in expressions
Unsupported scenarios
-
Referencing the entire Variant directly, such as
SELECT v FROM t. -
Combining direct reference with field access, such as
SELECT v, v['a'] FROM t. -
Array index access, such as
v[0]. -
The table schema contains nested Row fields. In this case, read pruning does not apply to any Variant field in the table.
Blob presigned URLs
Paimon can generate short-lived public HTTPS GET presigned URLs for Blob column data stored in OSS. External services, such as multimodal model services, can use the URLs to retrieve the file content directly, without transferring the file bytes inside a Flink job.
Limits
-
Only Flink compute engine VVR 11.9-preview1 and later support this feature.
-
Only Paimon tables stored in OSS are supported. The metastore type does not matter. Filesystem, DLF, and other metastore types are all supported.
-
The generated URL is a public HTTPS URL and must be consumed by an external service that has public network access. If the catalog uses an internal endpoint, we recommend that you change it to a public HTTPS endpoint so that the external service can access the generated URL.
-
To read a regular Blob column, you must first obtain the BlobDescriptor bytes by using the
blob-as-descriptoroption, and then pass the bytes to the function.
Syntax
descriptor_to_presigned_url(source_table, descriptor, validity)
try_descriptor_to_presigned_url(source_table, descriptor, validity)
Parameters
|
Parameter |
Type |
Description |
|
|
STRING literal |
The Paimon table that contains the Blob column, in the format |
|
|
BYTES |
The BlobDescriptor bytes, which are the result of reading the Blob column after |
|
|
INTERVAL |
The validity period of the URL. The value must be a positive number of seconds, for example, |
Return value
|
Type |
Description |
|
STRING |
A short-lived public HTTPS GET presigned URL. |
Example
-- Use an SQL hint to read the Blob column as a descriptor and generate a presigned URL
SELECT sys.descriptor_to_presigned_url(
'default.image_table',
image,
INTERVAL '5' MINUTE)
FROM image_table /*+ OPTIONS('blob-as-descriptor'='true') */;
-- Fault-tolerant mode: returns NULL on row-level errors, with other behavior unchanged
SELECT sys.try_descriptor_to_presigned_url(
'default.image_table',
image,
INTERVAL '5' MINUTE)
FROM image_table /*+ OPTIONS('blob-as-descriptor'='true') */;
-
descriptor_to_presigned_urlthrows an exception on a row-level error.try_descriptor_to_presigned_urlreturns NULL on a row-level error. -
The generated URL does not contain a file extension, and the bytes that the URL points to are identical to the Blob data. Downstream services must identify the file format based on the returned bytes instead of the URL suffix.
-
A presigned URL is equivalent to temporary access credentials. After a presigned URL is generated, submit it to the downstream service immediately. Do not write it to logs or result tables, and do not persist it.
Chain table
A chain table reconstructs a full data set from a snapshot branch and the changes that follow it in a delta branch. Write periodic full data to the snapshot branch and the changes between two full loads to the delta branch. A read uses the latest snapshot partition as an anchor and merges later delta partitions to rebuild the full data, so you no longer need to rewrite the entire table. This is useful when only a small fraction of rows changes between runs, such as an ODS binlog dump. For more information, see Chain Table in the Apache Paimon documentation. Chain tables are supported only in VVR 11.9 and later.
Enable a chain table
Set chain-table.enabled to true when you create the table, create the snapshot and delta branches, and then configure the branch fallback options on the primary table and on both branch tables.
CREATE TABLE `default`.`t` (
`t1` STRING,
`t2` STRING,
`t3` STRING,
`date` STRING,
PRIMARY KEY (`date`, `t1`) NOT ENFORCED
) PARTITIONED BY (`date`) WITH (
'chain-table.enabled' = 'true',
'sequence.field' = 't2',
'bucket-key' = 't1',
'bucket' = '2',
'partition.timestamp-pattern' = '$date',
'partition.timestamp-formatter' = 'yyyyMMdd'
);
-- mycat is the name of the Paimon catalog.
CALL `mycat`.sys.create_branch('default.t', 'snapshot');
CALL `mycat`.sys.create_branch('default.t', 'delta');
ALTER TABLE `default`.`t` SET (
'scan.fallback-snapshot-branch' = 'snapshot',
'scan.fallback-delta-branch' = 'delta');
ALTER TABLE `default`.`t$branch_snapshot` SET (
'scan.fallback-snapshot-branch' = 'snapshot',
'scan.fallback-delta-branch' = 'delta');
ALTER TABLE `default`.`t$branch_delta` SET (
'scan.fallback-snapshot-branch' = 'snapshot',
'scan.fallback-delta-branch' = 'delta');
Write data
Write full data to t$branch_snapshot. Each partition of this branch is a complete view at a point in time. Write incremental data to t$branch_delta. Each partition of this branch contains only the changes of that period.
-- Full write
INSERT OVERWRITE `default`.`t$branch_snapshot` PARTITION (date = '20250810') VALUES ('1', '1', '1');
-- Incremental write
INSERT INTO `default`.`t$branch_delta` PARTITION (date = '20250811') VALUES ('2', '1', '1');
Read data
|
Read path |
Operation |
Result |
|
Full query |
Query a partition of the primary table. |
The full data of the partition. The partition is read directly if the snapshot branch contains it. Otherwise, it is rebuilt from an anchor snapshot and the later delta partitions. |
|
Incremental query |
Query a partition of |
Only the changes stored in the selected delta partition. The full state is not rebuilt. |
|
Hybrid query |
Combine a full query and an incremental query with |
The combination of full and incremental results. The same row may appear twice. |
|
Streaming read |
Read the primary table in streaming mode. |
The result of the initial load phase is emitted first, then new commits of the delta branch are continuously consumed. |
|
Lookup join |
Use the chain table as a dimension table. |
The merged state of the latest snapshot partition and the later delta partitions of each group, refreshed incrementally based on |
-- Full query
SELECT t1, t2, t3 FROM `default`.`t` WHERE date = '20250811';
-- Incremental query
SELECT t1, t2, t3 FROM `default`.`t$branch_delta` WHERE date = '20250811';
-- Hybrid query
SELECT t1, t2, t3 FROM `default`.`t` WHERE date = '20250811'
UNION ALL
SELECT t1, t2, t3 FROM `default`.`t$branch_delta` WHERE date = '20250811';
Streaming read
A streaming read job runs in two phases:
-
Initial load phase: by default, the job reads the latest snapshot partition of each group and the delta partitions that come after it. Older snapshot partitions are excluded. To merge snapshot and delta data by anchor in this phase and resolve cross-branch DELETE records, set
chain-table.streaming.merge-snapshottotrue. The trade-off is a heavier scan at job startup. -
Incremental phase: the job continuously reads new commits from the
deltabranch.
SET 'execution.runtime-mode' = 'streaming';
INSERT INTO downstream_sink SELECT * FROM `default`.`t`;
The incremental phase monitors only the delta branch. Writes to the snapshot branch are not detected until the job is restarted. Run the compact_chain_table procedure on a regular basis to merge incremental data of a partition into the snapshot branch. After compaction, only the delta changes that arrived after compaction need to be merged in the initial load phase.
Multiple partition keys
If the table has more than one partition key, use chain-table.chain-partition-keys to separate the chain dimension from the group dimension. The value must be a contiguous suffix of the partition keys. The suffix fields form the time-ordered chain, and the partition fields before them become the group dimension. Each group of group partition values maintains its own chain and anchor. If this option is not set, all partitions belong to a single implicit group.
CREATE TABLE `default`.`t` (
`t1` STRING,
`t2` STRING,
`t3` STRING,
`region` STRING,
`date` STRING,
PRIMARY KEY (`region`, `date`, `t1`) NOT ENFORCED
) PARTITIONED BY (`region`, `date`) WITH (
'chain-table.enabled' = 'true',
'sequence.field' = 't2',
'bucket-key' = 't1',
'bucket' = '2',
'partition.timestamp-pattern' = '$date',
'partition.timestamp-formatter' = 'yyyyMMdd',
-- Only date forms the chain. Each region has its own chain.
'chain-table.chain-partition-keys' = 'date'
);
Limitations
-
Chain tables are supported only for primary key tables. You must specify
bucketandbucket-key. -
The schema of each branch must be consistent.
-
A chain table requires
sequence.field. Thefirst-rowandaggregationmerge engines are not supported. -
Streaming read and lookup join require the
deltabranch to use thededuplicatemerge engine (default value). Batch read is not affected. -
Deletion vectors are supported only for chain tables that use the
deduplicatemerge engine. -
Only
none(default value) andinputare supported forchangelog-producer. Thelookupandfull-compactionmodes are not supported. If the value isnone, records are normalized by the full primary key that includes the chain partition, and-Dand-Urecords whose original records reside in another partition are dropped. Set the value toinputif downstream consumers must receive cross-partition changelog records. -
Streaming read and lookup join support only the default startup mode (
latest-full). An error is thrown if you specify a starting position such asscan.snapshot-id,scan.timestamp-millis,scan.modeset tolatest, orconsumer-id. To perform a standard streaming read without chain table logic, read a specific branch table such ast$branch_delta. -
Partition filters are not supported in chain table streaming reads and lookup joins. This includes partition columns in the
WHEREclause and thescan.partitionsoption, because a filter interferes with the cross-branch merge logic. Use batch mode to read a specific partition. -
The join key of a lookup join must not contain partition keys. Otherwise, an error is returned.
-
partition.timestamp-formattermust specify a complete date (year-month-day or year-day-of-year), use only round-trippable date and time fields (y/u,M/L,d,D,H,k,m, ands), and have a minimum step of at least one second.
Mosaic file format
Mosaic is a columnar-bucket hybrid format optimized for wide tables that contain thousands or tens of thousands of columns. Columns are sorted by name and distributed into buckets by range, stored in column-oriented format within each bucket, and compressed independently by using ZSTD. Projection pushdown works at bucket granularity, so a query reads only the buckets that contain the requested columns. Columns with similar name prefixes land in the same bucket, which improves both the compression ratio and the locality of projections. For more information, see Mosaic. The Mosaic format is supported only in VVR 11.9 and later.
Enable the format
Set file.format to mosaic when you create the table. In a data ingestion job, the corresponding option is table.properties.file.format.
CREATE TABLE `default`.`wide_table` (
`id` BIGINT,
`col_a` STRING,
`col_b` DOUBLE
) WITH (
'file.format' = 'mosaic'
);
Format options
|
Option |
Default value |
Description |
|
|
auto |
The number of column buckets for parallel I/O. If the value is 0 or the option is not specified, the number of buckets is automatically determined. |
|
|
(empty) |
The names of the columns for which min/max statistics are collected for predicate pushdown, separated by commas (,). No statistics are collected if the value is empty. |
|
|
8 |
The number of row groups that a reader opens ahead of the row group being consumed. Prefetching overlaps range read latency with decoding. Each prefetched row group keeps its decoded batch in memory and uses a separate input stream. The value 0 disables prefetching. |
|
|
64 mb |
The upper bound of the estimated decoded size of the prefetched row groups, which is calculated from their row counts and the projected column types. A wide projection or a large row group lowers the effective value of |
Limitations
Mosaic does not support complex types: ARRAY, MAP, MULTISET, ROW, VARIANT, BLOB, and VECTOR.
Use Paimon as a dimension table
Paimon tables can be used as dimension tables. For the JOIN syntax, see Dimension table JOIN statements.
By default, lookup loads all data in each parallel instance. This approach is suitable only for small dimension tables. For large dimension tables, use the Shuffle Lookup solutions described below.
Partitioned dimension tables
If your dimension table is partitioned and you only need data from the latest one or two partitions, you can use the dynamic partition loading feature:
SELECT * FROM T
JOIN DIM /*+ OPTIONS('lookup.dynamic-partition'='max_pt()', 'lookup.dynamic-partition.refresh-interval'='1 h') */
FOR SYSTEM_TIME AS OF T.proc_time AS D
ON T.col = D.col;
|
Parameter |
Data type |
Default value |
Description |
|
lookup.dynamic-partition |
String |
N/A |
|
|
lookup.dynamic-partition.refresh-interval |
Duration |
1 h |
The interval at which the system checks for partition updates in the dimension table. |
Large dimension tables: fixed-bucket tables
Supported only in VVR 8.0.8 and later. For fixed-bucket tables (bucket > 0), you can use Shuffle Lookup to distribute data by bucket key across parallel instances, so each instance loads only the data in its assigned buckets:
SELECT /*+ LOOKUP('table'='D', 'shuffle'='true') */ T.col1, D.col2
FROM T
JOIN DIM FOR SYSTEM_TIME AS OF T.proc_time AS D
ON T.col1 = D.col1;
-
The join key must be the bucket key. The bucket key defaults to the primary key.
-
Only fixed-bucket tables (bucket > 0) support this feature.
Large dimension tables: non-fixed-bucket tables
Supported only in VVR 8.0.10 and later. For dynamic-bucket tables or append tables, you can use SHUFFLE_HASH or REPLICATED_SHUFFLE_HASH so that each parallel instance reads all data but retains only the portion it needs:
-- Shuffle Hash
SELECT /*+ SHUFFLE_HASH(D) */ T.col1, D.col2
FROM T
JOIN DIM FOR SYSTEM_TIME AS OF T.proc_time AS D
ON T.col1 = D.col1;
-- Replicated Shuffle Hash
SELECT /*+ REPLICATED_SHUFFLE_HASH(D) */ T.col1, D.col2
FROM T
JOIN DIM FOR SYSTEM_TIME AS OF T.proc_time AS D
ON T.col1 = D.col1;
For more information about SHUFFLE_HASH and REPLICATED_SHUFFLE_HASH, see Dimension table JOIN statements.
Data ingestion
You can use the Paimon connector as a sink in data ingestion YAML jobs.
Syntax
sink:
type: paimon
name: Paimon Sink
catalog.properties.metastore: filesystem
catalog.properties.warehouse: /path/warehouse
Parameters
|
Parameter |
Description |
Required |
Type |
Default |
Notes |
|
type |
The connector type. |
Yes |
STRING |
None |
The value must be |
|
name |
The name of the sink. |
No |
STRING |
None |
|
|
catalog.properties.metastore |
The type of the Paimon catalog. |
No |
STRING |
filesystem |
Valid values:
|
|
catalog.properties.* |
Parameters for creating a Paimon catalog. |
No |
STRING |
None |
For more information, see Manage Paimon Catalog. |
|
table.properties.* |
Parameters for creating a Paimon table. |
No |
STRING |
None |
For more information, see Paimon table options. |
|
table.properties.file.format |
The storage format of data files in the table. |
String |
No |
parquet |
Valid values:
|
|
catalog.properties.warehouse |
The root directory for file storage. |
No |
STRING |
None |
This parameter applies only when |
|
commit.user-prefix |
The username prefix for committing data files. |
No |
STRING |
None |
Note
We recommend setting different usernames for different jobs. This makes it easier to identify the job that causes a commit conflict. |
|
partition.key |
The partition keys for a partitioned table. |
No |
STRING |
None |
Different tables are separated by |
|
sink.cross-partition-upsert.tables |
Lists tables that require a cross-partition upsert, where the primary key does not include all partition keys. |
No |
STRING |
None |
Applies to tables with cross-partition updates.
Important
|
|
sink.commit.parallelism |
Specifies the parallelism of the Commit operator. |
No |
INTEGER |
None |
If the Commit operator is a bottleneck, use this parameter to increase its parallelism and improve performance. This parameter is supported only in Realtime Compute for Apache Flink 11.6 and later. Note
Setting this parameter changes the operator parallelism. When restarting a stateful job, you must specify |
|
sink.existing-table.schema-check.enabled |
Specifies whether to check the compatibility of the existing table schema and automatically adapt the schema when data is written to a pre-created table. |
Boolean |
No |
true |
The following list describes the valid values and the behavior when data is written to a pre-created table:
Note
This parameter is supported only in Realtime Compute for Apache Flink VVR 11.8 and later. |
Reuse an existing catalog
Starting with Realtime Compute for Apache Flink 11.5, you can directly reference a built-in Paimon catalog from the Data Management page in a Flink CDC data ingestion job. This reduces manual configuration.
sink:
type: paimon
using.built-in-catalog: paimon_dlf_catalog
catalog.properties.fs.oss.endpoint: oss-cn-beijing-internal.aliyuncs.com
Data ingestion jobs can automatically reuse all parameters in a Paimon catalog. This is equivalent to manually configuring parameters prefixed with catalog.properties. in a YAML job.
To override an automatically reused parameter, explicitly set it in the YAML job. The explicit YAML configuration has a higher priority. For example, in the preceding sample, the fs.oss.endpoint parameter uses the value from the YAML job, overriding the one in paimon_dlf_catalog.
Schema Changes
A CDC YAML pipeline job can use different strategies to handle schema changes. To configure the strategy, set the pipeline-level parameter schema.change.behavior as described in Schema evolution configurations. schema.change.behavior can be set to IGNORE, LENIENT, TRY_EVOLVE, EVOLVE, or EXCEPTION. The following sections describe how the LENIENT and EVOLVE modes, which involve schema changes, handle different schema change events.
LENIENT (default)
The following schema changes are supported in LENIENT mode:
-
Add a nullable column: the column is automatically appended to the schema of the result table, and data of the new column is automatically synchronized.
-
Drop a nullable column: the column is not dropped from the result table. Instead, the data of the column is automatically replaced with NULL values.
-
Add a NOT NULL column: the column is automatically appended to the schema of the result table, and data of the new column is automatically synchronized. The new column is nullable by default, and the data generated before the column is added is set to NULL.
-
Rename a column: the operation is treated as adding a column and dropping a column. The renamed column is appended to the end of the result table, and the data of the original column is automatically replaced with NULL values. For example, if col_a is renamed as col_b, col_b is appended to the end of the result table, and the data of col_a is replaced with NULL values.
-
Change a column type: changes the data type of the corresponding column in the result table.
-
The following schema changes are not supported:
-
Changes to primary key types.
-
Changes from NULLABLE to NOT NULL.
-
EVOLVE
The following schema changes are supported in EVOLVE mode:
-
Add a nullable column: supported.
-
Drop a nullable column: not supported.
-
Add a NOT NULL column: a nullable column is added to the result table.
-
Rename a column: supported. The original column is renamed in the result table.
-
Change a column type: changes the data type of the corresponding column in the result table.
-
The following schema changes are not supported:
-
Changes to primary key types.
-
Changes from NULLABLE to NOT NULL.
-
For an example of how to enable the EVOLVE mode, see Enabling the EVOLVE mode.
If the downstream Paimon table already exists, the job writes data to the existing table schema and does not attempt to create the table again. If the existing table schema differs from the schema of the upstream table, the nullable columns that exist in the upstream table but are missing in the downstream table are automatically added to the existing table. If a case cannot be automatically adapted, such as an incompatible column type or a missing or redundant NOT NULL column, the check fails and an error is reported. You can run the SQL statements that are recommended in the error message to modify the table schema. You can also set sink.existing-table.schema-check.enabled to false in the sink to skip the check and automatic column addition, so that data is written based on the schema of the downstream table.
Examples
The following examples show configurations in typical scenarios.
Write to a Rest Catalog
The following example shows how to write data to Data Lake Formation (DLF) when the Paimon catalog is of the rest type:
source:
type: mysql
name: MySQL Source
hostname: ${secret_values.mysql.hostname}
port: ${mysql.port}
username: ${secret_values.mysql.username}
password: ${secret_values.mysql.password}
tables: ${mysql.source.table}
server-id: 8601-8604
#(Optional) Synchronize data from newly added tables during the incremental phase.
scan.binlog.newly-added-table.enabled: true
#(Optional) Synchronize table and field comments.
include-comments.enabled: true
#(Optional) Prioritize the distribution of unbounded shards to prevent potential Out Of Memory (OOM) issues on TaskManagers.
scan.incremental.snapshot.unbounded-chunk-first.enabled: true
#(Optional) Enable parsing and filtering to accelerate data reads.
scan.only.deserialize.captured.tables.changelog.enabled: true
sink:
type: paimon
name: Paimon Sink
catalog.properties.metastore: rest
catalog.properties.uri: dlf_uri
catalog.properties.warehouse: your_warehouse
catalog.properties.token.provider: dlf
#(Optional) Specify a commit user. We recommend that you use a different commit user for each job to avoid conflicts.
commit.user: your_job_name
#(Optional) Enable deletion vectors to improve read performance.
table.properties.deletion-vectors.enabled: true
For information about the parameters prefixed with catalog.properties, see Flink CDC Catalog configuration parameters.
Write to a FileSystem Catalog
Example configuration for writing to Object Storage Service (OSS) with a filesystem Paimon catalog:
source:
type: mysql
name: MySQL Source
hostname: ${secret_values.mysql.hostname}
port: ${mysql.port}
username: ${secret_values.mysql.username}
password: ${secret_values.mysql.password}
tables: ${mysql.source.table}
server-id: 8601-8604
#(Optional) Synchronize data from newly added tables during the incremental phase.
scan.binlog.newly-added-table.enabled: true
#(Optional) Synchronize table and field comments.
include-comments.enabled: true
#(Optional) Prioritize the distribution of unbounded shards to prevent potential Out Of Memory (OOM) issues on TaskManagers.
scan.incremental.snapshot.unbounded-chunk-first.enabled: true
#(Optional) Enable parsing and filtering to accelerate data reads.
scan.only.deserialize.captured.tables.changelog.enabled: true
sink:
type: paimon
name: Paimon Sink
catalog.properties.metastore: filesystem
catalog.properties.warehouse: oss://default/test
catalog.properties.fs.oss.endpoint: oss-cn-beijing-internal.aliyuncs.com
catalog.properties.fs.oss.accessKeyId: xxxxxxxx
catalog.properties.fs.oss.accessKeySecret: xxxxxxxx
#(Optional) Specify a commit user. We recommend that you use a different commit user for each job to avoid conflicts.
commit.user: your_job_name
#(Optional) Enable deletion vectors to improve read performance.
table.properties.deletion-vectors.enabled: true
For information about the parameters prefixed with catalog.properties, see Create a Paimon Filesystem Catalog.
Write to a partitioned table
Converts the create_time field of the TIMESTAMP type in the upstream table into a DATE field, which is used as the partition key of the Paimon table.
source:
type: mysql
name: MySQL Source
hostname: <yourHostname>
port: 3306
username: flink
password: ${secret_values.password}
tables: test_db.test_source_table
server-id: 5401-5499
#(Optional) Synchronize data from newly added tables during the incremental phase.
scan.binlog.newly-added-table.enabled: true
#(Optional) Synchronize table and field comments.
include-comments.enabled: true
#(Optional) Prioritize the distribution of unbounded shards to prevent potential Out Of Memory (OOM) issues on TaskManagers.
scan.incremental.snapshot.unbounded-chunk-first.enabled: true
#(Optional) Enable parsing and filtering to accelerate data reads.
scan.only.deserialize.captured.tables.changelog.enabled: true
sink:
type: paimon
name: Paimon Sink
using.built-in-catalog: paimon_dlf_catalog
#(Optional) Specify a commit user. We recommend that you use a different commit user for each job to avoid conflicts.
commit.user: your_job_name
#(Optional) Enable deletion vectors to improve read performance.
table.properties.deletion-vectors.enabled: true
transform:
- source-table: test_db.test_source_table
projection: \*, DATE_FORMAT(CAST(create_time AS TIMESTAMP), 'yyyy-MM-dd') as partition_key
primary-keys: id, create_time, partition_key
partition-keys: partition_key
description: add partition key
pipeline:
name: MySQL to Paimon Pipeline
Synchronize a single table
We recommend that you configure single-table synchronization jobs for tables that have high requirements for real-time performance and stability. Single-table synchronization avoids the complexity of multi-table synchronization. The following example shows the configuration:
source:
type: mysql
name: MySQL Source
hostname: <yourHostname>
port: 3306
username: flink
password: ${secret_values.password}
tables: test_db.test_source_table
server-id: 5401-5499
#(Optional) Synchronize table and field comments.
include-comments.enabled: true
#(Optional) Prioritize the distribution of unbounded shards to prevent potential Out Of Memory (OOM) issues on TaskManagers.
scan.incremental.snapshot.unbounded-chunk-first.enabled: true
#(Optional) Enable parsing and filtering to accelerate data reads.
scan.only.deserialize.captured.tables.changelog.enabled: true
sink:
type: paimon
name: Paimon Sink
using.built-in-catalog: paimon_dlf_catalog
#(Optional) Specify a commit user. We recommend that you use a different commit user for each job to avoid conflicts.
commit.user: your_job_name
#(Optional) Enable deletion vectors to improve read performance.
table.properties.deletion-vectors.enabled: true
pipeline:
name: MySQL to Paimon Pipeline
Synchronize an entire database
Synchronizing an entire database reduces the number of jobs that you need to configure and run, which lowers O&M costs. The following example shows the configuration:
source:
type: mysql
name: MySQL Source
hostname: <yourHostname>
port: 3306
username: flink
password: ${secret_values.password}
tables: test_db.\.*
server-id: 5401-5499
#(Optional) Synchronize data from newly added tables during the incremental phase.
scan.binlog.newly-added-table.enabled: true
#(Optional) Synchronize table and field comments.
include-comments.enabled: true
#(Optional) Prioritize the distribution of unbounded shards to prevent potential Out Of Memory (OOM) issues on TaskManagers.
scan.incremental.snapshot.unbounded-chunk-first.enabled: true
#(Optional) Enable parsing and filtering to accelerate data reads.
scan.only.deserialize.captured.tables.changelog.enabled: true
sink:
type: paimon
name: Paimon Sink
using.built-in-catalog: paimon_dlf_catalog
#(Optional) Specify a commit user. We recommend that you use a different commit user for each job to avoid conflicts.
commit.user: your_job_name
#(Optional) Enable deletion vectors to improve read performance.
table.properties.deletion-vectors.enabled: true
pipeline:
name: MySQL to Paimon Pipeline
Synchronize an entire database and replace database and table names
To replace database and table names in a batch, add the route section to an entire-database synchronization job. The following example shows the configuration:
source:
type: mysql
name: MySQL Source
hostname: <yourHostname>
port: 3306
username: flink
password: ${secret_values.password}
tables: test_db.\.*
server-id: 5401-5499
#(Optional) Synchronize data from newly added tables during the incremental phase.
scan.binlog.newly-added-table.enabled: true
#(Optional) Synchronize table and field comments.
include-comments.enabled: true
#(Optional) Prioritize the distribution of unbounded shards to prevent potential Out Of Memory (OOM) issues on TaskManagers.
scan.incremental.snapshot.unbounded-chunk-first.enabled: true
#(Optional) Enable parsing and filtering to accelerate data reads.
scan.only.deserialize.captured.tables.changelog.enabled: true
sink:
type: paimon
name: Paimon Sink
using.built-in-catalog: paimon_dlf_catalog
#(Optional) Specify a commit user. We recommend that you use a different commit user for each job to avoid conflicts.
commit.user: your_job_name
#(Optional) Enable deletion vectors to improve read performance.
table.properties.deletion-vectors.enabled: true
route:
# Synchronize all tables in the MySQL test_db database to the Paimon test_db2 database and retain the table names.
- source-table: test_db.\.*
sink-table: test_db2.<>
replace-symbol: <>
pipeline:
name: MySQL to Paimon Pipeline
Merge sharded databases and tables
To write data from multiple sharded tables into a single downstream table, use the following template:
source:
type: mysql
name: MySQL Source
hostname: <yourHostname>
port: 3306
username: flink
password: ${secret_values.password}
tables: test_db.user\.*
server-id: 5401-5499
#(Optional) Synchronize data from newly added tables during the incremental phase.
scan.binlog.newly-added-table.enabled: true
#(Optional) Synchronize table and field comments.
include-comments.enabled: true
#(Optional) Prioritize the distribution of unbounded shards to prevent potential Out Of Memory (OOM) issues on TaskManagers.
scan.incremental.snapshot.unbounded-chunk-first.enabled: true
#(Optional) Enable parsing and filtering to accelerate data reads.
scan.only.deserialize.captured.tables.changelog.enabled: true
sink:
type: paimon
name: Paimon Sink
using.built-in-catalog: paimon_dlf_catalog
#(Optional) Specify a commit user. We recommend that you use a different commit user for each job to avoid conflicts.
commit.user: your_job_name
#(Optional) Enable deletion vectors to improve read performance.
table.properties.deletion-vectors.enabled: true
route:
# All sharded tables in the MySQL test_db database are merged into a single Paimon test_db.user table.
- source-table: test_db.user\.*
sink-table: test_db.user
pipeline:
name: MySQL to Paimon Pipeline
Add existing tables on job restart
To add existing tables to the synchronization, set scan.newly-added-table.enabled = true and restart the job.
If the job first sets scan.binlog.newly-added-table.enabled to true to capture new tables, do not then set scan.newly-added-table.enabled to true and restart the job to capture existing tables. Otherwise, data is delivered more than once.
source:
type: mysql
name: MySQL Source
hostname: <yourHostname>
port: 3306
username: flink
password: ${secret_values.password}
tables: test_db.\.*
server-id: 5401-5499
scan.startup.mode: initial
# On job restart, check the new tables captured by the tables parameter and take a snapshot.
# Note: Must be used together with scan.startup.mode: initial.
scan.newly-added-table.enabled: true
#(Optional) Synchronize data from newly added tables during the incremental phase.
scan.binlog.newly-added-table.enabled: true
#(Optional) Synchronize table and field comments.
include-comments.enabled: true
#(Optional) Prioritize the distribution of unbounded shards to prevent potential Out Of Memory (OOM) issues on TaskManagers.
scan.incremental.snapshot.unbounded-chunk-first.enabled: true
#(Optional) Enable parsing and filtering to accelerate data reads.
scan.only.deserialize.captured.tables.changelog.enabled: true
sink:
type: paimon
name: Paimon Sink
using.built-in-catalog: paimon_dlf_catalog
#(Optional) Specify a commit user. We recommend that you use a different commit user for each job to avoid conflicts.
commit.user: your_job_name
#(Optional) Enable deletion vectors to improve read performance.
table.properties.deletion-vectors.enabled: true
pipeline:
name: MySQL to Paimon Pipeline
Exclude specific tables in entire-database synchronization
To exclude specific tables, such as test tables or tables that cause repeated job failovers, from an entire-database synchronization job, use the following configuration:
source:
type: mysql
name: MySQL Source
hostname: <yourHostname>
port: 3306
username: flink
password: ${secret_values.password}
tables: test_db.\.*
# Tables that match this regular expression are not synchronized.
tables.exclude: test_db.table1
server-id: 5401-5499
#(Optional) Synchronize data from newly added tables during the incremental phase.
scan.binlog.newly-added-table.enabled: true
#(Optional) Synchronize table and field comments.
include-comments.enabled: true
#(Optional) Prioritize the distribution of unbounded shards to prevent potential Out Of Memory (OOM) issues on TaskManagers.
scan.incremental.snapshot.unbounded-chunk-first.enabled: true
#(Optional) Enable parsing and filtering to accelerate data reads.
scan.only.deserialize.captured.tables.changelog.enabled: true
sink:
type: paimon
name: Paimon Sink
using.built-in-catalog: paimon_dlf_catalog
#(Optional) Specify a commit user. We recommend that you use a different commit user for each job to avoid conflicts.
commit.user: your_job_name
#(Optional) Enable deletion vectors to improve read performance.
table.properties.deletion-vectors.enabled: true
pipeline:
name: MySQL to Paimon Pipeline
Synchronize an entire database to Paimon tables in the Lance file format
Lance is a data storage format designed for vector and multimodal data. Set the file.format table creation parameter to write data to Paimon tables that use the Lance file format. The following example shows the configuration:
source:
type: mysql
name: MySQL Source
hostname: <yourHostname>
port: 3306
username: flink
password: ${secret_values.password}
tables: test_db.\.*
server-id: 5401-5499
#(Optional) Synchronize data from newly added tables during the incremental phase.
scan.binlog.newly-added-table.enabled: true
#(Optional) Synchronize table and field comments.
include-comments.enabled: true
#(Optional) Prioritize the distribution of unbounded shards to prevent potential Out Of Memory (OOM) issues on TaskManagers.
scan.incremental.snapshot.unbounded-chunk-first.enabled: true
#(Optional) Enable parsing and filtering to accelerate data reads.
scan.only.deserialize.captured.tables.changelog.enabled: true
sink:
type: paimon
name: Paimon Sink
using.built-in-catalog: paimon_dlf_catalog
#(Optional) Specify a commit user. We recommend that you use a different commit user for each job to avoid conflicts.
commit.user: your_job_name
#(Optional) Enable deletion vectors to improve read performance.
table.properties.deletion-vectors.enabled: true
# Specify the file format of output files as lance
table.properties.file.format: lance
pipeline:
name: MySQL to Paimon Pipeline
Enable EVOLVE mode
The EVOLVE mode keeps the schema of the downstream table strictly consistent with the schema of the upstream table, including operations such as dropping fields and dropping tables. In this mode, if the downstream cannot apply all schema change events, the job may fail over and cannot recover without intervention. The following example shows the job configuration:
source:
type: mysql
name: MySQL Source
hostname: <yourHostname>
port: 3306
username: flink
password: ${secret_values.password}
tables: test_db.test_source_table
server-id: 5401-5499
#(Optional) Synchronize data from newly added tables during the incremental phase.
scan.binlog.newly-added-table.enabled: true
#(Optional) Synchronize table and field comments.
include-comments.enabled: true
#(Optional) Prioritize the distribution of unbounded shards to prevent potential Out Of Memory (OOM) issues on TaskManagers.
scan.incremental.snapshot.unbounded-chunk-first.enabled: true
#(Optional) Enable parsing and filtering to accelerate data reads.
scan.only.deserialize.captured.tables.changelog.enabled: true
sink:
type: paimon
name: Paimon Sink
using.built-in-catalog: paimon_dlf_catalog
#(Optional) Specify a commit user. We recommend that you use a different commit user for each job to avoid conflicts.
commit.user: your_job_name
#(Optional) Enable deletion vectors to improve read performance.
table.properties.deletion-vectors.enabled: true
pipeline:
name: MySQL to Paimon Pipeline
schema.change.behavior: evolve
FAQ
-
Why might a Paimon job fail with "Heartbeat of TaskManager timed out"?
-
Why might a Paimon job fail with "Sink materializer must not be used with Paimon sink"?
-
Why might a Paimon job fail with "File deletion conflicts detected" or "LSM conflicts detected"?
-
Why might a Paimon job fail with "File xxx not found, Possible causes"?
-
Is Paimon connector data visibility related to checkpoint intervals?