一、架构总览
二、项目亮点
2.1 架构选型解决 Lambda 架构痛点
| 维度 | 传统 Lambda 架构 | Paimon 流批一体架构 | 改善幅度 |
|---|---|---|---|
| 存储 | 实时表 + 离线表分开(~20+ 张冗余表) | 统一一张 Paimon 表 | 表数量减少 60%+ |
| 查询 | 需合并两份数据,逻辑复杂 | 下游直接查一张表 | 查询链路缩短 50% |
| 口径 | 实时/离线一致率 ~92% | 修正后 100% 一致 | 口径问题归零 |
| 维护成本 | 两套链路、两套代码(~40+ Flink 任务) | 共享存储,各自写入(~25 个任务) | 任务数减少 35%+ |
2.2 分层修正的数据质量保障机制 *
设计了完整的逐层修正触发链路,修正覆盖率 100%,修正完成时间 T+1 06:00前:
2.3 查询层架构设计
Paimon 做统一存储,Doris 做 OLAP 加速查询:
- Paimon 负责流批写入、snapshot 隔离、changelog 产出
- Doris 通过 Catalog 直接读取 Paimon 数据,提供高性能聚合查询
- 查询性能:单表聚合查询 P99 <1.5s
- 并发能力:支持100+ QPS的 OLAP 查询(满足看板高峰期需求)
2.4 成本收益
| 指标 | 改造前 | 改造后 | 收益 |
|---|---|---|---|
| 系统运维成本 | 4 套(Hive + Kafka + ClickHouse + ES) | 2 套(Paimon + Doris) | 运维成本降低 50% |
| 存储成本(月) | ~15TB 冗余存储 | ~9TB | 节省 ~40% |
| 数据链路故障率 | 月均 3-4 次口径不一致告警 | 月均 0-1 次 | 故障率降 80% |
三、难点与解决方案
3.1 小文件问题 *
问题描述:
Flink 流式写入 Paimon 时,每个 checkpoint 都会产生新的数据文件。写入并发高时,短时间内产生大量小文件(单文件 KB~几 MB 级别),严重影响下游查询性能。
影响:
- Doris 查询 Paimon 表时 P99 延迟从预期 3s 飙升到 10s+
- HDFS NameNode 压力增大
解决方案:
① 减少文件产生频率— 调整 Flink checkpoint 间隔:
-- 从 1min 调整为 3min,每个 bucket 产生文件的频率降低 3 倍SET 'execution.checkpointing.interval' = '3min';② 自动 compaction— 配置合并触发阈值:
-- 单个 bucket 积累 5 个文件时触发异步合并'compaction.min.file-num' = '5',-- 合并后文件数不超过 10 个'compaction.max.file-num' = '10',-- 合并后单文件目标大小 256MB'compaction.target-file-size' = '256MB'③ 定时 full-compaction— 兜底清理未触发自动合并的 bucket:
-- DolphinScheduler 每小时调度一次,强制合并全表CALL sys.compact('db_name.ad_traffic_detail');④ 快照清理— 及时释放已合并旧文件的存储空间:
-- 保留 5~10 个快照,过期快照关联的旧文件自动删除'snapshot.num-retained.min' = '5','snapshot.num-retained.max' = '10'业务收益:
| 指标 | 优化前 | 优化后 | 改善 |
|---|---|---|---|
| Doris 查询 P99 延迟 | 8-10s | 1.5s | 降低 80%+ |
| 单分区文件数 | ~500+ 个小文件 | ~100 个 | 减少 80% |
| 单文件平均大小 | 2-5MB | 200MB+ | 提升 ~50 倍 |
| DWS 聚合任务耗时 | ~15min | ~6min | 缩短 60% |
| HDFS NameNode RPC | 高峰期告警 | 正常水位 | 压力消除 |
3.2 流批并发写入冲突
问题描述:
流任务 7×24 持续写入 Paimon 表,批修正任务在 T+1 凌晨对同一张表做 OVERWRITE,两者并发时可能出现 snapshot 冲突(write conflict),导致某一方任务失败。
解决方案:
利用 Paimon 分区隔离,从设计上规避冲突:
时间线:7月29日白天 → 流任务写入 dt = '2026-07-29'(当天分区)7月29日凌晨 → 批任务修正 dt = '2026-07-28'(昨天分区)流写当天,批修昨天,操作不同分区,天然无 snapshot 冲突为什么有效:
- Paimon 的 snapshot 冲突发生在同一分区被多个 writer 并发提交时
- 流和批写不同分区,各自独立提交 snapshot,互不影响
3.3 数据一致性窗口期
问题描述:
从实时写入到批修正完成之间,存在一个数据不够准确的时间窗口。业务侧如果在这个窗口内做结算或对账,可能产生偏差。
解决方案:
- 数据状态标记:在表中增加
data_status字段
realtime— 实时写入,未经修正corrected— 批修正后的最终值
- 修正完成通知:批修正任务完成后发送消息通知下游系统,标记该分区数据已修正。
3.4 Schema 变更与数据回填 *
问题描述:
业务迭代中,上游经常会新增字段(如新增广告创意类型、投放策略标签等)。在流批一体架构下,一次 schema 变更会影响整条链路:
- Paimon 表需要加字段— 表结构要跟着变
- 流任务需要重启— Flink SQL 中要新增字段的处理逻辑,重启意味着中断实时写入
- 历史数据缺失新字段— 已写入的数据该字段为 null,业务查询可能出错
- 多层级联动— DWD 加了字段,DWS 的聚合逻辑可能也要改
核心矛盾:改表和重启流任务期间,Kafka 数据在持续产生但无人消费,堆积越久恢复越慢。
解决方案:
按是否需要修改 Flink SQL 逻辑,分两种情况处理:
情况一:新字段只需透传,不涉及业务逻辑变更
无需停任务,利用 Paimon 自动 Schema Evolution:
-- 数据源配置开启自动 schema 变更'schema-change.enabled' = 'true'- 上游新增字段后,Paimon 表自动加列
- 流任务使用
SELECT *或 CDC 全量同步模式,新字段自动透传写入 - 全程不停任务,零影响
情况二:新字段需要参与计算/聚合,必须改 Flink SQL
需要重启任务,通过 savepoint 机制最小化影响:
# ① 停任务并触发 savepoint(一步完成,记录 Kafka offset 和算子状态)flink stop --savepointPath hdfs:///savepoints/ <jobId># ② ALTER TABLE 加字段(如果没开自动 schema evolution)# ALTER TABLE ad_traffic_detail ADD COLUMN new_field STRING;# ③ 修改 Flink SQL(加入新字段的处理/聚合逻辑)# ④ 从 savepoint 恢复任务(从断点 offset 继续消费,不丢数据)flink run -s hdfs:///savepoints/savepoint-xxx -d new_job.jar# ⑤ 任务恢复后自动追赶堆积数据,追完后恢复实时学AI大模型的正确顺序,千万不要搞错了
🤔2026年AI风口已来!各行各业的AI渗透肉眼可见,超多公司要么转型做AI相关产品,要么高薪挖AI技术人才,机遇直接摆在眼前!
有往AI方向发展,或者本身有后端编程基础的朋友,直接冲AI大模型应用开发转岗超合适!
就算暂时不打算转岗,了解大模型、RAG、Prompt、Agent这些热门概念,能上手做简单项目,也绝对是求职加分王🔋
📝给大家整理了超全最新的AI大模型应用开发学习清单和资料,手把手帮你快速入门!👇👇
学习路线:
✅大模型基础认知—大模型核心原理、发展历程、主流模型(GPT、文心一言等)特点解析
✅核心技术模块—RAG检索增强生成、Prompt工程实战、Agent智能体开发逻辑
✅开发基础能力—Python进阶、API接口调用、大模型开发框架(LangChain等)实操
✅应用场景开发—智能问答系统、企业知识库、AIGC内容生成工具、行业定制化大模型应用
✅项目落地流程—需求拆解、技术选型、模型调优、测试上线、运维迭代
✅面试求职冲刺—岗位JD解析、简历AI项目包装、高频面试题汇总、模拟面经
以上6大模块,看似清晰好上手,实则每个部分都有扎实的核心内容需要吃透!
我把大模型的学习全流程已经整理📚好了!抓住AI时代风口,轻松解锁职业新可能,希望大家都能把握机遇,实现薪资/职业跃迁~