# 计算存储调优、业务实现与经验收尾

> 收尾页：从统一配置入口出发解读 Flink 状态与检查点、微批、Paimon 文件策略等系统优化点，说明元数据在 Hive Metastore、数据落 HDFS 的存储形态，介绍 StarRocks 外部分区目录直读与可视化，最后把交易、流量、用户等业务域如何落到各层做一次收束，并沉淀可复用经验。

- Repository: fengyu-eng/paimon-datalake
- GitHub: https://github.com/fengyu-eng/paimon-datalake
- Human wiki: https://grok-wiki.com/public/wiki/fengyu-eng-paimon-datalake-3aa80a38b533
- Complete Markdown: https://grok-wiki.com/public/wiki/fengyu-eng-paimon-datalake-3aa80a38b533/llms-full.txt

## Source Files

- `src/main/java/com/fy/warehouse/config/FlinkConfigUtil.java`
- `src/main/resources/core-site.xml`
- `src/main/resources/hdfs-site.xml`
- `pom.xml`
- `dependency-reduced-pom.xml`
- `StarRocks.sql`

---

<details>
<summary>相关源文件</summary>

以下文件用于生成此维基页面:

- [src/main/java/com/fy/warehouse/config/FlinkConfigUtil.java](src/main/java/com/fy/warehouse/config/FlinkConfigUtil.java)
- [src/main/resources/core-site.xml](src/main/resources/core-site.xml)
- [src/main/resources/hdfs-site.xml](src/main/resources/hdfs-site.xml)
- [pom.xml](pom.xml)
- [StarRocks.sql](StarRocks.sql)
- [src/main/java/com/fy/warehouse/ODS/ods_log_inc.java](src/main/java/com/fy/warehouse/ODS/ods_log_inc.java)
- [src/main/java/com/fy/warehouse/ODS/ods_order_info_full.java](src/main/java/com/fy/warehouse/ODS/ods_order_info_full.java)
- [src/main/java/com/fy/warehouse/DIM/dim_user_zip_full.java](src/main/java/com/fy/warehouse/DIM/dim_user_zip_full.java)
- [src/main/java/com/fy/warehouse/DIM/dim_sku_full.java](src/main/java/com/fy/warehouse/DIM/dim_sku_full.java)
- [src/main/java/com/fy/warehouse/dwd/dwd_trade_pay_detail_suc_full.java](src/main/java/com/fy/warehouse/dwd/dwd_trade_pay_detail_suc_full.java)
- [src/main/java/com/fy/warehouse/dwd/dwd_user_login_full.java](src/main/java/com/fy/warehouse/dwd/dwd_user_login_full.java)
- [src/main/java/com/fy/warehouse/dws/dws_trade_user_order_nd_full.java](src/main/java/com/fy/warehouse/dws/dws_trade_user_order_nd_full.java)
- [src/main/java/com/fy/warehouse/dws/dws_traffic_page_visitor_page_view_nd_full.java](src/main/java/com/fy/warehouse/dws/dws_traffic_page_visitor_page_view_nd_full.java)
- [src/main/java/com/fy/warehouse/ads/ads_activity_stats_full.java](src/main/java/com/fy/warehouse/ads/ads_activity_stats_full.java)
- [src/main/java/com/fy/warehouse/udf/JsonActionsArrayParser.java](src/main/java/com/fy/warehouse/udf/JsonActionsArrayParser.java)
- [src/main/datagenerator/business_code/README.md](src/main/datagenerator/business_code/README.md)
- [src/main/datagenerator/userlog_code/README.md](src/main/datagenerator/userlog_code/README.md)

</details>

# 计算存储调优、业务实现与经验收尾

这是电商实时数据湖仓系列的收尾页。前面的页面分别讲清了架构、各数据分层的数据流与明细实现,本页把这些线索合到一处:先从代码里唯一的统一配置入口出发,解释 Flink 状态与检查点、微批、Paimon 文件策略这些"计算与存储调优点"到底解决了什么问题;再说明"元数据进 Hive Metastore、数据落 HDFS"的双层存储形态,以及 StarRocks 如何以外部分区表的方式直读并交给 DataRT 可视化;最后把交易、流量、用户、互动、工具等业务域在 ODS→DIM→DWD→DWS→ADS 每层落成哪些表做一次收束,沉淀一批可直接复用的经验。读完本页,应该能回答"这套湖仓为什么这样配、每个业务域的指标是怎么一层层算出来的"。

