共计 2079 个字符,预计需要花费 6 分钟才能阅读完成。
传统方案的业务痛点
某电商推荐系统曾遇到这样的困境:用户夜间加购的商品直到次日中午才会影响推荐结果。排查发现现有架构存在三大缺陷:

- 数据一致性差 :T+ 1 批量导入导致数仓与业务库存在最大 18 小时差异
- 处理延迟高 :每小时运行的 Spark 作业面临 45 分钟调度积压
- 资源消耗大 :全量扫描 ODS 层消耗了集群 60% 的计算资源
技术选型:流批一体架构对比
针对上述问题,我们对比了三种主流方案的核心指标:
| 维度 | Spark Streaming | Flink | ClickHouse |
|---|---|---|---|
| 实时性 | 分钟级 | 秒级 | 秒级 |
| Exactly-Once | 微批模式保障 | 完整支持 | 不支持 |
| 运维成本 | 需要单独维护批流 | 统一批流 API | 需配合 Kafka 使用 |
最终选择 Flink 1.15 + Iceberg 0.13 组合,因其同时满足:
- 通过 Checkpoint 机制保证端到端精确一次
- 利用动态表特性实现 CDC 数据自动映射
- 内置 SSE4.2 指令集优化 Parquet 编码
核心实现细节
Flink CDC 数据捕获示例
-- 创建 MySQL 源表连接
CREATE TABLE mysql_orders (
order_id STRING,
user_id INT,
amount DECIMAL(10,2),
update_time TIMESTAMP(3),
PRIMARY KEY (order_id) NOT ENFORCED
) WITH (
'connector' = 'mysql-cdc',
'hostname' = 'mysql.prod',
'port' = '3306',
'username' = 'flink_user',
'password' = '******',
'database-name' = 'order_db',
'table-name' = 't_orders',
'server-time-zone' = 'Asia/Shanghai'
);
-- 定义 Iceberg 目标表
CREATE TABLE iceberg_orders (
order_id STRING,
user_id INT,
amount DECIMAL(10,2),
update_time TIMESTAMP(3),
PRIMARY KEY (order_id) NOT ENFORCED
) PARTITIONED BY (days(update_time))
WITH (
'format-version'='2',
'write.upsert.enabled'='true',
'write.metadata.compression-codec'='gzip',
'write.parquet.compression-codec'='zstd'
);
-- 启动持续同步作业
INSERT INTO iceberg_orders
SELECT * FROM mysql_orders;
Iceberg 优化关键配置
- 分区策略 :
- 按天分区:
PARTITIONED BY (days(update_time)) -
每月自动创建新分区文件夹
-
压缩算法 :
- Parquet 文件采用 ZSTD(level=3)
-
元数据文件使用 GZIP 压缩
-
小文件合并 :
# 每天凌晨合并小于 128MB 的文件 CALL catalog.system.rewrite_data_files( table => 'db.orders', strategy => 'binpack', options => map('min-input-files', '5', 'target-file-size-bytes', '134217728') );
性能测试数据
在 16 核 /64GB 内存的 Worker 节点上:
- 写入吞吐 :
- 10 亿条订单数据(约 2TB)写入耗时 83 分钟
-
平均吞吐量达到 200,000 records/s
-
查询性能 :
SELECT user_id, SUM(amount) FROM iceberg_orders WHERE update_time >= '2023-07-01' GROUP BY user_id; - 冷查询:12 秒(首次扫描)
- 热查询:1.3 秒(元数据缓存生效后)
生产环境避坑指南
- 敏感字段加密 :
- 使用 AWS KMS 加密信用卡号等字段
-
在 Flink SQL 中通过 UDF 实现动态加解密
-
版本回溯 :
-- 查看历史快照 SELECT * FROM iceberg_orders FOR VERSION AS OF 123456789; -- 回滚到指定版本 CALL catalog.system.rollback_to_snapshot( table => 'db.orders', snapshot_id => 123456789 ); -
Kerberos 认证 :
<!-- flink-conf.yaml --> security.kerberos.login.keytab: /etc/security/keytabs/flink.keytab security.kerberos.login.principal: flink@PROD.COM
开放性问题思考
当 AI 特征工程需要访问 1 小时前的近线数据时:
- 直接查询数仓可能引发资源争抢
- 采用增量导出到特征存储又会增加延迟
- 是否有更优的混合方案?比如:
- 将 Flink State 直接暴露为查询接口
- 使用 Alluxio 构建内存加速层
- 实施分级存储策略(热 / 温 / 冷数据分离)
正文完
