广告投流 SAAS 产品的实时数据仓库,采用纯流式 Kappa 架构,接入抖音官方 Kafka 数据源与自有爬虫采集数据,基于 Flink + Kafka + StarRocks 构建实时分层数仓并同步给业务端。整个 ETL 链路由 Flink + 微服务组合编排,单节点峰值负载可达 2-3K QPS。
从 0 到 1 做了什么
主导全链路架构。向团队引入 Flink,负责实时数仓分层表设计(DWD / DWM / DWS)和数据端整体架构,最终与团队架构师协同敲定方案。
自研核心中间件。独立开发了两个 StarRocks 周边基础设施:StarRocksSource——实时扫描 SR 数据变更(填补 SR 缺乏 binlog 的原生短板),在高 QPS 下保证变更捕获的实时性与准确性;StarRocksRepo——基于官方源码改进"部分更新"特性,大幅降低内存占用和写入延迟。这两个中间件在计算节点 2-3K QPS 的高负载下稳定运行,对并发控制、内存管理和容错设计都有很高要求。
以 SQL 为核心的同步框架。将 SR → SR、SR → ES 的数据同步模式抽象为 SQL 驱动,同类型任务从手写代码变为配置 SQL,开发效率提升 10 倍以上。
全链路可观测和运维。自研监控报警服务覆盖消息积压、Flink 异常重启和 ETL 数据量异常;用 Python 写了一个命令行运维工具,一键打包部署、checkpoint 管理、集群状态巡检。
核心成果
- 同类型数据同步任务从 2 天缩短到 2 小时
- 两个自研中间件成为团队实时数仓的基础设施组件
- SQL 同步框架降低了一线工程师的数据开发门槛