仓库的全部 Flink 作业共用同一份配置,因此看懂了它,就同时看懂了所有作业的运行与存储方式。

## 统一配置入口:状态、检查点、微批与文件策略

所有作业的入口逻辑一致:`main` 方法先取 `FlinkConfigUtil.getFlinkConfig()` 得到一份带满默认项的 `Configuration`,再以 `.withConfiguration(config)` 创建流式 `TableEnvironment`,之后所有 DDL/DML 都跑在这个环境上([FlinkConfigUtil.java:9-23](src/main/java/com/fy/warehouse/config/FlinkConfigUtil.java))。这份集中配置是按"边踩坑边修"的方式沉淀出来的,注释里保留了大量"关键修改"标记,本身就是最有价值的排障记录。

### Flink 状态与检查点

配置里最核心的一组开关围绕"状态能否持久化、能否恢复":

| 配置项 | 取值 | 作用与背景 |
| --- | --- | --- |
| `execution.checkpointing.enabled` | `true` | 注释明确:原为 `false`,导致状态无法持久化,这是第一个被修复的点 |
| `state.backend` | `filesystem` | 状态后端选文件系统;注释提示大状态可换 RocksDB(并预留了 `rocksdb` 与对应 checkpoint 目录两行) |
| `execution.checkpointing.interval` | `60s` | 检查点间隔 |
| `execution.checkpointing.timeout` | `300s` | 5 分钟超时,避免检查点长时间卡住 |
| `execution.checkpointing.tolerable-failure-number` | `3` | 容忍 3 次连续失败,网络抖动不至于直接让作业失败 |
| `execution.checkpointing.mode` | `EXACTLY_ONCE` | 精确一次语义 |
| `table.exec.state.ttl` | `86400000` | 状态 TTL 24 小时,控制状态体积 |
| `taskmanager.memory.process.size` / `task.heap.size` | `1024m` / `512m` | 明确 TaskManager 进程与堆内存,配合文件系统状态后端控制状态占用 |

关键结论:检查点目录直接指到 HDFS(`hdfs://192.168.10.102:8020/flink/checkpoints`),状态与检查点都落在分布式文件系统上;这正是"解决了状态超限"的手段——状态不再只存在于 TaskManager 本地内存,而是周期性快照到 HDFS,作业重启后可恢复。

### 微批(Mini-Batch)配置

```java
config.setString("table.exec.mini-batch.enabled", "true");
config.setString("table.exec.mini-batch.allow-latency", "500ms");
config.setString("table.exec.mini-batch.size", "1000");
```

微批把同一 key 上 500ms 内到达、或攒够 1000 条的数据攒批处理,让聚合、去重等高成本算子减少逐条触发,以亚秒级延迟换取吞吐。由于运行时是 `STREAMING`(见 21 行),微批是在流上叠加的批化优化,不是改成批处理。

### Paimon 侧的文件写入策略

```java
config.setString("paimon.sink.batch-size", "100");   // 小批次快速写入
config.setString("paimon.sink.buffer-time", "0s");   // 关闭缓冲等待
config.setString("paimon.file.size", "64mb");        // 小文件即可,避免合并
```

这三行体现了教学场景的文件策略取舍:批量 100、缓冲 0s 是为了"来一条写一条、尽快可见",`file.size = 64mb` 的注释直说"小文件即可,避免合并",即不做 aggressive 的 small-file 合并,换取写入即时性和实现简单。这与生产库追求大文件、低文件数的方向不同,是本项目定位决定的。

### 与 Hive / HDFS 协同的必要配置

