多彩编程 多彩编程MZPH · CODE BLOG
ARTICLE DETAIL

文章详情

深耕前端与后端开发技术的一线实战笔记与踩坑复盘。

大数据数据科学实战:从Hadoop到Spark的完整学习路线

大数据数据科学实战:从Hadoop到Spark的完整学习路线 1. 大数据数据科学到底是什么先打破炼丹和报表两个极端三年前我接了一个网约车平台脱敏数据集的综合项目数据量不算夸张五千多万笔订单但要从集群部署一路做到数据清洗、指标分析和可视化交付。做完那件事我才真正理解大数据数据科学不是很多人以为的跑一个高级模型出结果也不只是写SQL给运营提数。它是一条完整链路数据从源头进入分布式存储经过清洗变成可信的数据资产再通过数仓和计算引擎支撑分析和应用最后落到业务理解上。那时候身边有两种声音。一种认为数据科学就是机器学习炼丹天天调参另一种觉得它就是提数工具人写一堆报表。实际上两方都对了一半也都不全对。大数据领域的数据科学是用工程方法解决数据问题的能力模型只是其中一种产出报表和指标体系也是数据科学的成果。掌握它意味着你既懂底层分布式原理又懂统计思维还能用工程手段把数据价值交付出去。1.1 用网约车项目说清整条链路先说清楚网约车项目到底做了什么因为它是本文后面所有章节的统一案例。原始数据是三个CSV订单表下单时间、上车点经纬度、订单状态、司机ID、乘客ID、预估金额、实际金额等字段、司机表司机ID、城市ID、注册时间、评分、城市表城市ID、城市名。五千多万条订单分散在几十个压缩文件里单机Pandas根本读不动第一反应也不是模型而是先让人能查得动。项目分五步走把CSV导入HDFS这一步锻炼的是对分布式文件系统的理解HDFS不适合存小文件所以先做合并上传。用MapReduce写第一版清洗程序熟悉分治思想后来改用Spark重写做更复杂的脏数据过滤。清洗后的结果写入Hive分区表按城市和日期分区后续查询不用全表扫描。用Hive和Spark做指标计算比如早高峰完单率、单均等待时间、城市热力分布。最后用Flask写后端APIECharts做前端图表把指标变成运营同学能直接看的页面。整条链路下来算法模型几乎没有出场但每个环节都是数据科学的实打实内容。学习大数据数据科学的很多人卡就卡在只盯着一环要么只会写SQL要么只会装环境要么只会调sklearn。真正完整跑通一个项目之后那些分散的知识才会连成一张网。1.2 数据科学的三种产出模型、指标、系统深入做数据科学项目我把它归纳成三种核心产出模型、指标、系统。模型是预测性产出。比如预测一个乘客未来三天是否还会叫车、预测某区域未来一小时的订单增量这部分靠机器学习算法也是炼丹派最关注的。指标是描述性产出。完单率、取消率、司机在线时长分布、用户复购间隔这些指标回答了现在发生了什么哪里有问题。它不需要复杂算法但对业务理解、SQL能力、统计口径有很高要求。系统是工程性产出。把指标和模型变成接口、页面、调度任务让数据自动流转这部分需要工程能力。对于刚入门的人我的建议非常直接从指标开始以系统结束中途再看模型。这套路线最贴近实际岗位需求给人带来的成就感也最即时。下面所有章节的内容就是按这个逻辑组织起来的。2. 学习路线第一站Python、SQL和数学基础怎么在有限时间内拿下网上问大数据学习路线的人特别多我观察到一个普遍现象大家喜欢囤资料。什么《Python科学计算和数据科学应用第2版》《Python数据科学手册》全套买齐然后从NumPy底层源码开始啃啃到第二周就坚持不下去了。这完全是在用错误的顺序做正确的事。2.1 Python三件套NumPy、SciPy、Matplotlib的真实权重先聊Python科学计算三件套的定位。NumPy处理向量化数值计算SciPy提供统计和优化算法Matplotlib画图。入门阶段别管源码怎么实现的先把三件事练会取数、聚合、画图。举个例子当时我在网约车项目里做订单金额异常检测第一版就是用NumPy快速算的import numpy as np amount np.array([23.5, 28.0, 45.5, 31.2, 9999.0, 34.8, 27.6]) mean np.mean(amount) std np.std(amount) threshold mean 3 * std print(异常阈值:, threshold) # 输出: 异常阈值: 1936.14这个逻辑极其朴素均值加三倍标准差当阈值一万块的订单立刻被识别出来。后来我改用SciPy里的z-score函数做更严谨的检验from scipy import stats z_scores stats.zscore(amount) outliers np.where(np.abs(z_scores) 3)[0] print(异常索引:, outliers)这类代码写几遍就够用了不需要背诵每个参数。本阶段的目标是给你一个CSV你能用Pandas读进来用NumPy做聚合统计用Matplotlib画出趋势线。至于《Python数据科学手册》里那些高级技巧等遇到实际问题再去翻目录就好。2.2 SQL是门槛Hive再快也得先会写查询很多人有个误区觉得学了Hadoop、Spark就等于会大数据了结果一写查询就露馅。我面试人的时候必问SQL因为这最能反映候选人有没有处理过真实数据。Hive和Spark SQL本质上是把SQL翻译成分布式任务SQL写不顺后面全都不顺。网约车项目里有一个真实需求统计每个城市早高峰的完单率。完单率定义为状态为已完成的订单数除以已完成取消的订单总数。看起来简单但涉及订单表关联城市表、按时间段过滤、按城市分组聚合。这段SQL我当时写的是WITH morning_orders AS ( SELECT o.city_id, COUNT(*) AS total_cnt, SUM(CASE WHEN o.order_status completed THEN 1 ELSE 0 END) AS completed_cnt FROM orders o JOIN cities c ON o.city_id c.city_id WHERE HOUR(o.create_time) BETWEEN 7 AND 9 GROUP BY o.city_id ) SELECT c.city_name, completed_cnt / total_cnt AS completion_rate FROM morning_orders mo JOIN cities c ON mo.city_id c.city_id ORDER BY completion_rate DESC;很多初学者会卡在CASE WHEN和JOIN的配合上。这里我强烈建议先花两周专门练窗口函数row_number()、rank()、sum() over (partition by ...)、lag()、lead()。网约车场景特别适合练窗口函数比如计算每个乘客相邻两次订单的时间间隔用lag就能搞定。这些在大数据面试里出现频率极高也是Hive分析里最常用的工具。2.3 数学分阶段补不用一上来啃大部头数学基础是劝退很多人的一大原因。我的观点是别从微积分开始先补概率统计线代随用随查微积分能看懂公式就行。大数据数据科学里概率统计用的频率远超其他分支。做指标分析要知道样本与总体的关系、均值方差的稳定性做AB实验要知道显著性检验做模型要理解损失函数其实就是极大似然估计的另一种表达。我推荐用《概率论与数理统计》的经典教材前六章打底配合Python里的scipy.stats做实验验证。线性代数的核心就一个向量和矩阵怎么算。拿网约车数据举例一辆车的经纬度坐标是向量多个城市的打车热度组成矩阵计算相似度就是在做向量运算。用到什么查什么完全不需要第一遍就啃完。别把数学当拦路虎。我见过英语专业转行做数据科学的朋友数学也一般但人家Excel表格精通SQL上手极快后来一样能独立完成项目。数据科学的本质是用数据说话工具和思维可以同步学。3. Hadoop生态四层架构集群部署策略与运行原理大数据这个词已经流行很多年了但大数据架构包括四个层次这个问题至今还是面试常客。原因很简单它检验你对整个体系有没有全局观。我当时部署网约车项目集群的时候因为对架构理解不到位差点在存储层上栽跟头。3.1 大数据架构四个层次拆解通常被反复提到的大数据架构四层是数据存储层、计算引擎层、资源调度层、应用层。我用自己的理解拆一下层次代表组件负责的事情类比数据存储层HDFS、HBase、Kafka海量数据的可靠存储与传输仓库和传送带计算引擎层MapReduce、Spark、Flink把怎么算变成分布式任务工厂流水线资源调度层YARN、Kubernetes分配CPU、内存给计算任务工厂调度室应用层Hive、Sqoop、可视化平台把底层能力包装成好用的工具对外销售的成品车间我记得第一次部署集群按网上的教程把三台服务器配好然后傻乎乎地把几千万行的原始CSV用命令一条条put到HDFS。那时我不懂HDFS适合大块数据而不是海量小文件结果NameNode的元数据压力巨大文件太多直接拖慢了整个集群。这就是不理解存储层特性的后果。3.2 集群部署策略单机、伪分布、完全分布怎么选部署策略我按场景给建议单机 Hadoop拿来跑通WordCount、验证命令行脚本没问题和个人Windows环境一样简单。伪分布式一台机器模拟所有节点适合个人开发、调试MapReduce和Spark逻辑网约车项目前期我一直在伪分布式上写代码。完全分布式至少三台机子起集群适合做真实数据量测试和毕设级项目。三台虚拟机就够了每台给2核4G不用上实体服务器。生产高可用双NameNode、多ResourceManager一般公司才会用到个人阶段了解原理即可。集群配置有几个关键项。第一个是副本数测试环境设成1生产一般3第二个是内存分配Hadoop和Spark都很吃内存我当时三台机器每台4G内存NameNode默认的堆内存根本不够任务动不动就卡死后来把HADOOP_HEAPSIZE调到了1024MB才缓解。这些坑只有自己部署过才有体感。3.3 YARN资源调度的最低必要知识YARN在四层架构里属于资源调度层很多初学者直接跳过它结果Spark任务跑不起来时完全不知道看哪里。你需要理解的最小模型是提交一个Spark任务它会向YARN申请若干Executor每个Executor占多少CPU核数和内存YARN再决定放在哪台机器的Container里运行。我当时提交网约车清洗任务用的命令类似这样spark-submit \ --master yarn \ --deploy-mode client \ --executor-memory 2G \ --executor-cores 2 \ --num-executors 3 \ clean_riders.py这三个参数内存、核数、个数不是越大越好。集群总共就8G内存如果申请9G那任务会永远排队等资源。判断资源分配是否合理的标准就一条任务稳定跑完且CPU和内存利用率不过高。如果你遇到任务一直处于ACCEPTED状态多半是YARN没给你分配Container查看日志里的资源请求是排查第一步。4. 数据清洗的进化路径MapReduce第一版Spark第二版数据科学家大概有八成时间花在清洗数据上这话一点不夸张。网约车原始订单表的脏数据有多离谱司机ID有字符串和数字混用、经纬度有0值、订单金额有负数、时间字段有的精确到秒有的精确到日。清洗这一步直接决定后续分析的可信度CPU全烧在垃圾上的感觉一点都不好。4.1 为什么先写MapReduce很多人觉得MapReduce是老古董不想学。我在实际项目中先跑了一版MapReduce去重任务才真正明白分治思想的威力。任务本身很单纯把订单表里同一个乘客ID司机ID下单时间完全重复的记录去掉。单机做这个操作读一遍文件记录保留集合就能完成但到了分布式环境数据分散在多台机器上没法用一个全局的保留集合。MapReduce的思路是Map阶段把每行数据变成(key去重字段, value整行)Shuffle阶段把相同key分到同一个Reduce任务Reduce阶段保留一条即可。当时我用Java写了个最简单的版public class DedupMapper extends MapperLongWritable, Text, Text, Text { Override protected void map(LongWritable key, Text value, Context context) throws IOException, InterruptedException { String line value.toString(); String[] fields line.split(,); // 用 乘客ID_司机ID_下单时间 做key String dedupKey fields[0] _ fields[1] _ fields[3]; context.write(new Text(dedupKey), value); } }Reducer里只输出第一条。就这么几十行代码跑在伪分布式上五千多万行数据我去重花了二十五分钟。慢但给了我极强的直觉数据会被分成很多块每个块独立处理再按key归并。这个直觉后来转到Spark上一眼就能看懂RDD的分区机制。4.2 Spark清洗实录与踩坑有了MapReduce打底Spark学起来就快多了。同样是去重Spark SQL里一行groupBy就能搞定但要写出高性能版本并没有那么简单。我当时把清洗脚本改成PySpark核心逻辑是这样from pyspark.sql import SparkSession, functions as F spark SparkSession.builder \ .appName(ride_order_clean) \ .enableHiveSupport() \ .getOrCreate() df spark.read.csv(/data/raw/orders/*.csv, headerTrue, inferSchemaTrue) df_clean df \ .filter(df.order_amount.isNotNull()) \ .filter(df.order_amount 0) \ .filter(df.longitude ! 0) \ .filter(df.latitude ! 0) \ .dropDuplicates([passenger_id, driver_id, create_time]) \ .withColumn(create_time_ts, F.to_timestamp(F.col(create_time), yyyy-MM-dd HH:mm:ss)) \ .withColumn(date_partition, F.to_date(F.col(create_time_ts))) df_clean.write \ .mode(overwrite) \ .partitionBy(date_partition) \ .parquet(/data/warehouse/orders_clean)这段代码跑起来明显比MapReduce快五千多万行从二十五分钟压缩到了六分钟。但真正的问题出现在写入环节。我一开始没有加partitionBy结果Spark默认生成了几百个小文件每个才几MBHDFS上文件数量爆炸后面的Hive查询速度奇慢。解决方案是用partitionBy按日期分区并按城市做repartition控制输出分区数df_clean.repartition(6) \ .write \ .mode(overwrite) \ .partitionBy(date_partition) \ .parquet(/data/warehouse/orders_clean)这个过程踩的坑很有代表性很多初学者把Spark当成更快的单机Pandas忽略了它是分布式引擎输出文件数这个看似不起眼的细节直接影响下游性能。4.3 从十五分钟到六分钟性能到底优化了什么我把清洗脚本优化过程记录下来分享一张真实的性能对比优化动作数据量耗时说明纯MapReduce5000万行24分钟基准版本熟悉原理Spark直接读CSV5000万行9分钟只换引擎未调优加inferSchema关闭5000万行12分钟schema推断反而拖慢后来手动指定类型优化过滤顺序分区写入5000万行6分钟先过滤减少shuffle数据量写入更合理注意第二行到第三行耗时反而变长了。原因是我自作聪明关闭了inferSchema却手动写了复杂的schema解析中间有个正则表达式性能极差。后来去掉后恢复正常。性能优化不是堆配置就行一定要自己跑任务看日志。5. 分析层实战Hive做稳定指标Spark做复杂计算数据清洗完后就到了相对舒服的分析环节。这里我坚持一条原则能上Hive的绝不上Spark简单查询别杀鸡用牛刀。Hive的优势是稳定、好写、适合批处理Spark的优势是复杂计算、内存计算快。网约车项目里正好两种场景都碰到了。5.1 Hive统计指标的口径管理口径管理是数据分析最基础也最要命的环节。什么叫早高峰是7点到9点还是7点半到9点半完单率里的取消包括乘客主动取消吗不同人定义不同出来的数字能差十几个百分点。我为项目定义了一套口径表并在Hive里建了视图统一调用关键口径如下早高峰7:00-9:59晚高峰17:00-19:59完单率 已完成订单数 /已完成订单数 取消订单数取消订单包括乘客取消和司机取消不包括未接单超时关闭口径定好之后分析SQL就非常干净。比如统计各城市早高峰完单率前面章节的SQL稍微简化一下就能直接跑。5.2 Spark分析的核心结构网约车项目里最复杂的一个分析是乘客流失预警。需求是找出那些过去两周活跃、最近三天没有下过单的乘客并按城市、历史客单价分层统计。用Hive写要绕过半天我干脆用Spark DataFrame的API来算乘客历史最后一笔订单的时间from pyspark.sql import Window w Window.partitionBy(passenger_id).orderBy(F.col(create_time_ts).desc()) df_with_rank df_clean \ .withColumn(rn, F.row_number().over(w)) \ .filter(F.col(rn) 1) active_three_days_ago df_with_rank.filter( F.col(create_time_ts) F.current_timestamp() - F.expr(INTERVAL 3 DAYS) )窗口函数在PySpark中写起来比SQL变扭一些但逻辑完全相同。做完这一步把结果join城市表按客单价分桶统计就得到了流失预警的最终表。这段经历让我体会到Hive和Spark不是二选一是左膀右臂。5.3 从跑数到结论的思维转变如果纯粹把计算结果导出来给运营那还停留在跑数阶段。数据科学真正值钱的是从数字里抽出结论。我在网约车项目里有过一次经典误判运营看完单率周环比下降5%立刻觉得是司机端出问题了。我扒数据发现那周正好下暴雨所有城市晚高峰完单率都在跌去除天气因素后居然还微涨0.8%。这个分析里涉及的其实很简单先看总体分组趋势再做同环比而不是看到一个变化就开始给结论。掌握这个习惯比掌握十个算法都重要。数据科学是脑力活也是体力活绝大多数时候拼的是严谨和常识。6. 可视化交付FlaskECharts把结果做成能点能看的页面分析完的指标如果不展示给业务方价值就少了一半。网约车项目的最后一步是做一个数据可视化页面给运营同学实时查核心指标。我用FlaskECharts组合一共花了三天把页面搭起来既没上重型BI工具也没写复杂前端框架效果却相当能打。6.1 为什么选FlaskECharts这套组合的最大优势是轻和熟。Flask是Python生态最轻量的Web框架写一个API只需要不到二十行代码ECharts是百度开源的前端图表库社区生态好中文资料多遇到不会配置的图表搜一下就有答案。对比一下其他方案Streamlit和Dash胜在快但自定义能力弱想做个好看的多页报表页面反而费劲Grafana做监控图是强项但业务口径灵活的报表它帮不上写纯前端VueHighcharts太重一个人做完整个项目周期太长。对个人项目、毕业设计、企业内部小报表来说FlaskECharts是性价比最高的组合。6.2 Flask接口怎么组织后端接口我这里讲究薄API原则不要在后端做复杂聚合让SQL和Spark跑好数据前端直接拿最终结果。我的项目里长期维护的是几张报表宽表比如city_daily_report表字段有城市、日期、订单量、完单率、平均等待时长。Flask代码就变得相当简单from flask import Flask, jsonify import pandas as pd app Flask(__name__) app.route(/api/city_report) def city_report(): df pd.read_parquet(/data/report/city_daily.parquet) result df.to_dict(orientrecords) return jsonify(result) if __name__ __main__: app.run(host0.0.0.0, port5000, debugFalse)这里有个细节把DataFrame直接to_dict返回数字字段默认会变成Python原生的float和intJSON序列化没问题。但如果你遇到nan值会序列化失败前端拿到的数据直接断掉。所以报表宽表在生成时一定要清理空值用0或者Null统一填充。6.3 ECharts图表的配置痛点ECharts配置本身不难真正让人崩溃的是数据格式不匹配。我踩过最大的坑是后端返回的时间字段是2023-06-01字符串前端X轴直接用没问题但要排大小序就得转成Date对象。后来我统一在后端就转成毫秒级时间戳前端配置里用type: time省了一堆转换代码。折线图加堆叠柱状图的双轴配置我当时改了很久核心代码长这样option { tooltip: { trigger: axis }, legend: { data: [订单量, 完单率] }, xAxis: { type: time }, yAxis: [ { type: value, name: 订单量 }, { type: value, name: 完单率, axisLabel: { formatter: {value}% } } ], series: [ { name: 订单量, type: line, smooth: true, data: orderData // [[时间戳, 订单量], ...] }, { name: 完单率, type: line, smooth: true, yAxisIndex: 1, data: rateData // [[时间戳, 完单率], ...] } ] };一个小建议数据量达到上千个点时开启sampling: lttb可以大幅提升渲染性能外卖平台那种分钟级长时序图尤其需要。我当时没开页面滚动前几十秒卡得没法看开了之后丝滑很多。7. 竞赛贴脸输出MATHORCUP大数据挑战赛的三个方法论网上的热搜词里有不少大数据竞赛比如MathorCup。我自己参加过类似的数据竞赛虽然名次不算顶尖但那段经历对构建数据科学思维帮助巨大。竞赛的意义不是拿奖而是逼你在有限时间内把一个陌生数据集从零研究到交付。7.1 先约束问题空间再提模型很多队友在一起做比赛第一反应就是哪个模型牛逼用哪个XGBoost、LightGBM、深度学习一顿上结果提交分数反而难看。我的经验完全相反前两小时不做任何建模只做三件事——搞懂评测指标怎么算看数据字典跑一个最简单的baseline比如全用均值或众数预测。拿一次用户留存预测的赛题举例评测指标是AUC。我们先用全0.5概率做baselineAUC是0.5然后只用历史活跃天数一个特征做逻辑回归AUC立刻到0.6。这时候的信息量比后面调十天参还要大说明特征工程方向对了。后续的模型选择、超参调优都是在0.6到0.7之间抠分数而没有方向的话可能三天还困在0.5。7.2 代码与报告的组织比赛里最重要的隐藏分竞赛的最终交付往往不只是预测结果还要提交一份分析报告。报告的核心不是我用了随机森林而是我发现了数据里的什么规律。我当时使用的报告结构值得参考业务问题重述预测谁会在未来流失、为什么预测这个数据处理策略缺失值填充逻辑、异常值处理、时间窗口切分特征与模型讲清楚为什么做这些特征、模型选择依据结果分析分群对比、误差分析、可落地建议代码则按标准工程结构组织config.py放路径和参数data_process.py单独放清洗和特征工程model.py放训练逻辑infer.py放预测输出。比赛结束回头再看这样的组织方式让复盘变得轻松也让它变成了简历上拿得出手的项目。7.3 从比赛到毕业设计的思路迁移很多做大数据毕业设计的同学问我怎么选题我的建议和竞赛思路一模一样先把最小闭环跑通再谈创新。所谓最小闭环就是从一份真实数据集出发完成清洗-存储-分析-可视化这条链路产出一个能看能点的页面或一份完整报告。我见过太多人选了特别宏大的题目比如基于深度学习的XX预测结果数据根本跑不动最后草草应付。选题目的时候先问自己三个问题数据能不能拿到机器能不能跑动工作量能不能在截止日期前完成三个问题都过关的题目才是好题目。竞赛和毕设本质都是同一个问题在资源有限的情况下如何交付高质量的数据科学成果。8. 学习路线图与避坑清单我亲手蹚过的路最后这部分直接上干货给一份以16周为周期的可执行学习路线图以及一张我真实踩坑汇总的清单。照着做比收藏十份全网最强学习路线都管用。8.1 一份16周的学习路线表周次学习主题核心任务交付物第1-2周Python基础 NumPy/Pandas读入一个CSV做过滤分组聚合输出分析报告代码第3周SQL 窗口函数完成10道经典SQL练习题练习题笔记第4周Linux基础安装Ubuntu虚拟机熟悉命令行能用vim编辑脚本第5-6周Hadoop单机HDFS部署单机和伪分布式跑通WordCount部署文档第7周Hive基础建表、导入、写查询Hive练习表第8周数据清洗实战用MapReduce做一次去重再用Spark重写对比总结第9-10周Spark核心编程完成订单数据的复杂聚合和窗口计算Spark分析脚本第11-13周综合项目网约车数据清洗分析可视化三件套完整可展示项目第14周竞赛实战参加一个数据竞赛跑通baseline提交参赛记录第15周复盘与文档把项目写成博客或README文档技术博客第16周查漏补缺针对薄弱点系统学习问题清单这张表有两点要注意。第一第8周必须做MapReduce再转Spark直接学Spark没有底层直觉第二第14周别追求好名次跑通完整流程本身就是目的。把这张表执行完你已经具备独立完成大数据数据科学项目的完整能力。8.2 亲测有效的避坑清单我把自己踩过的坑整理成了表格每条都是真实代价坑原因解法CSV直接Pandas读取卡死单机内存不够先走分布式存储或抽样预览HDFS小文件过多put原文件没合并用Spark读后重写一次并控制文件数时间字段类型混乱CSV推理成字符串清洗阶段统一用to_timestamp转类型Hive查询慢到不可忍没做分区和分桶按日期分区大表按常用key分桶Spark任务一直排队申请内存超过集群总量按yarn资源总量反推executor参数接口返回NaN导致前端报错DataFrame空值未清理分析层统一处理空值再写入报表表图表数据多到卡顿前端渲染上万点后端聚合到分钟级前端开sampling这些坑没有一个是靠看书能提前避开的都必须亲手撞一次。所以我的建议永远是边做边学先跑通一个小项目再把项目做大。最后聊一点我的个人体会。很多人问我入门数据科学最难的是什么我觉得不是算法、不是数学、也不是工具而是动手这件事。网上下载的数据集永远有脏数据集群部署永远有莫名其妙的问题报表永远有人提新需求。这些真实任务堆一起能力才会长在你身上。我自己就是从装Hadoop开始一点点磨的磨到某个瞬间发现知识终于连成了片。那时你会觉得之前所有的弯路都值得。
返回列表