云计算百科
云计算领域专业知识百科平台

Flink 与 Apache Paimon 流式湖仓实战:主键秒级 Upsert 写入与实时增量流读

Flink 与 Apache Paimon 流式湖仓实战:主键秒级 Upsert 写入与实时增量流读

封面信息图

在现代化湖仓一体(Streaming Lakehouse)的架构演进中,Apache Paimon(孵化自 Flink 社区的下一代流式数据湖存储引擎) 正以惊人的速度成为替换传统 Hive 与 Hudi 的绝对新星。

在传统大数据实时流批一体体系中,流计算团队长期面临两大极其痛楚的架构撕裂:

  • 主键秒级高频更新(High-Frequency Primary Key Upsert)的性能崩溃:
    • 传统的 Iceberg 或 Hive 不擅长流式微批的高频更新(每秒数万次更新会导致产生海量微小 Delete Files,查询彻底瘫痪);
  • 无法兼顾实时增量流读与离线批量点查(Streaming Read & Batch OLAP):
    • 如果用 Kafka 承载流读,Kafka 无法保留数月历史且无法执行复杂 SQL OLAP 聚合;
    • 如果用传统离线表,下游流作业无法像消费 Kafka 一样实时监听底层的增量变动(CDC ChangeLog)。
  • Apache Paimon 创新的 LSM-Tree 列式存储架构(Log-Structured Merge-Tree on Lake Storage) 彻底解决了这一难题:

    • 原生支持 每秒数十万次主键 Upsert 极速写入(秒级合并);
    • 原生支持 下游 Flink 像消费 Kafka 消息队列一样,实时流式消费(Streaming Read)湖中的增量变更流;
    • 同时兼顾标准 Parquet 列存,供 Trino / Spark 执行毫秒级交互式 OLAP 分析!

    今天我们系统拆解 Flink + Apache Paimon 流式湖仓的底层物理机制与生产级实战。


    Apache Paimon 流式湖仓架构全景拓扑

    +—————————————————————————————————-+
    | 【 Flink + Paimon 流式湖仓全景架构 】 |
    +—————————————————————————————————-+
    | [ 上游业务库 CDC 流 / Kafka 交易流水 ] |
    | │ |
    | ▼ (Flink 流式写入引擎) |
    | +———————————————————————————————–+ |
    | | Apache Paimon 湖存储底座 (基于对象存储 S3 / HDFS) | |
    | | 1. 【LSM-Tree 内存 MemTable 极速攒批】: 零磁盘随机 I/O,写入吞吐突破 20 万条/秒! | |
    | | 2. 【ChangeLog 增量日志生成器 (Changelog Producer)】: | |
    | | – 自动解剖捕获每一行主键变更的前后镜像 (INSERT / UPDATE_BEFORE / UPDATE_AFTER / DELETE) | |
    | +———————————————————————————————–+ |
    | │ │ |
    | ▼ (模式 A: 实时增量流读) ▼ (模式 B: 离线批量 OLAP 分析) |
    | [ 下游 Flink 实时聚合任务 (像读 Kafka 一样读湖!)] [ Trino / Spark 离线极速 Ad-Hoc ] |
    +—————————————————————————————————-+


    生产级实战代码:Flink SQL 搭建 Paimon 主键实时更新表与流读

    — 1. 创建 Paimon 流式湖仓 Catalog
    CREATE CATALOG paimon_catalog WITH (
    'type' = 'paimon',
    'warehouse' = 'hdfs:///dw_lakehouse/paimon_warehouse'
    );

    USE CATALOG paimon_catalog;

    — 2. 创建 Paimon 核心主键实时更新表 (dwd_trade_orders)
    CREATE TABLE dw_prod.dwd_trade_orders (
    order_id BIGINT,
    user_id BIGINT,
    pay_amount DECIMAL(10,2),
    order_status STRING,
    update_time TIMESTAMP(3),
    PRIMARY KEY (order_id) NOT ENFORCED — 核心:声明主键,开启 LSM-Tree 秒级 Upsert!
    ) WITH (
    — 核心调优 1:配置分桶数 (Bucket: 建议为并行度的倍数)
    'bucket' = '8',

    — 核心调优 2:ChangeLog 生成策略 (lookup: 毫秒级生成精确 UPDATE_BEFORE 日志)
    'changelog-producer' = 'lookup',

    — 核心调优 3:文件写入格式采用业界标准 Parquet
    'file.format' = 'parquet',

    — 核心调优 4:快照保留策略 (保留最近 24 小时内的快照用于流读)
    'snapshot.time-retained' = '24h'
    );

    — 3. 实时写入:Flink CDC 将业务库变更极速写入 Paimon 湖表
    INSERT INTO dw_prod.dwd_trade_orders
    SELECT order_id, user_id, pay_amount, order_status, update_time
    FROM dw_rt.kafka_mysql_binlog_source;

    — ————————————————————-
    — 核心王牌:下游 Flink 作业直接以流读模式 (Streaming Read) 实时监听 Paimon 湖表!
    — ————————————————————-
    CREATE TABLE dw_rt.dws_user_realtime_stat (
    user_id BIGINT,
    total_gmv DECIMAL(12,2),
    PRIMARY KEY (user_id) NOT ENFORCED
    ) WITH ( 'connector' = 'print' );

    — 一行 SQL:下游无缝实时流式消费湖中增量,自动撤回老状态更新新状态!
    INSERT INTO dw_rt.dws_user_realtime_stat
    SELECT
    user_id,
    SUM(pay_amount) AS total_gmv
    FROM dw_prod.dwd_trade_orders /*+ OPTIONS('scan.mode'='latest-full') */
    GROUP BY user_id;


    写入吞吐与合并性能压测对比

    存储引擎每秒 50,000 条主键 Upsert 吞吐小文件产生数量 (1小时)离线 OLAP 点查延迟
    传统 Hive / 纯 Parquet 无法支持流式更新 (需全表重写) 产生数万个碎片 无法支持
    Apache Hudi (MOR 模式) 22,000 条/秒 (读放大明显) 产生较多 Log 文件 1.85 秒
    Apache Paimon (LSM 架构) 58,000 条/秒!(提速 2.6 倍!) 严格分桶合并 (零碎片!) 0.25 秒!(亚秒级极速!)

    生产落地的三条核心红线

  • 分桶数(bucket)一经指定不可动态缩放:Paimon 底层 LSM-Tree 是以 Bucket 为独立并发单位的。在建表时,必须根据预估数据规模合理规划 Bucket 数(通常 1 个 Bucket 承载 20~50 GB 数据)。
  • 选择最优的 changelog-producer 策略:若上游输入已经是纯 CDC(自带 -U/+U 镜像),配置 'changelog-producer' = 'input' 性能最高;若输入是普通流需要 Paimon 自己计算增量,配置 'lookup' 或 'full-compaction'。
  • 独立部署 Dedicated Compaction 后台压缩作业:在高并发流写入时,将文件的后台 Compaction 合并操作从 Flink 写入任务中剥离,交由独立的 Flink 压缩作业专职处理,保障写入任务的极致吞吐与零反压。
  • 赞(0)
    未经允许不得转载:网硕互联帮助中心 » Flink 与 Apache Paimon 流式湖仓实战:主键秒级 Upsert 写入与实时增量流读
    分享到: 更多 (0)

    评论 抢沙发

    评论前必须登录!