同一份配置里还埋了几个容易踩的环境类配置:Hive 版本必须与集群一致(`hive.version=3.1.3`);开启 Hive 兼容读取器;关闭分区剪枝以免 Flink 扫描分区卡住;把 `HADOOP_USER_NAME` 与 `hadoop.user.name` 设为 `root` 并对齐 HDFS 地址(`fs.defaultFS`)、Hive Metastore(`thrift://192.168.10.102:9083`);HDFS RPC 超时通过 `env.java.opts` 注入(`dfs.client.socket-timeout`、`ipc.client.connect.timeout`、`ipc.client.connect.max.retries=5`),注释记录这是为"解决 RPC 中断"加的。也就是说,一份配置同时决定了**运行时语义(流/检查点/微批)、落盘文件行为(Paimon)、以及外围依赖(Hive/HDFS)怎么连**。Sources:[FlinkConfigUtil.java:24-73](src/main/java/com/fy/warehouse/config/FlinkConfigUtil.java)

## 存储形态:元数据在 Hive Metastore,数据落 HDFS

湖仓的表不把数据存在 Hive 表里,而是"借壳":建一个 Paimon Catalog,`metastore = hive` 指向 Hive Metastore,`warehouse` 指向 HDFS 路径——表结构(库、表、字段、分区元数据)注册进 Hive Metastore,真正的列式数据文件写在 HDFS 的 warehouse 目录下,由 Paimon 管理:

```sql
CREATE CATALOG paimon_hive WITH (
    'type' = 'paimon',
    'metastore' = 'hive',
    'uri' = 'thrift://192.168.10.102:9083',
    'warehouse' = 'hdfs://192.168.10.102/user/hive/warehouse'
);
```

这段 Catalog DDL 在几乎每个作业里原样出现(如 [ods_log_inc.java:72-79](src/main/java/com/fy/warehouse/ODS/ods_log_inc.java)、[dim_user_zip_full.java:19-24](src/main/java/com/fy/warehouse/DIM/dim_user_zip_full.java))。落到 Paimon 表的数据文件统一用 `parquet` 格式,并普遍开启两类存储参数:写缓冲 `write-buffer-size = 512mb` + `write-buffer-spillable = true`(缓冲写不下就溢写到磁盘);以及**分区生命周期管理**:

```java
'partition.expiration-time' = '1 d',
'partition.expiration-check-interval' = '1 h',
'partition.timestamp-formatter' = 'yyyy-MM-dd',
'partition.timestamp-pattern' = '$k1'
```

分区字段是每个表都有的 `k1`(业务日期字符串,如 `2024-05-01`),Paimon 按 `$k1` 从分区名解析时间戳,1 小时检查一次、自动清理超过 1 天的分区,避免 HDFS 上堆积过多历史分区。ODS 层业务表建表时即声明 `PRIMARY KEY (id, k1) NOT ENFORCED` 与 `PARTITIONED BY (k1)`,形成"主键 + 日期分区"的通用骨架([ods_order_info_full.java:78-116](src/main/java/com/fy/warehouse/ODS/ods_order_info_full.java))。

配套的 Hadoop 环境配置在 XML 中:`core-site.xml` 声明默认文件系统为 `hdfs://hadoop102:8020`、本地数据目录,并把 `root` 配成全局可代理用户(`hadoop.proxyuser.root.*` 均为 `*`,对应代码里统一用 root 访问 HDFS);`hdfs-site.xml` 配置 NameNode/SecondaryNameNode 的 Web 地址与节点白名单、黑名单。换言之:**代码里的 IP/端口与 XML 里的集群地址是一体两面**,迁移集群时两处要同步改。Sources:[core-site.xml:20-49](src/main/resources/core-site.xml)、[hdfs-site.xml:20-39](src/main/resources/hdfs-site.xml)

从依赖侧看,工程同时引入 Paimon 的 Flink 连接器与 Hive 3.1 连接器(`paimon-flink-1.17` + `paimon-hive-connector-3.1`,见 [pom.xml:91-100](pom.xml)),并统一了 Flink 1.17.2、Hadoop 3.1.3、Hive 3.1.3、MySQL CDC 2.4.2 等版本;打包用 shade 插件合并 `META-INF/services`,正是为了让 Hive Catalog 等工厂类在打成 fat jar 后仍能被 SPI 找到([pom.xml:205-245](pom.xml))。

## StarRocks 外部分区表直读与可视化

Hive 元数据 + HDFS 数据的形态带来一个直接收益:**下游分析引擎可以不改数据、以外部目录方式直读**。仓库的 `StarRocks.sql` 演示了在 StarRocks 侧建一个 Paimon External Catalog:

