ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

混合计算模式实战:Hive、Spark、Flink与ClickHouse如何协同构建批流一体大数据体系

混合计算模式实战:Hive、Spark、Flink与ClickHouse如何协同构建批流一体大数据体系 做大数据这行时间长了你会发现一个特别有意思的现象以前聊分布式计算大家纠结的是用Hadoop还是Spark仿佛选定了某个框架就万事大吉。可这几年项目复杂度上来了数据量从GB涨到TB甚至PB实时性要求也从“跑一天出报表”变成了“分钟级甚至秒级响应”再指望拿一套计算框架通吃所有场景基本是痴人说梦。我自己最近在带一个网约车综合数据分析项目就遇到了这种典型的“既要又要”困境。业务方既要跑T1的离线报表又要看实时的订单监控大屏既要对海量历史订单做深度分析又要对刚产生的数据做即时清洗。折腾了一圈下来最后落地方案的核心思路就是标题里说的——混合计算模式。说白了就是让不同的计算引擎干自己最擅长的事再通过一套统一的数据链路把它们串联起来形成一个能批能流、能跑SQL也能写代码的分布式计算体系。这篇文章我不打算讲虚的就结合这个项目从零到一的搭建过程聊聊混合计算模式怎么设计、怎么选型、怎么落地以及中间踩过哪些坑。内容会覆盖Hive、Spark、Flink这些核心引擎的定位协作集群部署的资源规划还有大家容易忽略的行列级数据权限设计以及从数据清洗到可视化展示的完整链路。如果你正在做大数据的架构设计、数据开发或者准备面试大数据岗位这篇文章应该能给你一些可以直接用的东西。1. 混合计算模式的核心思路与方案选型1.1 为什么一套计算框架搞不定先说个我自己的认知转变。早几年做数仓一套Hive基本能解决90%的问题慢是慢点但胜在稳定、简单SQL写起来也顺手。后来接触了Spark感觉打开新世界的大门内存计算快得不是一星半点于是有一段时间恨不得所有任务都用Spark跑。但真正到了复杂业务场景你会发现单引擎的局限性特别明显。拿我们的网约车项目来说数据源是源源不断的订单流水、GPS轨迹、司机信息和乘客行为日志。这些数据里有结构化的MySQL业务表有半结构化的JSON日志还有非结构化的文本。如果只用离线计算那实时大屏的数据永远是昨天的业务方看监控就像在看延迟直播如果只用实时计算那需要回溯历史几个月甚至几年的数据做趋势分析时内存和成本又扛不住。这就是混合计算模式存在的意义不是用一套引擎去适配所有场景而是根据数据的时效性、计算复杂度、资源成本把任务拆解后分发给不同的计算引擎。离线批处理交给Hive和Spark实时流处理交给Flink再加上MPP数仓或OLAP引擎做交互式查询各司其职配合起来才能形成完整的大数据计算能力。1.2 Lambda架构与Kappa架构的取舍说到混合计算绕不开两个经典架构——Lambda和Kappa。Lambda架构是我最初的首选因为它成熟稳健。它的核心逻辑是双轨制批处理层用Spark或Hive处理全量数据保证结果的准确性和完整性速度层用Flink或Storm处理增量实时数据保证低延迟。最后在服务层把两套结果合并输出。这套架构在工业界应用极广几乎所有头部互联网公司的大数据平台早期都是Lambda架构。但它的痛点也很明显同一套业务逻辑要写两遍一套批处理代码一套流处理代码维护成本翻倍。而且批流结果合并时经常出现数据对不上的情况排查起来很头疼。Kappa架构则是精简版只保留速度层所有数据都走实时流处理需要回溯时用Kafka重放历史数据。好处是逻辑统一一套代码搞定批和流缺点是当数据量极大、需要回溯的时间窗口特别长时Kafka重放的成本和时效性都不可控。我在实际项目里发现Kappa对实时链路的基础设施要求很高中小团队玩转起来有难度。最终我采用的是改良版Lambda架构核心生产链路用Flink做实时处理离线批量任务用Spark SQL处理历史数据回溯和补数用Hive对外提供数据服务时再叠加一个ClickHouse或Doris做实时OLAP查询。1.3 混合计算模式的项目落地架构先画一下这个项目的整体技术栈数据采集层用Flume和Canal分别收集日志文件和MySQL的binlog消息队列用Kafka做数据缓冲和分发计算层就是混合计算的核心离线用Spark on Yarn实时用Flink on YarnAd-hoc查询用Hive on Tez存储层则使用HDFS作为数据湖底座Kudu或Iceberg管理有更新需求的数据ClickHouse承载结果表和即席查询最后通过Flask加ECharts做可视化大屏。这套结构看起来很常规没什么新奇的但真正常考验人的是计算资源的分配和任务的优先级管理。我后面会在集群部署部分详细展开。这里先讲一个关键原则离线任务和实时任务必须做资源隔离不能因为跑一个耗时数小时的批量ETL就把实时大屏的Flink任务挤死。2. 核心引擎解析Hive、Spark、Flink的分工与协作2.1 Hive离线数仓的定海神针很多人觉得Hive慢就把它淘汰了这是误解。Hive慢是因为它默认的MapReduce执行引擎慢但Hive的优势从来不是计算速度而是它作为数仓元数据管理和SQL解析层的能力。在混合计算模式里我把Hive当作数据仓库的统一入口和元数据中心。所有落库到HDFS的数据都会在Hive里建立对应的库表结构。不管是Spark任务还是Flink任务产出的结果表最终都注册到Hive的元数据服务里这样后续的即席查询、数据治理、权限管理都有了统一的依托。特别是Hive的StorageHandler机制可以让我们通过SQL直接查询HDFS上的文件、Kafka里的消息甚至关联其他存储系统中的数据。这种“一切皆表”的抽象能力是Hive在混合计算模式中立稳脚跟的根本。执行引擎上我把Hive的默认引擎从MapReduce换成了Tez。Tez能把MapReduce的多个步骤合并成一个DAG减少中间结果落盘在跑复杂SQL时性能提升明显而且不用改任何SQL代码只是修改hive-site.xml里的配置项。实测下来原来跑一个多小时的多表Join任务用Tez之后半小时内就能跑完。2.2 Spark离线批处理和轻量ETL的主力Spark的角色在混合计算模式里是“全能型选手”。我用Spark SQL处理日常的离线ETL数据清洗、维度表关联、聚合计算都是它的活。之所以选择Spark而不是继续用Hive跑ETL是因为Spark的内存计算模型在处理需要多轮迭代的复杂计算时效率比Hive高太多。举个例子我们做订单数据清洗时要对经纬度做逆地理编码判断订单所在的商圈。这个计算属于典型的“逐行处理外部字典关联”用Hive写UDF也能跑但Spark用mapPartitions算子配合广播变量性能能提升5到6倍。而且Spark的DataFrame API写起来比Hive SQL更灵活适合做日志解析、JSON字段提取这类半结构化数据处理。不过这里有个心得Spark虽然快但不是所有任务都适合用Spark跑。简单统计、大表小表关联这类常规SQL操作我仍然优先让Hive on Tez执行因为Hive的SQL优化和成本控制更成熟。Spark则专注于Hive跑不动的复杂计算和需要精细控制的任务。2.3 Flink实时计算链路的核心引擎实时这一块我用Flink来接。选Flink而不是Spark Streaming核心原因是Flink的流处理模型是原生流式的延迟更低而且它的事件时间处理、状态管理和精确一次语义支持在订单交易这类对数据准确性要求极高的场景里太重要了。我们项目里Flink干的活有两类。一类是实时数据清洗从Kafka消费原始的订单日志清洗掉字段缺失、格式异常、经纬度越界的脏数据然后写入ClickHouse供大屏实时展示。另一类是实时指标计算比如每五分钟统计一次各区域的接单率、完单率、高峰期订单密度这类窗口计算Flink跑得干净利落。Flink和Spark在混合模式下不是竞争关系而是配合关系。Flink负责把实时的数据算好写入结果存储Spark/Hive负责处理批量历史数据两者产出的结果在服务层合并最终由ClickHouse统一对外提供查询服务。2.4 ClickHouseOLAP查询的加速器ClickHouse在我这套混合计算模式里是最后临门一脚的加速器。不管是Spark产出的离线结果表还是Flink产出的实时结果表最终都会写入ClickHouse。这样才能让可视化大屏在查询千万级乃至亿级数据聚合结果时做到秒级响应。ClickHouse的列式存储和向量化执行引擎让它在聚合查询场景下的性能非常夸张。我们有一张按天的订单指标表每天新增几百万条记录做按区域、按时段的多维聚合用ClickHouse查询基本都在200毫秒以内。有一点要特别注意ClickHouse的分布式表和高可用需要提前设计好。如果你只是单机部署查起来也快但数据量一上来或者节点挂了整个可视化链路就断了。我在项目里用了三副本的ReplicatedMergeTree表引擎加上分片集群基本保证了查询服务的稳定性。3. 集群部署策略与行列级数据权限设计3.1 资源规划离线与实时的资源隔离方案集群部署这块我踩的坑最多也最想分享。刚开始做混合计算的时候我的做法是把Spark和Flink的任务都提交到同一个Yarn集群上跑想着反正都是动态分配资源没问题。结果一到业务高峰期Spark的离线ETL任务和Flink的实时任务抢占资源实时任务的延迟从几百毫秒飙升到十几秒大屏数据直接不动了。后来我调整了思路物理上共用一个Hadoop集群逻辑上用Yarn的队列机制严格隔离。具体做法是在Yarn上创建两条独立的调度队列。一条命名为offline配置集群60%的资源专门跑Spark和Hive任务另一条命名为realtime配置40%的资源专门跑Flink任务。两条队列之间做了资源上限控制就算离线队列有大量任务堆积最多也只能使用它配额内的60%资源不能侵占实时队列的份额。这里要说明的是Yarn的Capacity Scheduler配起来并不复杂关键是要设置好yarn.scheduler.capacity.root.offline.capacity和yarn.scheduler.capacity.root.realtime.capacity同时开启yarn.scheduler.capacity.root.offline.maximum-capacity做硬性上限。这样配置完成后离线任务再猛也不会拖垮实时链路。核心指标是Flink任务的Checkpoint超时次数从每天几十次降到了零。3.2 任务优先级与调度策略实战除了队列隔离任务本身的优先级管理也重要。我在项目中给Flink的实时任务设置了yarn.resource-manager.scheduler.monitor.enabledtrue配合Flink的原生支持让实时任务在队列内具有抢占优先权。另外调度策略我推荐使用Capacity Scheduler而不是Fair Scheduler。前者更适合“不同业务线明确分资源”的场景后者则适合“多个用户公平共享资源”的场景。我们这种混合计算模式显然前者更合适。还有一个容易忽略的点Yarn的节点标签功能。如果条件允许可以给部分节点打上realtime标签让Flink的TaskManager只在这部分节点上启动。这样实时计算不受离线任务的数据本地性影响网络和磁盘IO也完全是独立的。我在项目里把三台配置较高的机器单独划给了Flink使用效果比纯靠队列隔离更好。3.3 行级与列级数据权限的开源方案数据权限这块我们在项目里踩的坑也不少。最开始数仓建好后数据分析师要查订单数据开发人员要查日志数据双方都能看到全表的数据。后果就是一个实习生误操作把正在调试的表全量刷了一遍搞得下游任务全挂。这暴露出来的问题是Hive和Spark本身自带的权限控制比较粗只有库表级别的权限。但实际的业务需求是精确到行和列的。比如订单表里的乘客手机号只有客服人员和风控人员能看数据分析师只能看脱敏后的数据而商圈的订单量数据不同区域的管理人员只能看自己管辖区域内的数据。我们调研了一圈开源方案最后敲定了Apache Ranger。Ranger可以统一管理Hive、Spark、HDFS的权限策略支持基于角色的访问控制。它最强大的地方在于可以在表级别定义行过滤策略和列脱敏策略。行级权限的配置逻辑是在Ranger里创建一个策略指定某个用户组只能访问省份字段值等于“上海”的订单数据。这样即使分析师执行了全表查询实际返回的数据也只有上海地区的记录。列级权限则用到了Ranger的Masking功能把手机号字段在查询结果中自动替换成138****1234的格式。这套方案配置好以后数据安全问题解决了业务方的抱怨也没了。数据分析师该看的能看到不该看的自动脱敏权限管理从“一刀切”变成了“按需分配”。如果你也在做大数据的权限设计Ranger值得花时间研究。4. 从数据清洗到可视化混合模式下的完整链路实战4.1 基于Spark的数据清洗实现前面说了那么多架构和规划现在落到具体代码上。我用网约车订单数据清洗这一段展示一下Spark在混合计算模式里的实际用法。我们的原始订单日志以JSON格式存储在HDFS的/raw/order_log目录下一天大概几千万条。清洗的目标是过滤掉脏数据、补全缺失字段、统一时间格式然后写入Hive的订单明细表中。val spark SparkSession.builder() .appName(order-log-cleaning) .enableHiveSupport() .config(spark.sql.shuffle.partitions, 200) .getOrCreate() val rawDF spark.read.json(/raw/order_log/dt20250112) val cleanedDF rawDF .filter(col(order_id).isNotNull col(driver_id).isNotNull) .filter(col(longitude).between(73, 136) col(latitude).between(3, 54)) .withColumn(event_time, to_timestamp(col(timestamp), yyyy-MM-dd HH:mm:ss)) .withColumn(date_str, date_format(col(event_time), yyyyMMdd)) .withColumn(mobile_masked, regexp_replace(col(passenger_mobile), (\\d{3})\\d{4}(\\d{4}), $1****$2))这份代码里比较关键的是经纬度过滤这个范围是对照中国地图的经纬度边界做的能把明显异常的GPS点洗掉。手机号脱敏是通过正则实现的这其实是给数据加上了一道在计算层的保护配合前面说的Ranger权限形成了双重防护。清洗完的数据我会再用Spark SQL把它写入Hive的订单明细表中INSERT OVERWRITE TABLE dwd_order_detail PARTITION(dt20250112) SELECT order_id, driver_id, passenger_id, city_id, area_id, start_lng, start_lat, end_lng, end_lat, order_amount, event_time FROM cleaned_tmp这里用INSERT OVERWRITE而不是INSERT INTO是为了保证当天数据重跑时不会产生重复记录。这在离线任务里是一个重要习惯。4.2 Flink实时清洗与指标计算实时链路用Flink做订单数据的即时清洗和指标聚合。消费Kafka里的实时订单流做基本的清洗后按区域开一个五分钟的滚动窗口计算接单量、完单量和平均应答时长。DataStreamOrderEvent orderStream env .addSource(new FlinkKafkaConsumer(order_topic, new JSONDeserializationSchema(), kafkaProps)) .map(new OrderCleanMapFunction()) .filter(order - order.isValid()); orderStream .keyBy(order - order.getAreaId()) .window(TumblingEventTimeWindows.of(Time.minutes(5))) .aggregate(new OrderCountAggregate(), new WindowResultFunction()) .keyBy(result - result.getWindowEnd()) .process(new AreaMetricProcessFunction()) .addSink(new ClickHouseSink());窗口这块我用的是事件时间指定了Watermark策略这样即使数据到达有延迟也不会影响到窗口计算的正确性。之前图省事用过处理时间结果数据一延迟实时指标就乱套了后来才改回事件时间。Flink产出的指标结果写入ClickHouse后大屏直接查询展示整个过程端到端延迟控制在十秒以内。4.3 Flask ECharts可视化大屏的快速构建可视化层用的是Flask加ECharts这也是很多大数据项目的标配组合。Flask提供一个轻量的后端服务接ClickHouse的查询接口ECharts负责前端的图表渲染。后端Flask的主要逻辑很简单从ClickHouse查出聚合好的指标数据转成JSON返回给前端from flask import Flask, jsonify import clickhouse_driver app Flask(__name__) client clickhouse_driver.Client(hostclickhouse-server, port9000) app.route(/api/area-order-trend) def area_order_trend(): sql SELECT area_name, window_end, order_count FROM area_order_metrics WHERE window_end now() - INTERVAL 1 HOUR ORDER BY window_end rows client.execute(sql) return jsonify({ metrics: [ {area: row[0], time: row[1], count: row[2]} for row in rows ] }) if __name__ __main__: app.run(host0.0.0.0, port5000, debugFalse)前端的ECharts只做一件事定时请求接口拿数据更新图表配置。大屏的核心图表包括订单量实时曲线、区域热力图、司机在线数、完单率排行。ECharts的动态更新能力很成熟配置好setInterval定时刷新即可。这里有个经验之谈大屏的数据刷新频率不是越快越好。刚开始我把刷新设成一秒一次ClickHouse压力大不说大屏上的曲线跳得人心慌。后来改成五秒一次视觉效果和数据压力都平衡了。4.4 链路联调与数据准确性校验混合计算链路里最难的不是单个环节而是联调。批数据和流数据产出的结果经常对不上两边口径不一致的情况特别多。我们的做法是建立一套数据校验机制每天晚上用Spark跑一遍全量订单统计和当天实时链路累计的指标做交叉比对。如果差异超过阈值就触发告警人工介入排查数据是在哪个环节出了问题。这种机制绝对能干因为离线链路和实时链路天然就是两套代码字段口径、时间窗口都可能不一致。比如时间字段一个是用的下单时间一个用的支付时间数据自然就对不上。混合计算模式的落地实际上很大一部分精力是花在这种一致性核对上。想清楚这点你才能有心理准备。5. 常见问题与排查技巧实录5.1 Flink与Spark资源抢占的排查这是我在混合计算模式里遇到频率最高的一个问题。现象是实时大屏数据延迟变大Flink任务的反压持续告警但看Yarn的资源使用率并不高。排查思路是先看Flink任务的GC情况如果频繁Full GC可能是内存配置不对。再看TaskManager的线程数如果业务高峰期线程数飙升可能是并行度配置过高。最后看Yarn的队列分配确认实时队列的资源是否被离线任务悄悄侵占了。我最终定位到的问题是并行度配置不合理。默认情况下Flink会根据任务复杂度自动调整并行度但自动调整并不一定适合真实集群环境。我后来手动把算子的并行度统一配置为集群CPU核数的1.5倍左右合理设置了TaskManager的slot数量反压问题基本消失。5.2 ClickHouse查询慢的调优实录ClickHouse单表查询通常很快但Join场景一多就容易变慢。我们的可视化大屏如果涉及多表关联就会出现查询超时的情况。优化方案有几层。首先是数据模型上尽量把查询需要的数据打平成一张宽表写入ClickHouse避免查询时才做Join。其次是写分布式表时不要用Distributed表引擎做物理写入而是直接写底层的ReplicatedMergeTree本地表减少一层转发开销。最后是建表时正确设置ORDER BY键把高频过滤字段放在排序键的前面这会直接影响查询性能。实践下来这三板斧能把查询速度提升至少一个量级。尤其是宽表设计是ClickHouse高性能的精髓但要求你在数据处理阶段就把维度关联的工作做完是典型的“以空间换时间”思路。5.3 批量补数引发的数据重复处理做离线数仓不可能避免批量和补数的场景。我们的离线任务有一次因为上游数据延迟当天凌晨的ETL没跑完需要手动补跑前几天的数据分区。结果因为调度配置没有处理好补跑任务和正常的每日任务产生了重叠导致订单明细表里出现了重复记录。排查后发现我们当时使用的是INSERT INTO语义任务重复执行自然会插入重复数据。解决方案是统一改用INSERT OVERWRITE语义同时配合Hive分区的动态覆盖功能确保任务跑哪一天的分区就只覆盖哪一天的数据任务可以放心重跑。这个坑也提醒了我在混合计算模式中一套任务调度平台是不可或缺的。调度平台的作用不只是定时触发更重要的是任务依赖管理、失败重试和幂等控制这些做好了离线任务才有稳定的运行环境。6. 实操总结与后续演进建议做这个网约车混合计算项目前后花了大半年从最开始只有一套Hive到后来形成Hive、Spark、Flink、ClickHouse各司其职的混合计算体系我最大的体会是架构选型没有绝对的好和坏关键是适配你的业务场景和数据规模。混合计算模式真正解决了两个问题一是资源利用效率离线任务和实时任务不再互相干扰二是数据时效性的覆盖既能跑T1的批量分析也能看秒级的实时指标。但它的复杂度也实实在在摆在眼前运维成本、排错难度都比单引擎模式高不少。如果你准备在自己的项目里尝试混合计算我建议从小处起步。先保证离线链路稳定运行再逐步接入实时链路优先选择数据价值最高、实时性需求最强的业务场景切入。不要一开始就想搭一个面面俱到的平台那样很容易在搭建阶段就耗尽精力。最后分享一个运维上的小技巧给所有计算任务加上完善的可观测性指标。我们给Spark任务接上了Spark History Server给Flink任务接上了Prometheus监控给ClickHouse的所有慢查询做了日志收集。这套监控体系在后续排障时简直是救命稻草。混合计算模式链路长、组件多有了可观测性才能快速定位问题出在哪一环。这个方向后续还可以继续延伸。比如统一批流SQL化用Flink SQL和Spark SQL把计算逻辑统一管理再比如引入数据湖框架让批和流共用同一份数据文件还有实时数仓的分层设计也都是值得深入研究的方向。混合计算模式不是终点而是通往更高效数据架构的一个里程碑。
返回列表