The Streaming Lakehouse: Apache Paimon, Iceberg and Hudi comparison
Introduction
Apache Paimon just graduated from the Apache Incubator.
It started as a project under the Apache Flink umbrella (known as Flink Table Store at the time) and within the last year went from an early-stage umbrella project to an Apache incubator project (to grow more and accommodate more engines like Apache Spark).
Within a year it grew a big community and is deployed in large-scale production environments; upgrading existing lakehouse infra.
It has created a rich ecosystem supporting many query engines like Apache Spark, Apache Hive, Apache Doris, Trino, Presto, and StarRocks, while also providing strong and sophisticated CDC streaming ingestion from many popular sources, including MySQL, Postgres, MongoDB, Apache Kafka, Redpanda and Apache Pulsar.
Apache Flink is a unified compute engine and Apache Paimon was created to provide a unified storage lake storage. Combined with Flink CDC, we can realize an end-to-end Unified (Batch & Streaming) data stack for real-time data analytics.
Apache Paimon is the first wave, that unlocks a true unified batch and streaming architecture.
One SQL, One table/dataset to support a (unified) streaming data platform.
Apache Paimon serves as a unified streaming lake storage. Let’s see first why we needed something like Apache Paimon is needed.
The following illustration depicts typical modern data architectures.
On the one hand, we have streaming, and on the other, we have the Lakehouse (Batch).
These are common scenarios for our customers.
- Deploying Apache Flink along with a streaming storage layer for stream processing allows low latencies (millisecond/second), but is also quite expensive in different dimensions.
- Deploying Apache Flink along with a table format to create a Lakehouse architecture is much cheaper, but is batch and comes with high latencies (hour to day).
If we look at this from a different dimension, as depicted below:
We can see this as a variation of the Lambda Architecture where we need to maintain two different layers i.e. the streaming layer and the batch layer.
The customers have to choose either streaming or lakehouse (batch), while also being unable to truly leverage Flink’s unified batch and streaming API.
What about an intermediate solution? Many streaming use cases don’t require that millisecond or second level of latency.
The unified stack that Apache Paimon complements, aims to allow this architecture that provides a single abstraction for both batch and stream.
Since Apache Paimon builds on the concept of the Lakehouse, many people can’t help but wonder - why create a project from scratch and how it compares to existing table formats.
This blog post aims to address these questions and also show how Apache Paimon compares to Apache Hudi since compared to Iceberg it can also handle streaming data.
Background
The Lakehouse is a relatively new open and scalable architecture that provides comprehensive data analysis and management capabilities. With stronger management capabilities, underlying file storage can better use object storage with lower costs and can accommodate more data.
- Better ACID transaction capability to ensure consistency when reading or writing data (usually SQL) at the same time.
- Cheap and scalable, while storing partitions, files, and statistics.
- Rich query functions, version management capabilities, support for Time Travel, rollback, etc.
- Supports file data skipping and provides good query performance
- Improved timeliness. More fine-grained data updates can improve the timeliness of data.
The original goal was to provide a streaming lake storage layer native for Apache Flink, which means, we needed:
- Strong support for handling streaming data
- Strong support for upserts and CDC data
- Ability to support streaming reads,
- Create a complete changelog; provide downstream apps with correct results.
And at the same time, it needed to be resource and cost-efficient.
Apache Iceberg is a high-performance format for huge analytic tables. By design, it was created for batch processing and support for query engines.
On the other hand, Apache Hudi was created for incremental updates and is a better fit for streaming data.
Let’s first take a high-level overview of all three projects
| Hudi | Iceberg | Paimon | |
|---|---|---|---|
| Core Focus | Incremental updates | Upgrade of Hive format | Streaming Lakehouse |
| Main Engine | Spark | Spark | Flink |
| Started By | Uber | Netflix | Ververica/Alibaba |
| Community Development | 7 years | 6 years | 2 years |
| Computing Ecosystem | Strong, integrated by mainstream engines | Strong, integrated by mainstream engines and a large number of vendors | Strong, integrated by mainstream engines with the strongest Flink integration |
| Version Management | Weak, old design for historical reasons | Strong | Strong |
| Update Mechanism | Medium | Weak | Strong |
| Stream write stream read | Medium | Weak | Strong |
Community development
Github Stars:
- Hudi has the earliest R&D and open source, currently ~5K stars
- Iceberg was open-sourced around 2019 and later came to the top, currently ~5.5K Stars
- Paimon started R&D in 2022, released its official version in 2023, and has received more extensive attention after joining the Apache Incubator. Currently, it has ~2k Stars.
Computing Ecosystem
| Hudi | Iceberg | Paimon | |
|---|---|---|---|
| Spark | Strong | Strong | Strong |
| Flink | Medium | Weak | Strong |
| Hive | Medium | Medium | Medium |
| Doris/StarRocks | Medium | Strong | Medium/Strong |
| Trino | Medium | Strong | Medium |
| Presto | Medium | Strong | Medium |
| Vendor Support | Medium | Strong | Weak |
Apache Iceberg
Apache Iceberg was designed as a Hive alternative to support large analytics. It has grown to become a popular choice for creating a Lakehouse and is well-supported by main query engines.
However because it was created for query engines and needs to maintain compatibility there, it misses the properties required for streaming data and it was also hard to add first-class support to meet the needs of Apache Flink.
Our team invested a lot in Iceberg and has also helped introduce the v2 format in an attempt to support streaming (and upserts), but it is only well suited for small to medium-sized workloads as its design is mainly for batch processing. Moreover, it must maintain compatibility with query engines, and it's hard to make the required changes on the kernel level, to support the needs of Apache Flink and near real-time latencies. Apache Iceberg is a great fit for use cases that require higher latency SLAs, typically about 1 hour or more.
Because the Apache Paimon team has many Apache Iceberg committers and Iceberg has an excellent foundation and file layout design, Paimon builds on this and shares the same file layout. It takes a different approach though with a focus on providing the primitives needed for (unified) streaming and updates.
Limitations that Apache Paimon addresses:
- Strong support for upserts (CDC data) via a variety of rich merge engines.
- Ability to create a complete changelog required for downstream consumers to be able to “see” correct results.
- Consumer-ids to handle snapshot expiration issues, that can result in FileNotFoundExceptions that you can often encounter with Iceberg.
- Strong integration with Flink CDC, for an end-to-end streaming flow, from source to business aggregate layers.
- Automatic (and optimized) compaction to deal with small file problems.
- Due to the LSM, tags for time travel can reuse data files, compared to Iceberg which needs to do full copies, which means way less storage costs.
At the same time, Apache Paimon provides lots of rich functionality and avoids the use of Flink’s state. You can achieve operations like partial updates to replace expensive joins, automatic aggregations, lookup joins, message queue functionality, and more.
Apache Hudi
Apache Hudi is a better fit for streaming data and before Apache Paimon many large-scale companies have adopted it along with Apache Flink while exploring the Streaming Lakehouse (at least for the ingestion layer).
The Apache Paimon team was also involved and committed to Apache Hudi for more than half a year.
Apache Hudi though has lots of complexity (due to historical reasons) and a hard-to-use API. The users need to know lots of internal details, how to decide and make trade-offs, i.e. using a State Index that's easy to use, but has poor performance or Bucket Index that’s high performant but hard to use? CopyOnWrite (CoW) or MergeOnRead (MoR)?
Moreover, the update efficiency needed improvements to meet the requirements, and even with 1-3 minute checkpointing it was prone to backpressure, while also requiring by default 5 (full) merges.
The compatibility between various query engines was somewhat poor, and bug fixes are difficult to converge (due to the complex system design), and need lots of tuning.
Hudi was born for a processing model, closer to Spark’s batch processing design.
It constantly makes detailed transformations on the batch-oriented architecture, and cannot fully adapt to the stream processing and update scenario.
During the investigation, it was found that Apache Hudi could only meet ~15-minute SLAs.
It forcibly improves the stream processing update capabilities on the batch processing architecture, resulting in the architecture becoming more and more complex and hard to maintain; Maintainability was another key consideration.
Another aspect is having first-class support for Apache Flink and the needs of the engine. If you check the Apache Hudi roadmap you will find that there are still very few Flink-related and stream-related things. Often the Hudi community creates a new feature that is only supported by Spark and not by Flink. If Flink wants to support it, it requires lots of redesign and reconstruction and is prone to bugs while also ending up not being supported very well.
Since projects have evolved over the years, we can see that the stability of Hudi has become much better in more recent versions, but overall couldn’t meet the streaming lakehouse requirements needed.
Let’s look next to how it might compare to Apache Paimon as both can be two different solutions for a Streaming Lakehouse Architecture.
Cluster Environment
Cluster Setup:
- 1 master node: m6i.2xlarge
- 2 worker nodes: m6i.4xlarge
The following components and versions are used:
- Paimon: 0.7-SNAPSHOT(Paimon community 0.6 release)
- Hudi: 0.14.0
- Flink: 1.15
- Spark: 3.3.1
- OSS-HDFS: 1.0.0
Streaming Data Ingestion
Streaming data ingestion with good data freshness is a crucial application scenario for table formats and serves as the first step in building a Streaming Lakehouse. The tests in this section refer to the paimon-cluster-benchmark.
We will explore how both perform in two common ingestion use cases, i.e:
- Upserts: a good fit for cdc data ingestion and handle the updates
- Append-only: a good fit for log stream data ingestion like Kafka topics
Apache Paimon and Hudi's read and write capabilities are tested in these respective scenarios.
Apache Flink is used for streaming data lake ingestion.
The deployment mode is Flink Standalone and the Flink configuration is as follows.
As the TM memory size greatly affects the test results, we will test with different memory settings, i.e. 8 GB, 16 GB, and 20 GB.
parallelism.default: 16
jobmanager.memory.process.size: 4g
taskmanager.numberOfTaskSlots: 1
taskmanager.memory.process.size: 8g/16g/20g
execution.checkpointing.interval: 2min
execution.checkpointing.max-concurrent-checkpoints: 3
taskmanager.memory.managed.size: 1m
state.backend: rocksdb
state.backend.incremental: true
table.exec.sink.upsert-materialize: NONEUpsert Scenario
The data lake upsert is used to update or insert new data. When performing an upsert, it checks whether the data to be written already exists in the data lake. If the data already exists, it will be updated; if the data does not exist, the new data will be inserted. Upsert is usually based on a unique identifier or a primary key to determine whether the data already exists.
We are using the datagen connector for Apache Flink to create a test data source. It randomly generates data with primary keys ranging from 0 to 100.000.000.
Apache Flink ingests data via Paimon and Hudi tables respectively, and we measure the total time consumed to write 500.000.000 records. (The total size of parquet files within a single bucket is ~2 GB).
At the same time, we also use Apache Flink to read the written Paimon and Hudi tables in batch mode and count the total duration.
To deal with upsert data scenarios, Apache Paimon uses a Primary Key table, while Apache Hudi makes use of Merge-On-Read tables.
Both Paimon and Hudi support automatic compaction, so we divide the tests further to test with compaction both disabled and enabled.
Disable Compaction
The configuration of the Paimon table is as follows:
'bucket' = '16',
'file.format' = 'parquet',
'file.compression' = 'snappy',
'write-only' = 'true'The number of buckets is the same as the parallelism of Flink and is set to 16. The default file format of Hudi is parquet. To be consistent with Hudi, the file output format is parquet and the compression method is set to snappy.
The configuration of the Hudi table is as follows:
'table.type' = 'MERGE_ON_READ',
'metadata.enabled' = 'false',
'index.type' = 'BUCKET',
'hoodie.bucket.index.num.buckets' = '16',
'write.operation' = 'upsert',
'write.tasks' = '16',
'hoodie.parquet.compression.codec' = 'snappy',
'read.tasks' = '16',
'compaction.schedule.enabled' = 'false',
'compaction.async.enabled' = 'false',
'compaction.max_memory' = '4096/8192/10240' -- Half of the TM process memoryWe are configuring a BucketIndex with a number of buckets set to 16, which is the same as the parallelism of Flink.
Because the reading of the Hudi MOR table is affected by the parameter compaction.max_memory, we configure it to be half of the taskmanager.memory.process.size.
The test results are as follows:
We can see that for upsert scenarios with the compaction disabled, Paimon has better read and write performance than Hudi, while also having less memory requirements.
Next, let’s test with the compaction enabled.
Enabled Compaction
The configuration of the Paimon table is as follows:
'bucket' = '16',
'file.format' = 'parquet',
'file.compression' = 'snappy',
'num-sorted-run.compaction-trigger' = '5' -- default configurationThe configuration of the Hudi table is as follows:
'table.type' = 'MERGE_ON_READ',
'metadata.enabled' = 'false',
'index.type' = 'BUCKET',
'hoodie.bucket.index.num.buckets' = '16',
'write.operation' = 'upsert',
'write.tasks' = '16',
'hoodie.parquet.compression.codec' = 'snappy',
'read.tasks' = '16',
'compaction.schedule.enabled' = 'true',
'compaction.async.enabled' = 'true',
'compaction.tasks' = '16',
'compaction.delta_commits' = '2'
'compaction.max_memory' = '4096/8192/10240' -- Half of the TM process memoryThe total amount of time required for testing is small (the number of checkpoints is correspondingly small), and as the number of uncompacted log files increases, the compaction memory required by Hudi increases a lot.
Therefore, the configuration of compaction.delta_commits should be set to 2 to ensure that compaction can be completed during writing.
The test results are as follows:
In upsert scenarios, when compaction is enabled, Paimon has better read and write performance than Hudi. Compared with the previous test with compaction disabled, the write performance of Paimon and Hudi is reduced, but the read performance is improved.
The compaction of Hudi consumes much memory, requires long running time, and is executed asynchronously.
When the writing task is completed, the still-running compaction will stop.
We observed that when the TaskManager memory reaches 20 GB, Hudi still has 4 delta commits that are not compacted (even if compaction.delta_commits=2 is configured).
Append Scenario
The other scenario of data ingestion is data append write, such as log ingestion.
The test data source is again created using the Apache Flink datagen connector, to write data into Paimon and Hudi tables.
Similarly, we measure the total time needed for Flink to ingest 500.000.000 records (in the append scenario, Paimon and Hudi buckets are not required. The configuration of the Paimon table is as follows:
'bucket' = '-1',
'file.format' = 'parquet',
'file.compression' = 'snappy'The configuration of the Hudi table is as follows:
'table.type' = 'COPY_ON_WRITE',
'metadata.enabled' = 'false',
'write.operation' = 'insert',
'write.tasks' = '16',
'hoodie.parquet.compression.codec' = 'snappy',
'read.tasks' = '16',
'write.insert.cluster' = 'false',
'clustering.schedule.enabled' = 'false',
'clustering.async.enabled' = 'false'The amount of data within a single batch is large enough, so we don’t have small file problems. This means we can disable clustering.
The test results are as follows:
In the append scenario, Paimon has better read and write performance than Hudi, and neither of them consumes high TaskManager memory.
Ease of Use and Functionality
Performance is just one aspect and is not always the decisive factor. The main reasons we have seen lots of customers (and community) choose (or migrate) from Apache Hudi to Paimon to create a streaming lakehouse are as follows:
- First-class support for Apache Flink
- Reliability, stability, and ease of use.
- No historical burden
- Strong integration with the complete unified stack
Along with rich functionality and optimizations for streaming data. Apache Paimon offers lots of rich functionality for streaming data; a variety of merge engines, different ways to deal with out-of-order data like sequence ids and sequence groups, partial updates to replace expensive streaming joins, aggregation engines, cross-partition updates, and more. Also for append-only scenarios, it has the Append-for Message Queue table that allows implementing message queue functionality directly on the lake.
Let’s also look at the ease of use more closely. One thing you might have noticed in previous examples is the difference in the number of configurations just for creating the respective tables. Apache Paimon aims for simplicity.
Let’s also look at another example. Let’s assume a common scenario; we want to create an aggregation layer on the lake, which contains tables with business aggregates ready to be consumed.
This is a really simple example of some basic aggregation, to calculate the uv, pv, and fee_sum columns. All we need to do is specify the aggregation function. At the same time, we can easily create a complete changelog (by specifying changelog-producer=lookup) for downstream readers more efficiently (rather than having to use the expensive upsert-materialize operator).
CREATE TABLE paimon_catalog.order_gold.shop_user_aggregates (
shop_id BIGINT,
user_id BIGINT,
ds STRING COMMENT 'hour',
uv BIGINT COMMENT 'the user's consumption times in the shop within this hour',
pv BIGINT COMMENT 'the user's consumption times in the shop within this hour',
fee_sum DECIMAL(20, 2) COMMENT 'the user's total spend in the shop within this hour'
) WITH (
'primary-key' = 'shop_id, user_id, ds',
,'bucket' = '8',
,'changelog-producer' = 'lookup',
,'file.format' = 'parquet',
,'file.compression' = 'snappy',
,'merge-engine' = 'aggregation',
,'fields.uv.aggregate-function' = 'sum',
,'fields.pv.aggregate-function' = 'sum',
,'fields.fee_sum.aggregate-function' = 'sum',
,'metadata.stats-mode' = 'none'
);All you need to do is specify the aggregation functions you require on the fields.
For Hudi tables on the other hand, if you want to implement similar aggregation operations, you need to use a custom Payload or Merger.
In this example, a custom Merger is used to aggregate the uv, pv, and fee_sum fields of records with the same key. The logic is as follows:
public class OrdersLakeHouseMerger extends HoodieAvroRecordMerger {
@Override
public Option<Pair<HoodieRecord, Schema>> merge(HoodieRecord older, Schema oldSchema, HoodieRecord newer, Schema newSchema, TypedProperties props) throws IOException {
// ...
Object oldData = older.getData();
GenericData.Record oldRecord = (oldData instanceof HoodieRecordPayload)
? (GenericData.Record) ((HoodieRecordPayload) older.getData()).getInsertValue(oldSchema).get()
: (GenericData.Record) oldData;
Object newData = newer.getData();
GenericData.Record newRecord = (newData instanceof HoodieRecordPayload)
? (GenericData.Record) ((HoodieRecordPayload) newer.getData()).getInsertValue(newSchema).get()
: (GenericData.Record) newData;
// merge uv
if (HoodieAvroUtils.getFieldVal(newRecord, "uv") != null && HoodieAvroUtils.getFieldVal(oldRecord, "uv") != null) {
newRecord.put("uv", (Long) oldRecord.get("uv") + (Long) newRecord.get("uv"));
}
// merge pv
if (HoodieAvroUtils.getFieldVal(newRecord, "pv") != null && HoodieAvroUtils.getFieldVal(oldRecord, "pv") != null) {
newRecord.put("pv", (Long) oldRecord.get("pv") + (Long) newRecord.get("pv"));
}
// merge fee_sum
if (HoodieAvroUtils.getFieldVal(newRecord, "fee_sum") != null && HoodieAvroUtils.getFieldVal(oldRecord, "fee_sum") != null) {
BigDecimal l = new BigDecimal(new BigInteger(((GenericData.Fixed) oldRecord.get("fee_sum")).bytes()), 2);
BigDecimal r = new BigDecimal(new BigInteger(((GenericData.Fixed) newRecord.get("fee_sum")).bytes()), 2);
byte[] bytes = l.add(r).unscaledValue().toByteArray();
byte[] paddedBytes = new byte[9];
System.arraycopy(bytes, 0, paddedBytes, 9 - bytes.length, bytes.length);
newRecord.put("fee_sum", new GenericData.Fixed(((GenericData.Fixed) newRecord.get("fee_sum")).getSchema(), paddedBytes));
}
HoodieAvroIndexedRecord hoodieAvroIndexedRecord = new HoodieAvroIndexedRecord(newRecord);
return Option.of(Pair.of(hoodieAvroIndexedRecord, newSchema));
}
}Conclusion
Apache Paimon was created to “complete” the unified stack. It brings the first wave of streaming data to the lake, while we work towards even millisecond latencies to bring true streaming data analytics to the lake. Currently, we are tackling all those customer use cases that require strong support for streaming and upsert data and have ~1-minute SLAs.
Depending on the volumes some might also get more down to the seconds > 30 seconds, but 1+ minute is what we recommend.
At Ververica we provide flexibility and can support your data architectures in many different dimensions; whether it's real-time stream processing or a lakehouse.
At the same time though as part of our streaming data platform and our vision for democratizing streaming data, we invest heavily into our Streamhouse.
Apache Paimon is what powers VMT - Ververica Materialized Tables - the core abstraction of Ververica’s Streamhouse. We will see more on this soon.
You can get started with Apache Paimon easily using the Ververica Cloud for free, here.