```sql
CREATE EXTERNAL CATALOG paimon_catalog
PROPERTIES (
    'type' = 'paimon',
    'paimon.catalog.type' = 'hive',
    'warehouse' = 'hdfs://192.168.10.102:8020/user/hive/warehouse',
    'hive.metastore.uris' = 'thrift://192.168.10.102:9083',
    'hadoop.username' = 'root'
);
use ads;
show tables;
select * from ads_coupon_stats_full;
```

要点是 `paimon.catalog.type = hive` + 与 Flink 侧完全相同的 warehouse 与 metastore 地址:StarRocks 通过 Hive Metastore 拿到各层分区表的元数据,再按 HDFS 路径直接读取 Paimon 写的 parquet 分区数据,整个链路**不复制、不导入、不经过 HiveServer2**,做到了"建一个 Catalog 即可查数仓各层"。示例里直接 `use ads` 查 ADS 应用层表,正是面向可视化取数的场景——README 中说明可视化层由 StarRocks(MPP 查询)+ DataRT(可视化工具)构成,Flink 只负责把各层数据写进 Paimon,查询压力由 StarRocks 承接。Sources:[StarRocks.sql:1-13](StarRocks.sql)

## 业务域落地收束:从数据生成到各层表

要讲清"业务如何落到各层",先明确最上游的数据从哪来:仓库自带两套 Spring 模拟生成器。业务数据生成器面向 MySQL 的 `gmall` 库,按基础数据、商品、营销、交易、用户、仓储等模块持续写入([business_code/README.md](src/main/datagenerator/business_code/README.md));用户行为日志生成器输出 `common / start / page / actions / displays / err / ts` 结构的 JSON,发往 Kafka 主题 `ODS_BASE_LOG`([userlog_code/README.md](src/main/datagenerator/userlog_code/README.md))。

一张图概括全链路与各层归属:

```mermaid
flowchart LR
    subgraph SRC[数据源]
        M[(MySQL gmall 业务库)]
        K[(Kafka ODS_BASE_LOG)]
    end
    subgraph FLINK[Flink 流式作业]
        CDC[MySQL-CDC 同步]
        KV[Kafka Connector 读取]
    end
    subgraph LAKE[Paimon 湖仓 · Hive Metastore 元数据 + HDFS 数据]
        ODS[ODS 原始层]
        DIM[DIM 维度层]
        DWD[DWD 明细层]
        DWS[DWS 汇总层]
        ADS[ADS 应用层]
    end
    subgraph BI[可视化]
        SR[StarRocks External Catalog 直读]
        RT[DataRT 报表]
    end
    M --> CDC --> ODS
    K --> KV --> ODS
    ODS --> DIM
    ODS --> DWD
    DIM --> DWD
    DWD --> DWS --> ADS
    ADS --> SR --> RT
```

### ODS:两条采集路径

ODS 层是两类源的分岔口。日志类(`ods_log_inc`)用 Kafka Connector 订阅 `ODS_BASE_LOG`,source 侧把 `common/start/page` 声明成嵌套 `ROW`、`actions/displays` 暂存为字符串;落到 Paimon 时通过 `SELECT` 把嵌套列全部展开成 `common_ar、common_mid、page_during_time……` 这样的平铺大宽列,主键用 `kafka_partition+offset+timestamp` 拼接保证唯一([ods_log_inc.java:24-47](src/main/java/com/fy/warehouse/ODS/ods_log_inc.java)、[ods_log_inc.java:87-162](src/main/java/com/fy/warehouse/ODS/ods_log_inc.java))。业务类(如 `ods_order_info_full`)则用 MySQL CDC(`scan.startup.mode = initial`,先全量后增量)把 `gmall.order_info` 等表搬到 Paimon,并按 `create_time` 生成 `k1` 天分区([ods_order_info_full.java:50-59](src/main/java/com/fy/warehouse/ODS/ods_order_info_full.java)、[ods_order_info_full.java:121-139](src/main/java/com/fy/warehouse/ODS/ods_order_info_full.java))。

### DIM:维表做宽表与拉链

