首页 / 资讯中心 / 文章详情

大数据平台选型与架构演进:从单体应用到实时计算链路的技术决策

大数据平台选型与架构演进:从单体应用到实时计算链路的技术决策 ★ FEATURED ARTICLE
简介针对创业公司不同发展阶段的大数据平台选型与演进难题这份PPT从产品验证、业务增长到成熟运营给出完整技术路线适合技术负责人、架构师及数据工程师参考。内容以实际业务监测场景为例对比了从简单Java应用搭配MySQL起步到引入Nginx、Kafka、SparkHDFS、Flume等组件的演进过程并包含RDD优化、Spark参数调优、留存用户计算、Elasticsearch实时查询等落地经验。压缩包内共有1个pptx文件文件大小仅406KB为完整演示文稿便于直接阅读和二次整理。已有113人学习下载可帮助读者快速掌握大数据平台从零到一的搭建思路与关键技术取舍。1. 大数据平台选型创业公司先想清楚“当下要什么”再谈技术栈说实话很多团队一上来就奔着 Hadoop、Spark 全家桶去最后往往死在运维和迭代速度上。这篇《大数据平台选型与演进》的 PPT 分享了我当年在一家做移动端监测和 Deep Link 唤醒服务的创业公司从零搭建数据平台的完整路径先是一个包含采集、计算、展示的 Java 单体应用硬扛了三个月再逐步演进到 Nginx Kafka Flume Spark HBase Elasticsearch CDH 的混合架构。它最有价值的不是最终那套技术栈而是里面贯穿的选型逻辑——每个阶段选什么、为什么选、什么时候该升级全部基于当时的业务数据量和团队资源来定。这份材料适合两类人一类是创业公司刚准备搭数据平台的技术负责人另一类是在成熟公司里想理解平台演进脉络、准备做技术方案汇报的工程师。它的核心关键词就两个选型和演进而这两件事的本质其实是“在正确的时间做正确的事”。2. 产品验证期的“反模式”架构为什么单体 Java 应用反而是最优解2.1 验证期的真实约束成本是第一优先级大家看 PPT 里那个阶段的技术选型原则很多人会觉得“这也太简陋了根本不是大数据平台”。但回到当时的业务场景就明白了用户只有几十个种子用户采集到的数据量根本撑不起“大数据”这三个字。同时计算的统计指标是否对用户有真正帮助完全无法确定——很可能整个功能会被市场反馈砍掉。这种情况下唯一理性的选择是尽最大可能缩小验证成本端到端跑通功能设计和实现上越简单越好不需要积累技术债务被毙掉了也不可惜。我当时做的架构就是典型的 monolithic 应用拆开来看只有三块一个移动端 SDK 的数据采集接口、一组跑在 MySQL 上的统计脚本、一张管理后台的数据展示页面。这个架构支撑了大约三个月。MySQL 在这里发挥了两个关键作用一是结构化特性让计算脚本非常好写日活、打开次数、流失用户、回流用户这些移动端常用指标一条 SQL 加上几个 group by 就能搞定二是一站式开发框架让业务修改极快因为指标计算逻辑和业务逻辑耦合在同一个应用里改一个接口、重启一次服务就完事了。这个阶段的决策其实只能用“划算”两个字来衡量所谓架构先进性在业务是否能活下来这个问题面前不重要。2.2 验证期架构的边界哪些信号告诉你必须升级这套架构的缺点也很明显它不是可水平扩展的数据量一上来MySQL 的连接数和查询性能会先出问题。我当时遇到的实际信号有两个第一采集数据时经常触发 MySQL 连接失效需要不断优化服务器端和客户端的连接参数比如 max_connections、wait_timeout、连接池的初始大小和最大上限第二当有流量稍大的种子用户进来后实时计算和离线计算的需求开始分化比如 Deep Link 短链的曝光和安装转化率需要较快看到结果而留存分析更适合离线批量跑。当这两个信号同时出现时就该启动下一阶段的架构改造了。这里顺带说一个我后来反复用到的判断方法每次准备做技术升级前先把“当前架构最痛的三件事”写下来。写不出来说明还没到升级的时候写出来了就按痛苦程度排序逐个在下一版架构里解决。PPT 里后面的演进路径本质上就是在解决“采集稳定性、存储可靠性、计算性能、运维效率”这四件事每件事对应一个新组件的引入。3. 采集与暂存链路重构Nginx 参数调优和 Kafka 的取舍逻辑3.1 采集端选 Nginx高吞吐不是说出来的是调出来的采集端选了 Nginx这个基本没争议。异步非阻塞模型天然适合高并发写入场景但很多工程问题恰恰出在默认参数上。我记得当时整理过一份需要重点调的参数清单参数默认值参考我的调整方向说明worker_processes1设为 CPU 核数Nginx worker 进程数通常等于物理核心数避免进程切换开销worker_connections512调到 10240 以上单个 worker 最大连接数和 worker_processes 相乘决定总并发能力keepalive_timeout75s调到 10~15sSDK 采集一般是短连接过长的 keepalive 会占着连接不放client_body_buffer_size8k/16k调到 32k~64k移动端上报 body 可能比默认缓冲大太小会触发磁盘临时文件写入拖慢响应access_log off开启关掉或异步高并发下 access log 写盘会成为瓶颈我们直接关了这里的核心思路是给 Nginx 一个足够的并发基座再让每个连接的生命周期尽量短保证采集接口的吞吐量能稳定抗住流量峰值。调完后建议用简单的并发脚本压一下比如用 ab 或 wrk 发几万请求观察 worker_connections 是否被耗尽、有没有报 502 或连接超时。我当时压测后调整过两轮参数才把采集端的稳定性拉起来。3.2 数据暂存区用 Kafka 而不是日志3 个决定性的理由PPT 里专门提到了一个和传统监测架构的区别没有把 Nginx 日志当数据暂存区而是直接上 Kafka。这个决策我后来在很多场合都推荐过理由有三点。第一是吞吐量和及时性比起磁盘 IOKafka 的吞吐量高得多而且生产者提供异步写入方式Nginx 采集到的数据能第一时间进入暂存区不需要等日志落盘再扫描。第二是分布式特性自带的高可用消息队列本身支持 Failover数据可以存放多份Partition 机制让写入和加载更高效而且天然能解决不同监测数据的区分问题——比如曝光、点击、安装转化这些不同类型的事件分别用不同的 Topic 隔离下游消费互不干扰。第三点是最容易被忽略但实际最值钱的数据续传问题。如果下游存储节点崩溃了纯日志方式很难知道该从哪条记录开始续传而 Kafka 客户端可以把自己的 Offset 保存在 Zookeeper 里崩溃恢复后直接从上次的 Offset 继续消费一条不多一条不少。我当时踩过日志采集的坑日志文件被 rotate 后丢了一批数据排查了整整一天才发现是续传位置算错了。换到 Kafka 之后这种问题基本绝迹了。这里补充一个参数经验Kafka 的 log.retention.hours 我会按业务需求设成 24~72 小时太短会导致凌晨补数时数据已过期太长会白白占用磁盘如果对乱序敏感可以把 max.in.flight.requests.per.connection 设成 1但吞吐会略有下降需要自己权衡。4. 离线计算与实时链路Spark 参数优化和 Elasticsearch 的定位4.1 Spark 离线计算从“能跑”到“跑得快”的 6 个调整习惯离线计算选了 Spark HDFS这部分大家都很熟重点在优化细节。PPT 里总结了 6 条心得我逐条展开说明因为这些每一条都是实际踩过的坑换来的。第一应用里要了解 RDD 的 partition 和执行中的 stage 情况。常见做法是打开 Spark UI 看每个 stage 的 task 数量和 shuffle 数据量如果出现大量极小的 task说明 partition 切得太碎需要适当增大每个 partition 的数据量。第二尽可能复用 RDD。如果同一个 RDD 要被多次 action 触发一定要做 cache并根据数据量和内存情况选择持久化策略——默认 MEMORY_ONLY 适合纯内存能放下的场景MEMORY_AND_DISK 适合数据量大于内存但还能接受落盘的场景。第三必要时候用 broadcast 和 accumulator。小维度表用 broadcast join 代替 shuffle join能省掉一大截网络传输accumulator 适合做全局计数器比如统计过滤掉的事件总数避免用 collect 把数据拉回 driver。第四资源类参数要按作业情况调不是公式化地照抄spark-submit \ --master yarn \ --deploy-mode cluster \ --num-executors 20 \ --executor-memory 8g \ --executor-cores 4 \ --conf spark.default.parallelism80 \ --conf spark.shuffle.file.buffer64k \ --conf spark.reducer.maxSizeInFlight96m \ --conf spark.yarn.executor.memoryOverhead2g \ --class com.example.MobileReportJob \ your-app.jar逻辑说明num-executors 决定总并发度executor-memory 和 executor-cores 决定单个 executor 的资源配额。一个经验值是单个 executor 的 cores 不超过 5因为超过 5 后 HDFS 读写吞吐会碰到瓶颈spark.default.parallelism 我一般设为 num-executors × executor-cores 的 2~3 倍让 task 数略多于核心数避免最后几个 task 拖慢整体。第五如果允许尝试官方推荐的 Kryo 序列化。默认 Java 序列化对象体积大、速度慢换成 Kryo 后 shuffle 数据量能降低 30%~50%对作业耗时改善非常明显。第六通过打印 GC 信息了解内存使用情况在 spark-submit 里加--conf spark.executor.extraJavaOptions-XX:PrintGCDetails -XX:PrintGCTimeStamps观察 Full GC 频率和耗时如果频繁 Full GC优先调整 executor-memory 或内存Overhead而不是盲目加大整个 executor 内存。4.2 Elasticsearch 引入的动机为什么离线平台之外还需要一个“实时查询层”这里有个容易被忽略的点Spark 离线计算再快也是分钟级到小时级的延迟。像 Deep Link 短链这种业务运营投放后需要较快看到曝光和转化情况所以还需要一条实时计算链路。我们选 Elasticsearch 而不是直接流式处理全量数据原因是它本身是一个近乎实时索引的分布式搜索引擎查询响应时间基本随节点增长线性下降而且结构化数据基于 JSON天然是半结构化的计算方式很灵活。配合查询模板化能力可以把查询条件和客户端代码解耦运维上所有功能都是 API 化的非常便于自动化管理。插件生态也够丰富支持从 HDFS、Kafka 等数据源双向导入数据。这个架构的定位不是替代 Spark而是在离线之外做一个“近实时查询层”。离线链路负责深度加工和批量聚合ES 链路负责“刚发生的事件能不能马上看到”。我一般在数据写入环节用 Kafka 直连 Logstash 再进 ES或者在代码里直接用 Bulk API 批量提交单批次建议控制在 1MB~5MB 之间太大容易触发 ES 的内存压力太小吞吐上不来。5. 存储计算选型避坑Flume、HBase 流式计算和 CDH 的边界5.1 为什么选了 Flume 而不是 Spring XD能被托管比功能丰富更重要在 Nginx 到 Kafka 这个环节我们有两种可选方案一个是 Flume另一个是 Spring XD。最终选了 Flume核心原因有两条第一使用足够简单Source、Channel、Sink 的组件化模型非常清晰各种现成的 source 和 sink 可以直接拼装第二它能被 CDH 托管而 Spring XD 只能被 Yarn 托管。CDH 的托管意味着监控、配置、告警都能统一走 Cloudera Manager这对一个只有两三个人的数据团队来说节省的运维成本是无可估量的。一个反直觉的结论是在资源有限的小团队里可运维性经常比功能丰富度更值得优先考虑。功能再强大如果没人能长期维护也是无效资产。5.2 用 HBase 做留存计算不引入 Storm靠 CAS 和行键设计解决问题留存计算这类指标回溯历史数据去做非常困难更好的思路是边采集边算。我们没有引入 Storm 或 Spark Streaming而是直接在 Flume 传输数据的链路中用一个 HBase 表做增量计算。原理是这样的HBase 表以 DeviceId 为行键维护首次访问时间和上次访问时间每条新事件进来时先判断是不是新设备——是就插入新行不是就计算距离上次访问的时间间隔更新时间戳然后通过 CAS(Compare And Set) 操作递增对应留存区间的计数器。比如一条事件tenant1|deviceId1|timeStamp1|action1如果距离上次访问 2 天就跨了 1 天说明 1 日留存用户加 1如果跨 7 天就是 7 日留存。周的留存逻辑类似只看上次访问时间是否跨周。这个方案的边界在于它只适合“单设备单次事件驱动”的留存场景如果单事件同时触发多个指标更新CAS 操作会变复杂需要自己设计行键或加列。这里有一个关键参数HBase 的hbase.client.retries.number默认是 10在网络抖动时会导致耗时翻倍我会把它调到 3hbase.client.write.buffer默认 2MB批量写入场景调到 5MB 可以显著减少 RPC 次数。rowkey 设计上尽量把 tenant 和日期前缀放前面避免把所有写入压力打在同一个 region server 上。5.3 发行版选型的避坑记录CDH 确实是运维最优解但代价要知道针对 Hadoop 发行版我们对比过 CDH、IDHIntel、HAWQPivotal、Hortonworks 四家。IDH 的特点是 HBase 提供 LOB 类型对二进制存储有帮助也优化了 Hive 让相关数据尽量落在同一 region但这些和我们的业务需求毫无关系而且 Intel 已经战略投资 ClouderaIDH 的功能会逐步移入 CDH。HAWQ 本质是 MPP 架构的数据库基于 HDFS 之上的 SQL 支持在 3 个 DataNode 的情况下上亿级别的 group by 聚合加子查询能在 10 秒左右返回适合异步近实时查询但我们没这个场景。Hortonworks 各方面和 CDH 很像但管理工具不如 CDH 强。所以最终选了 CDH因为它的管理工具提供了安装、维护、监测、预警等一系列运维能力是最适合小团队的。这里整理几个 CDH 使用中的实际踩坑记录每条都是亲历现象Cloudera Manager 里 HDFS 容量告警但实际业务数据并不多。原因默认开启了 HDFS 回收站删除的文件在.Trash里保留 N 天没清理。解决按业务需求调短fs.trash.interval比如从默认 6 小时调到 1 小时并定期跑hdfs dfs -expunge强制清空。现象Spark 作业在 CDH 上频繁报Container killed on request. Exit code: 143。原因executor 内存超过 Yarn 容器上限被 ResourceManager 主动杀掉。解决确认spark.yarn.executor.memoryOverhead是否被默认值低估按 executor 内存的 10%~20% 手动调大。现象Flume 在 CDH 托管下频繁重启。原因默认 agent 堆内存太小source 到 channel 的并发一高就 OOM。解决在 Flume Agent 的配置里把flume.child.heap.memory调到 1G 以上并关闭无关的 source 和 sink。现象Kafka Offset 存储压力大Zookeeper 线程飙高。原因客户端默认每 1 秒自动提交一次 offset。解决调大auto.commit.interval.ms到 10 秒并改成手动提交降低 ZK 写压力。现象HAWQ 试用时发现聚合查询确实快但运维成本高。原因它是一个独立 MPP 数据库需要单独维护一套元数据和资源组和 HDFS 的整合不像 CDH 那样开箱即用。解决没有实际引入只在文档里留有评估记录结论是“能力很好但不匹配当前阶段”。6. 架构验证与技术债清理一套可持续演进的健康度检查方法不管选型多合理演进多次后架构一定会有隐性债务。我后来养成一个习惯每季度强制做一次“架构健康度检查”用的是从这套选型经验沉淀出来的标准动作。最核心的动作是压测验证先用wrk或ab对 Nginx 采集接口做一次基线压测记录吞吐量和 P99 延迟然后停掉 HBase 的 region server 节点观察 Kafka 消费者是否按 Offset 续传、数据有没有丢接着重跑一次离线聚合任务对比 Spark UI 里每个 stage 的 shuffle 数据量和耗时有没有异常增长。这套动作的优点是它不测“系统能不能跑”而是专门测“故障时数据会不会丢、恢复后是否一致、性能是否随时间退化”。从那次以后我每接手一个数据平台第一件事就是跑这三项检查。另一件值得做的事是给每个组件写清“选型前提条件”。比如 Kafka 需要 Zookeeper 保持奇数节点ES 需要保留至少一个副本HBase 的 region split 策略依赖 rowkey 设计如果业务变了前提不成立那对应的组件也应该重新审视。而不是因为当年选了它就一直默认它是对的。这种“写清楚为什么选它”的习惯对我们这种需要持续演进的小团队尤其有用因为新人接手时能直接看到当时的决策边界而不是把现状当作理所当然。从那次项目之后我每次做数据平台选型都会先回答三个问题当前阶段最痛的点是什么、为了这个点愿意付出多少运维成本、这个方案在半年后还能不能继续迭代。这三个问题的答案比任何技术栈的优劣对比都重要。希望这份演进思路能帮你在选型时少走一段弯路。本文还有配套的精品资源点击获取
阅读完成 · 觉得有帮助?
咨询建站