从Lambda到Flink实时数仓:批流一体替代双轨架构的工程实践

发布时间:2026/10/8 9:07:32
从Lambda到Flink实时数仓:批流一体替代双轨架构的工程实践 先说一个我自己的背景在真正把 Flink 用进实时数仓之前我有将近三年的时间都在跟 Lambda 架构打交道。批处理一套逻辑、实时处理一套逻辑听起来是各司其职实际上做过的人都懂——两套代码、两套口径、两套调度数据对不上就焦头烂额。后来我们团队决定全面转向 Flink 构建实时数仓把这套最经典的双轨制彻底干掉才真正体会到什么叫批流一体。这篇文章我不会跟你堆架构概念而是把我们做 Flink 实时数仓时真实的选型逻辑、链路搭建过程、踩过的 JDBC 连接器相关的坑以及 Flink 与周边系统整合时的工程化细节都捋一遍。适合正在评估要不要从 Lambda 迁移到 Flink 的团队也适合那些已经决定用 Flink 但还在为分层设计、同步链路、任务运维发愁的人。1. 为什么我最终放弃了 Lambda 架构的两套代码1.1 Lambda 架构的经典分工批层与实时层各干各的Lambda 架构刚火起来那阵子几乎成了大数据领域的标配答案。它的思路很清晰把数据链路切成两条一条走批处理用 Hive 或者 Spark 定期跑全量或增量作业产出精确的历史报表另一条走实时计算用 Storm 或者早期的 Flink 处理流式数据提供秒级甚至毫秒级的指标。两条线最终在服务层汇合对外提供统一的查询接口。这套架构在最早期确实解决了既要又要的问题批处理保证数据的准确性和完整性实时处理保证数据的时效性。但你如果真的拿它来支撑一个快速迭代的业务很快就会发现不对劲。第一是两套代码维护成本极高业务指标一旦调整批处理 SQL 要改实时计算任务也要改而且改完还得保证两边结果一致。第二是两套计算引擎的语义天然有差异比如去重用户数这个指标批处理可以精确计算实时窗口却只能用近似算法兜底最后出来的数据总差那么一点。1.2 维护两套计算逻辑的真实成本到底有多高我说一个特别典型的例子。我们当时有一个核心看板需要统计当日新增付费用户数。批处理链路用的是 Hive从 HDFS 上读全量日志跑一个标准的 JOIN 和 COUNT DISTINCT实时链路用的是 Flink 1.x 的老版本通过滚动窗口计算窗口结束再输出结果。表面上看逻辑差不多实际上两个任务的代码风格、状态管理方式、甚至字段命名规范都不一样。批处理任务归数据组维护实时任务归实时组维护两边各自迭代。有一次业务方加了通过活动页进入的用户才计入这个条件两个组分别改代码上线后一看批处理口径是对的实时任务却因为窗口状态没有清理干净把前一天的数据也带进来了差了好几个百分点。这种双轨不一致的问题根本不是靠测试能完全兜住的因为数据量一大很多边界情况在测试环境根本复现不了。1.3 数据不一致的经典场景与定位过程数据对不上是最让人崩溃的环节。我记得有一次周报数据对账批处理那边显示新增用户数是 12.8 万实时这边显示 13.2 万差了四千多。我们熬夜排查最后定位到三个原因一是批处理任务存在重复读取上游文件的问题因为文件分区策略是每小时一个目录但凌晨补数任务重跑了前一天的数据又没有做幂等处理二是实时任务的动态水位线设置过于激进高峰期乱序数据被直接丢弃三是两边对新增的定义本身就不一致批处理用的是注册时间归属日实时用的是首次启动时间归属日。这三个原因没有一个是在架构设计阶段能预见到的它们都是两层架构长期并行演化的产物。那次之后我们就开始认真考虑一件事能不能只有一套计算引擎、一套逻辑代码同时兼顾批量场景和实时场景这就是后来转向 Flink 实时数仓的直接动因。2. Flink 批流一体如何改变数仓的玩法2.1 Kappa 架构的核心思想一套代码走天下Kappa 架构本质上就是对 Lambda 的减法。它的核心观点是所有数据都当成源源不断的流来处理过去的数据其实就是已经流过的数据批处理只是流处理的一种特殊形式——把数据从某个历史位置重新回放一遍而已。Kappa 架构下你只需要维护一套 Flink 作业消息队列比如 Kafka保留足够长的历史数据需要重算的时候直接从 Kafka 的指定 offset 开始消费重新跑一遍逻辑输出到新的结果表。这套思想在理论上很漂亮落地的时候最大的拦路虎不是 Flink 本身而是你喊不喊得动 Kafka 的存储成本。所以我们做架构选型的时候没有走纯 Kappa的极端路线而是用 Flink 的批流一体能力作为折中实时链路照常跑流式作业离线批量场景直接复用同一套 Flink SQL 逻辑只是把读取源从 Kafka 换成 HDFS 或者 Iceberg计算引擎和 SQL 口径完全保持一致。2.2 Flink 为什么能同时扛起批处理和流处理Flink 能在架构层替代两套引擎根本原因在于它对流和批做了统一抽象。在 Flink 眼里批就是有界流流就是无界流两者的底层执行引擎是同一个。你可以把同一个 Flink SQL 作业放在流模式下运行也可以切换到批模式下运行SQL 逻辑一行都不用改。另一个关键点是 Flink 的状态管理机制。流处理作业需要在内存里维护各种状态比如窗口聚合的中间结果、去重集合等Flink 通过 StateBackend 把状态持久化到本地 RocksDB 或 HDFS配合 Checkpoint 机制实现精确一次语义。这个能力让 Flink 在处理实时数据的时候也能做到算对而不是像早期的 Storm 那样只能靠外部存储兜底。正是这种有状态流处理的能力让 Flink 敢说实时结果可近似于离线结果也让 Kappa 架构真正有了落地的底气。2.3 从 Lambda 到 Flink 实时数仓的架构演进路线我们实际的演进路径分了三步走。第一步是存量不动、增量叠加先把新增的实时指标全部用 Flink 实现跟已有的 Lambda 链路并行跑通过一个月的对账来验证 Flink 计算的正确性。第二步是核心链路切换把最核心的十几个实时看板全部迁到 Flink 实时数仓架构下切掉老的实时计算任务保留批处理层作为每日对账的压舱石。第三步才是真正意义上的减负把离线批处理任务里能改造的部分也逐步迁移到 Flink 批模式Hive 任务只剩极少数的历史数据回溯场景。这套演进路线的好处是平滑不需要搞Big Bang式重写。我强烈建议任何想从 Lambda 迁移的团队都按这个节奏来先并行、再切换、最后收敛。一步到位的架构升级往往会在业务高峰期把你拖垮。3. 实时数仓的分层设计与核心链路搭建3.1 ODS 层实时数据接入的选型与细节ODS 层是实时数仓的地基核心任务是把业务库的变更数据和埋点日志实时采集到消息队列里供下游消费。我们当时的基础设施是这样的业务数据库是 MySQL日志数据走的是自研的 Agent 上报到 Kafka。ODS 层的第一个难点是 MySQL 的实时同步这里直接决定了后面所有链路的稳定性。我们对比过三种方案Canal 监听 binlog、Debezium 做 CDC、以及 Flink CDC 的直连模式。前两种都需要额外部署同步组件把 binlog 转换成统一的 JSON 消息写入 KafkaFlink CDC 则可以借助 Flink 自带的连接器直接解析 binlog配合 YAML 配置启动同步任务。考虑到我们已经全面拥抱 Flink最终选了 Flink CDC 方案既省掉了中间组件又能直接用 Flink SQL 处理同步逻辑一举两得。3.2 DWD 层统一口径的清洗与关联DWD 层是实时数仓里最见功夫的一层。它要解决的核心问题有两个第一把 ODS 层原始数据的脏数据清洗掉比如字段缺失、时间戳格式不一致、枚举值异常等第二把多个来源的数据做实时关联形成符合业务口径的明细宽表。清洗逻辑我用 Flink SQL 写主要用WHERE过滤脏数据、用CASE WHEN做字段标准化、用FROM_UNIXTIME统一时间格式。这部分没有太多技术难度真正的难点在实时关联——两个 Kafka Topic 之间的数据怎么保证在某个时间窗口内能正确 JOIN 到一起我的经验是两条路。第一条路是使用 Flink SQL 的INTERVAL JOIN给两个流设定一个时间窗口比如订单流和用户维表流在 5 分钟内做等值连接。这条路实现简单但窗口边界容易丢数据。第二条路是把维表数据维护到 Flink 的状态里通过Lookup Join实时关联这条路更稳但需要自己管理状态的生命周期。最终我们用的是第二条路用一条独立的维表同步作业把用户维度数据实时刷进 Flink 的 State明细流在做关联时直接查状态。3.3 DWS 层与 ADS 层宽表构建和存储引擎的选择DWS 层是聚合层按业务主题对 DWD 层明细做汇总。比如我们按用户这个主题构建实时宽表把用户的注册信息、最近一次登录时间、今日订单数、累计消费金额等指标都汇总进来供下游直接查询。ADS 层是应用层直接对接报表和各种数据产品。我们在 ADS 层的存储选型上踩过不少坑最后稳定下来的组合是ClickHouse 承担高性能查询和分析型报表Elasticsearch 承担需要全文检索或者关键词过滤的搜索场景Redis 承担大促页面的实时热数据。ClickHouse 的使用频率最高后面我会重点讲它跟 Flink 对接时遇到的 JDBC 连接器异常问题。4. MySQL 同步 ClickHouse——JDBC 连接器实战与避坑4.1 同步链路的方案选型为什么不直接全量灌数据把 MySQL 的数据同步到 ClickHouse看似简单实际上链路选型非常关键。很多人第一反应是写一个定时任务全量导数据或者直接用 ClickHouse 的 MySQL 引擎表去查线上库。但这两种方案都有硬伤全量导数据在高并发下会把 MySQL 打爆而且数据实时性完全没法保证MySQL 引擎表则会把查询压力直接透传到线上库ClickHouse 这边一跑大查询MySQL 那边就会告警。我们最终采用的链路是MySQL binlog 通过 Flink CDC 捕获经过 Flink 作业清洗转换后通过 JDBC 连接器写入 ClickHouse。这条链路的好处是实时性高、对源库压力小、可扩展性强。但这里就牵出了一个绕不开的问题——JDBC 连接器在实际使用中的异常处理。4.2 高频故障flink JDBC 连接器异常根因分析先说我们遇到的第一个高频异常Communications link failure。这个报错通常在 ClickHouse 服务端主动断开连接时出现。排查过程是这样的第一步确认 ClickHouse 是否正常运行。我们用clickhouse-client手动执行查询发现服务端没有任何异常内存和 CPU 也正常。第二步查看 Flink 任务日志发现报错集中在某个整点时间段而且每次报错后任务会自动重启。这里就看出问题了——我们的 ClickHouse 集群在整点会执行一批定时任务比如分区合并和 TTL 清理这些操作会导致部分节点短暂不可用Flink 这边如果连接池里的连接刚好被断掉就会触发异常。第三步检查连接池配置。我们最初用 Flink JDBC 连接器默认配置连接池里的连接存活时间很长ClickHouse 服务端的wait_active_tasks_timeout又设置得比较短两边参数不匹配。解决方案是调整连接池配置把sink.buffer-flush.max-rows和sink.buffer-flush.interval调小让连接不要长时间处于空闲状态同时在 ClickHouse 端的connect_timeout_with_failover和receive_timeout上做了适配。再有一个高频异常是Code: 60. DB::Exception: Table ... doesnt exist。这个报错很迷惑人因为表明明存在。后来发现原因是 Flink 作业里用了数据库前缀而 ClickHouse 连接用户权限只开放了默认数据库的访问权限导致 JDBC 驱动在解析表名时因为大小写或逃逸字符问题找不到表。解决方案有二一是在 SQL 里显式声明USE语句或完整库表路径二是重新规划 ClickHouse 账号的 GRANT 权限避免依赖默认库。4.3 稳定运行 sink 到 ClickHouse 的最佳实践经过几个月的磨难我把 Flink JDBC sink 到 ClickHouse 的最佳实践总结为四条第一连接池参数不是越大越好。Flink 的sink.buffer-flush.max-rows默认值是 100sink.buffer-flush.interval默认是 1 秒对于 ClickHouse 这种列式存储来说写入不宜太碎但也绝不能把缓冲调太大——一旦任务崩溃未 flush 的数据会全部丢失。我们最终压测出一个比较合适的组合max-rows1000interval3s单批次大小控制在 5000 行以内。第二ClickHouse 表引擎的选择直接决定写入稳定性。同步到 ClickHouse 的明细表我们统一使用ReplicatedMergeTree引擎并且配合SummingMergeTree做聚合表的预聚合。这个选择的好处是副本节点能自动同步数据某个节点宕机后 Flink 写入不受影响聚合表的重复键会被自动合并即便 Flink 端发生了重试导致重复写入最终查询结果也是正确的。第三写入失败一定要配合重试机制但重试要有上限。我们最初把 Flink 的 sink 重试次数设成了无限次结果 ClickHouse 节点故障恢复后积压的重试请求瞬间把负载打到 100%。后来改成最多重试 3 次、间隔指数退避同时配合死信 Topic 把最终写入失败的记录转发到 Kafka 备查既保证不丢数又避免雪崩。第四SQL 里的数据类型映射必须逐字段核对。MySQL 的DATETIME、DECIMAL和 ClickHouse 的DateTime64、Decimal之间需要显式 cast特别是DECIMAL精度不一致会导致写入报错。这个坑非常隐蔽因为 Flink 在 DDL 阶段往往不会报错写入阶段才会抛异常。5. Flink 与 SpringBoot 整合的工程化落地5.1 为什么要让 Flink 作业融入应用生态很多团队用 Flink 时有一个误区觉得 Flink 任务是独立部署的flink run进程跟业务应用系统没什么关系各自独立就好。这种想法初看没毛病实际运行起来问题很多。第一独立部署的 Flink 任务没有统一的应用管理入口。任务启动、停止、查看状态、修改参数全得开发手搓命令行或者依赖 YARN 的 UI运维同学接手成本很高。第二Flink 任务跟 SpringBoot 服务之间往往有共享配置的需求比如数据库连接串、Kafka 地址、业务规则参数如果两套系统各管各的配置很容易出现改了一边忘了另一边的情况。第三复杂的实时业务往往需要规则热更新或者参数动态调整这又很难脱离应用系统独立实现。我们最终定下的整合思路是把 Flink 作业作为 SpringBoot 应用内的一个可管理模块来对待而不是让两者完全解耦。SpringBoot 工程负责配置管理、资源协调、启动触发Flink 客户端负责将具体的 SQL 作业提交到 Flink 集群。5.2 SpringBoot 整合 Flink 的核心设计具体的整合方案我推荐用flink-sql-client的方式而不是直接引入 Flink 的底层 Java API 去手写各种算子。为什么呢底层 API 虽然灵活但代码量巨大维护成本高而且团队里不是所有人都具备 DataStream API 的编码能力。通过 Flink SQL 文件加配置参数的方式业务逻辑完全收敛在 SQL 模板里开发效率提升非常明显。我们的做法是SpringBoot 服务定义一个任务注册表里面存着每个 Flink 任务的 SQL 模板路径、Kafka topic、checkpoint 路径、并发度、sink 表等信息。服务启动时根据注册表自动生成提交命令调用 Flink 的 REST API 把任务提交到集群。这样运维只需要关心 SpringBoot 服务本身不用直接跟 Flink 集群交互。下面给一个典型的 SpringBoot 方法代码示例用于提交一个 Flink SQL 任务Service public class FlinkJobSubmitter { private final RestClusterClientString flinkClient; public FlinkJobSubmitter(RestClusterClientString flinkClient) { this.flinkClient flinkClient; } public String submitSqlJob(FlinkJobConfig config) throws Exception { String sql loadSqlTemplate(config.getSqlTemplatePath()); MapString, String params buildParamMap(config); // 使用 Flink SQL 解析器生成 Plan TableEnvironment tableEnv createTableEnvironment(config); tableEnv.executeSql(sql); // 提交到远程集群 JobClient jobClient tableEnv.getJobClient().orElseThrow(() - new RuntimeException(Flink job submission failed)); return jobClient.getJobID().toString(); } }代码本身不复杂但有几个细节要提醒你createTableEnvironment的时候一定要记得指定StreamExecutionEnvironment的并行度和 checkpoint 配置否则提交到集群后会使用默认值任务跑起来才发现资源分配不对。另外SQL 文件里不建议硬编码连接信息全部使用${variable}占位符由 SpringBoot 配置中心统一注入这样换环境的时候不用改动 SQL 文件。5.3 任务提交与状态管理的几个实战要点实战里我总结出三个要点每一个都是踩过坑才换来的。第一个要点是 SpringBoot 管理 Flink 任务的生命周期核心在于 checkpoint 目录不要随便改。Flink 的 checkpoint 路径一旦改变任务恢复时就找不到历史状态了等于重新起了一个新任务。我们统一用作业名 时间戳生成 checkpoint 路径配合StateProcessorFunction做周期性快照确保任务重启后能从最近的检查点恢复。第二个要点是任务提交不能做成同步阻塞操作。初始化 Flink 作业的开销非常大尤其是涉及状态恢复的作业可能阻塞几十秒甚至几分钟。如果 SpringBoot 接口同步等待作业启动完成接口超时是必然的。我们做法是异步提交作业提交请求放进线程池立即返回通过独立线程监听 Flink 的 JobStatus 变化状态变更时回调更新本地任务状态表。第三个要点是 Flink 任务日志要跟应用日志区分开。Flink 任务跑在集群上日志默认输出到 TaskManager 的本地目录排错的时候经常要登录到不同节点去翻日志非常痛苦。我们通过 Log4j 配置把 Flink 的作业日志集中发送到 Kafka再由 Logstash 收集写入 Elasticsearch配合 Kibana 做日志检索排查问题的效率提升一个量级。这个配置看着繁琐但长期收益非常大。6. 实时链路稳定运行的性能调优与资源规划6.1 并发度与分区数不匹配带来的背压隐患Flink 实时链路在运行一段时间后最容易出现的性能问题就是背压。背压的本质是下游处理能力跟不上上游数据产生速度Flink 会自动把压力往上游传导表现出来就是 Flink 监控面板里的背压百分比不断攀升。我们遇到过一个非常典型的背压场景上游 Kafka Topic 有 12 个分区Flink 作业的并发度却只配了 4 个。每个 Flink 子任务要消费 3 个 Kafka 分区的数据分区之间数据量不均匀导致某个子任务成为瓶颈整个作业的吞吐就被卡在了这个短板上。解决方式不复杂但需要对 Kafka 分区和 Flink 并发的关系有清晰认知Kafka 的单个分区是消息有序性的最小单位Flink 的单个并发度如果消费多个分区跨分区的数据顺序就无法保证而并发度超过分区数又会有部分并发度空转。最佳实践是让 Flink 并发度与 Kafka 分区数保持一致或者让并发度是分区数的整数倍。我们最终把并发度调整为 12背压从 90% 以上降到了 20% 以下。6.2 Checkpoint 调优与状态大小的控制Checkpoint 是 Flink 可靠性的基石但同时也是最容易拖垮性能的环节。我们的经验是三个参数需要重点调优execution.checkpointing.interval、execution.checkpointing.min-pause-between-checkpoints、以及 StateBackend 的选择。interval决定 checkpooint 频率太频繁会增加磁盘 I/O太稀疏则会导致恢复时间边长。我们的生产环境一般设置在 60 秒到 120 秒之间这个区间能在故障恢复速度和正常吞吐影响之间取得平衡。min-pause-between-checkpoints则用来避免连续 checkpoint 挤在一起至少保留 30 秒的间隔。状态大小的控制更有讲究。我们一开始把所有明细数据都压在 Flink 状态里导致 RocksDB 的磁盘占用直逼 100GBcheckpoint 一次就要十几分钟。后来做了两步优化一是用 Flink SQL 的TTL配置给状态设置过期时间把超过 7 天的状态自动清理掉二是把不需要精确查询的明细数据直接沉到 ClickHouseFlink 里只保留聚合结果和必要的索引状态。优化后状态大小降到 5GB 以内checkpoint 耗时降到 30 秒左右。6.3 资源规划的经验公式与常见瓶颈我这几年做 Flink 资源规划总结出一个粗线条的经验公式。一个 Flink 作业建议预留资源大致是单个 TaskManager 的 slot 数 并发度 \* (1 0.3)。多出的 0.3 部分用于处理状态恢复时的临时负载以及 Flink 框架本身的辅助线程。如果状态比较大或者有维表关联这个系数还要再往上加一点我见过最夸张的场景0.5 的额外预留都不够。常见的性能瓶颈还有三个。第一个是序列化开销。如果 DWD 层的宽表字段过多比如超过 50 个字段每次序列化和网络传输的开销会非常可观。建议尽量减少 VO 类字段数量用二进制格式传输能显著降低 CPU 开销。第二个瓶颈是 ClickHouse 的写入 merge 压力。Flink 如果以高频率小批量写入 ClickHouse会造成 ClickHouse 端 parts 数量爆炸后台 merge 任务持续高负载。解决办法是让 Flink 侧汇集更大的批次再写入或者在 ClickHouse 端设置optimize_on_insert参数控制写后合并。第三个瓶颈是 Kafka 的 rebalance。一旦某个 Flink 作业运行不稳定导致频繁重启Kafka 消费组会反复触发 rebalance整个消费链路会被暂停数秒。我们刻意控制同一消费组下作业数量不要太多并且为关键作业配置独立的消费组避免相互影响。6.4 为什么这套方案能真正取代 Lambda写到这里我想回头总结一下为什么 Flink 实时数仓这套方案能真正取代我开头说的 Lambda 架构。核心不在于 Flink 比 Hive 快多少也不在于它比 Spark Streaming 强多少而在于它让一套代码成为了现实。DWD 层的数据清洗逻辑是同一份 Flink SQL你把它跑在流模式上它就是实时链路你把它跑在批模式上它就是离线链路。两边口径天然一致再也不用对账对到深夜再也不用担心活动页用户这个条件只加在了批处理那边。当然Flink 实时数仓也不是银弹。如果你的业务以 T1 离线分析为主、实时性要求极低那 Lambda 架构的简单直接仍然是一种优势。但如果你像我一样被两套代码的双重维护折磨过被口径不一致的数据坑过那么用 Flink 构建实时数仓、用批流一体替代 Lambda 架构绝对是一条值得走的路。最后分享一个实操中的小建议迁移的节奏一定从并行对账开始让新旧两套架构跑至少一个月把所有指标的对账通过后再做切换。数据正确性这件事在实时数仓里永远需要敬畏心哪怕 Flink 的精确一次语义再强大也敌不过业务定义的变化和人的疏忽。