DIM 层把分散在多个 ODS 业务表的属性收敛成可直接 join 的宽维表。`dim_sku_full` 是典型:**sku 主体 left join spu、三级/二级/一级分类、品牌、属性聚合**,一次拼出带全链路分类名的商品维表([dim_sku_full.java:64-140](src/main/java/com/fy/warehouse/DIM/dim_sku_full.java) 附近,insert 段见文件尾部多表 left join);`dim_user_zip_full` 则是**拉链表**:对 `ods_user_info_full` 按 `id` 开窗取最新一条,`start_date` 取 operate_time 当天、`end_date` 置为 `9999-12-31`,并以 `md5()` 处理姓名、手机、邮箱等敏感字段([dim_user_zip_full.java:33-86](src/main/java/com/fy/warehouse/DIM/dim_user_zip_full.java));`dim_date_full` 这类日期维表由 MySQL 侧 `dim_date` 表 CDC 同步而来([dim_date_full.java:20-39](src/main/java/com/fy/warehouse/DIM/dim_date_full.java))。

### DWD:按业务域建模的事实明细

DWD 是"业务域"归属最明显的一层。每个 Java 类开头的注释即声明域归属,汇总如下:

| 业务域 | 事实表 | 计算要点 |
| --- | --- | --- |
| 交易域 | 下单 `dwd_trade_order_detail_full`、支付成功 `dwd_trade_pay_detail_suc_full`、退单 `dwd_trade_order_refund_full`、退款支付成功 `dwd_trade_refund_pay_suc_full`、取消 `dwd_trade_cancel_detail_full`、加购 `dwd_trade_cart_add_full`、购物车快照 `dwd_trade_cart_full` | 以订单明细为主体,多路 left join 支付、订单、活动、优惠券、字典表,算各类**金额分摊**(如 `split_original_amount = sku_num * order_price`),并把支付类型、来源类型等编码关联字典表翻译成名称 |
| 流量域 | 页面浏览 `dwd_traffic_page_view_full`、启动 `dwd_traffic_start_full`、动作 `dwd_traffic_action_full`、曝光 `dwd_traffic_display_full`、错误 `dwd_traffic_error_full` | 由 `ods_log_inc` 展开而来,`actions/displays` 这类数组用自定义 UDF `JsonActionsArrayParser`(Jackson 解析数组首元素成 `action_id/item/item_type/ts` 行)配合维表补地区、品牌 |
| 用户域 | 登录 `dwd_user_login_full` | 从日志里按"设备 + 会话起点"拼 `session_id`、开窗取会话首个事件界定登录,再 join 省份维表补省份;并对单表 `ALTER TABLE ... SET ('sink.parallelism'='10')` 单独调大写入并行度 |
| 互动域 | 收藏 `dwd_interaction_favor_add_full`、评论 `dwd_interaction_comment_full` | 收藏/评论业务事实 |
| 工具域 | 优惠券领取 `dwd_tool_coupon_get_full`、下单用券 `dwd_tool_coupon_order_full`、支付用券 `dwd_tool_coupon_pay_full` | 券的生命周期分阶段建模 |

以支付成功事实为例,其 SQL 结构是"明细子查询 × 多路维表/字典 left join + 字典过滤(parent_code='11' 支付方式、'24' 来源类型)"([dwd_trade_pay_detail_suc_full.java:94-185](src/main/java/com/fy/warehouse/dwd/dwd_trade_pay_detail_suc_full.java))。注意这里 DWD 仍然在 Flink SQL 里用 `left join` 直查 ODS Paimon 表——维表关联走的是**批量小表 join 明细流**的简化路线,这也是为什么全局并行度默认为 1、多数表按分区追加。

### DWS:按粒度的 N 日/累计汇总

DWS 把 DWD 明细按"分析粒度 + 时间跨度"聚合,命名规律是 `域_粒度_指标_nd/td`:用户粒度最近 30 日(`dws_trade_user_order_nd_full`)、用户粒度累计(`dws_trade_user_order_td_full`,带首末次下单日期 `order_date_first/last` 与 `_td` 累计口径)、用户 × SKU 粒度 30 日、省份粒度 30 日(`dws_trade_province_order_nd_full`)、活动/优惠券粒度订单 30 日,以及流量侧的访客 × 页面粒度 30 日。统计口径统一为"从明细里取 `k1` 与基准日(代码中硬编码如 `2024-05-31`)相差 0~30 天的记录再 `group by` 聚合",例如下单次数 `count(order_id)`、金额 `sum(split_*)`([dws_trade_user_order_nd_full.java:60-82](src/main/java/com/fy/warehouse/dws/dws_trade_user_order_nd_full.java));流量侧同理,按 `mid_id/brand/model/operate_system/page_id` 分组算 30 日浏览时长与访问次数([dws_traffic_page_visitor_page_view_nd_full.java:55-84](src/main/java/com/fy/warehouse/dws/dws_traffic_page_visitor_page_view_nd_full.java))。

### ADS:面向指标的极薄应用层

ADS 直接消费 DWS 汇总,做最轻的二次加工,产出可直接画图的口径。例如活动补贴率表(`ads_activity_stats_full`)就是把 DWS 的活动 30 日数据转成"补贴率 = 活动优惠金额 / 原始金额"的比率指标,字段只有 `dt / activity_id / activity_name / start_date / reduce_rate`([ads_activity_stats_full.java:42-49](src/main/java/com/fy/warehouse/ads/ads_activity_stats_full.java));优惠券统计同理落 `ads_coupon_stats_full`。由于 ADS 建表不带 `k1` 分区,它更像"轻量结果集",也正是 StarRocks `use ads; select * from ads_coupon_stats_full;` 直接可视化的对象。

## 经验沉淀

把前面所有细节收拢,可复用的经验大致有四条:

1. **集中配置、注释留痕排障过程。** 全仓库只有一个配置入口,每处调优都带着"原值是什么、导致什么问题、改成什么"的注释(如开启 checkpoint、换文件系统状态后端、加 HDFS RPC 超时)。新环境复刻时,逐条对注释检查即可,不必重新踩坑。
2. **"读时展开、写时平铺"的日志处理套路。** Kafka 日志源以嵌套 ROW + JSON 字符串进 ODS,落地时一次 SELECT 展开成平铺列;数组类字段(actions/displays)留给 DWD 用 UDF 拆。既保住了原始数据,又让下游好查。
3. **统一的"主键 + k1 天分区 + 分区过期"表骨架。** 每层表都以 `(业务键, k1)` 为主键、按天分区,配 `partition.expiration-time=1 d` 自动清历史;Flink 侧文件走 `file.size=64mb` 小文件快速写。这套骨架被 ODS/DIM/DWD/DWS 数百行建表 SQL 原样复用,是一致性的来源。
4. **多引擎共享一份湖上数据。** 写入用 Flink + Paimon(Hive Catalog),查询用 StarRocks External Catalog 指向同一 metastore/warehouse 直读,可视化在 DataRT。切引擎只需改连接声明,数据文件不复制——这是"湖"相对传统数仓最直接的收益。

一点提醒(同样来自代码):DWS 的 30 日窗口基准日(`2024-05-31`、`2024-05-01`)与 `application.properties` 的 `mock.date=2024-05-01` 是**硬编码的教学日期**,业务日期推进时必须同步改,否则聚合窗口会空;这是把本仓库迁移到真实日期口径时最容易漏的改动点。Sources:[dws_trade_user_order_nd_full.java:80](src/main/java/com/fy/warehouse/dws/dws_trade_user_order_nd_full.java)、[src/main/datagenerator/business_jar/application.properties:16](src/main/datagenerator/business_jar/application.properties)

综上:这套实时数据湖仓以"Flink 流式作业 + Paimon 湖表 + Hive 元数据/HDFS 存储"为底座,用一个统一配置入口管住状态与检查点、微批、文件策略等全部计算存储调优点;业务数据与日志数据分别在 ODS 汇聚,经 DIM 宽表与拉链、DWD 分域事实、DWS 分粒度汇总,最后在 ADS 形成指标,交给 StarRocks 外部分区直读并可视化。各业务域(交易、流量、用户、互动、工具)在每一层都有清晰、命名规律统一的落点,配合集中配置与通用表骨架,整条链路从数据生成到报表呈现保持了一致的可复现性——这套"配置入口 + 统一表骨架 + 分层分域"的组织方式,是比任何单张表都更值得带走的东西。
