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

文章详情

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

网约车大数据综合项目实战:Spark清洗+Hive分析+Flask可视化全流程

网约车大数据综合项目实战:Spark清洗+Hive分析+Flask可视化全流程 做软件方向的人一定绕不开“综合项目”这四个字。最近我把自己那份网约车大数据综合项目的笔记重新整理了一遍这条链路踩过的坑、做过的取舍、留下的代码都在这里了。项目本身规模不算大但麻雀虽小五脏俱全基于Spark的数据清洗、Hive的数据分析、Flask加ECharts的数据可视化外加架构设计、文档沉淀和最后的作品包装走通之后对软件项目的理解会完全不一样。这篇笔记写给两类人。一类是正在做软件综合项目、课程设计或者毕业设计拿到题目却不知道从何下手的同学另一类是已经跑通单个技术Demo、但没想明白怎么把Spark、Hive、Web端串成一个整体工程的人。另外如果有人在准备软考软件设计师中级你会发现综合项目里的很多知识点——数据库设计、软件工程流程、数据流图——都能在备考时派上用场两者可以互为补充。1. 综合项目的难点不在技术而在把链路串起来1.1 为什么选网约车数据做综合项目我在定项目方向时先列了几个候选电商订单分析、日志分析、网约车数据。最后选了网约车原因是它的数据字段足够“像真实业务”订单ID、乘客ID、司机ID、上下车经纬度、时间戳、金额、里程、订单状态。这些字段既有数值型、又有时间型和字符串型天然适合做数据清洗和分析演示也方便后续做可视化。更关键的是网约车数据背后有明确的业务问题高峰期订单量如何变化、哪些区域叫车需求最集中、完单率受什么因素影响、平均客单价怎么分布。有了业务问题Hive分析才不会变成“为了写SQL而写SQL”可视化也才有话可说。1.2 一个经过验证的项目拆法综合项目最忌讳的是拿到题目就写代码。我见过太多同学直接从Python脚本开始结果做到第三天发现数据格式乱得没法继续又回头重来。我的做法是先拆成四层数据获取层一份脱敏后的网约车订单数据CSV格式约60万行。数据量不大但足够让Spark跑出意义。数据清洗层用Spark做ETL处理空值、重复、异常经纬度、时间格式统一结果落地到Hive表。数据分析层用Hive SQL做指标统计把聚合结果导出到MySQL供上层查询。数据可视化层Flask提供HTTP接口ECharts在前端渲染图表。这个拆法的好处是每一层都有明确的输入和输出调试时只要盯着层与层之间的接口就行。比如清洗层输出的是Hive表分析层读的就是这张表哪一层出了问题直接在这一层排查不会牵一发动全身。提示如果你做的是毕设或者课程设计强烈建议先把分层图画出来再写代码。后面写文档、画架构图、做答辩PPT这张图直接能复用。2. 先画架构图再动手每一层都要想清楚干什么2.1 架构分层与工具选型网上有很多软件架构图的画法但综合项目不需要过度设计重点是让人一眼看懂数据从哪来、经过什么处理、最终展示成什么。我当时画的架构图分四层用到的工具也按层固定下来层次职责技术选型说明数据源层原始数据存储HDFS原始CSV按天分目录存放清洗层脏数据处理SparkPySpark分布式处理为后续分析打基础分析层指标计算Hive MySQLHive算数仓指标结果导出MySQL展示层数据可视化Flask EChartsFlask提供APIECharts渲染图表这个选型思路是学习和使用成本低、生态成熟、资料多。Spark和Hive是大数据项目的标准组合Flask是Python生态里最轻量的Web框架和数据分析天然衔接ECharts在前端图表库里颜值和上手度都排得上号。有人可能会问为什么分析结果要导到MySQL直接让Flask查Hive不行吗技术上可以但实践上不推荐。Hive查询延迟以秒计接口响应会很慢而且并发能力弱。MySQL做结果存储查询是毫秒级Flask接起来也简单。这就是分层的好处数仓负责算业务库负责查。2.2 环境准备清单工具定下来之后环境搭建又是一关。我的实际环境是单机部署一台8核16G的Linux服务器装了以下组件HadoopHDFS YARN负责分布式存储和资源调度。Spark用PySpark做数据清洗。Hive元数据存MySQL跑分析SQL。MySQL既存Hive元数据又存分析结果。Flask ECharts纯Python后端加静态页面。版本选择上我踩过一个小坑Hadoop、Spark、Hive三个版本必须互相兼容。一开始我图省事各装各的最新版结果Spark连接Hive时老是报元数据不匹配。后来统一换成Cloudera发行版的对应版本一次就过。如果你不想在环境上折腾直接用Docker搭一套Hadoop Hive Spark的镜像也没什么问题但记得把端口映射和内存限制提前配好。提示单机部署时YARN的容器内存很容易把机器拖垮。我最后是把YARN的单容器内存上限压到2GSpark执行内存也做了对应限制任务跑得慢一点但至少不会OOM。3. 基于Spark的数据清洗脏数据往往比预期多一倍3.1 原始数据里都有哪些坑数据清洗这步教科书里一句话就带过了实际动手才知道什么叫“脏”。我用PySpark读进原始CSV做了一次全景扫描问题比想象中多空值乘客ID、司机ID、订单状态都存在空值。乘客ID为空还能理解司机ID为空就得琢磨是不是系统异常。重复记录同一个order_id出现多次有的是完全重复有的是时间戳不同、其他字段一样。经纬度越界start_lat跑到了90度以上lng跑到180度以上明显是设备上报错误。时间格式不统一有的“2025-01-05 08:12:33”有的“2025/01/05 08:12:33”还有的把时分秒写成了毫秒时间戳。状态字段乱写订单状态字典是0、1、2数据里冒出一个5。这些脏数据如果不处理后面分析出的指标全是错的。比如完单率分母里混进非法状态结果直接失真。3.2 清洗脚本的完整思路我的清洗脚本大致是下面这个流程from pyspark.sql import SparkSession from pyspark.sql.functions import col, when, to_timestamp, row_number from pyspark.sql.window import Window spark SparkSession.builder \ .appName(order_etl) \ .enableHiveSupport() \ .getOrCreate() # 1. 读取原始数据开启schema推断 df spark.read.csv(hdfs://localhost:9000/raw/order/, headerTrue, inferSchemaTrue) # 2. 去重按order_id排序保留最早一条 window Window.partitionBy(order_id).orderBy(col(start_time).asc()) df df.withColumn(rn, row_number().over(window)) \ .filter(col(rn) 1) \ .drop(rn) # 3. 过滤空值关键字段 df df.filter( col(order_id).isNotNull() col(passenger_id).isNotNull() col(driver_id).isNotNull() col(status).isNotNull() ) # 4. 时间格式统一 df df.withColumn(start_time, to_timestamp(col(start_time), yyyy-MM-dd HH:mm:ss)) \ .withColumn(end_time, to_timestamp(col(end_time), yyyy-MM-dd HH:mm:ss)) # 5. 经纬度和价格合法性校验 df df.filter( (col(start_lat).between(-90, 90)) (col(start_lng).between(-180, 180)) (col(end_lat).between(-90, 90)) (col(end_lng).between(-180, 180)) (col(price) 0) ) # 6. 枚举值校验 df df.filter(col(status).isin(0, 1, 2))这里有三个细节值得说。第一去重时用窗口函数的row_number而不是dropDuplicates是因为可以显式控制保留哪一条比如保留最早下单的那条。第二时间格式化先统一成标准格式再转Timestamp避免格式混杂导致解析失败。第三经纬度校验用的是between闭区间极限值保留但超过经纬度物理范围的值一律剔除。清洗完我再做一次行数对比。原始数据60万行清洗后大概剩52万行丢了13%左右。这个比例对网约车数据来说很正常也提醒了我“数据清洗不是走过场”它真的会改变最终分析结论。3.3 清洗结果落Hive前的表设计清洗完成后要把数据写进Hive。这里的表设计直接影响后续SQL好不好写。我建了一张按天分区的外部表存储格式用ParquetCREATE EXTERNAL TABLE dwd_order_detail ( order_id STRING, passenger_id STRING, driver_id STRING, start_time TIMESTAMP, end_time TIMESTAMP, start_lat DOUBLE, start_lng DOUBLE, end_lat DOUBLE, end_lng DOUBLE, price DOUBLE, distance_km DOUBLE, status STRING ) PARTITIONED BY (dt STRING) STORED AS PARQUET LOCATION /warehouse/dwd/order_detail;写入时按dt分区df.withColumn(dt, date_format(col(start_time), yyyy-MM-dd)) \ .write \ .mode(overwrite) \ .insertInto(dwd_order_detail)选Parquet而不是普通文本是因为Parquet是列式存储在按时间范围、按状态做聚合时扫描的数据量会少很多。按天分区的好处更直接分析时只需要处理某几天的数据不用全表扫描。比如查1月6日高峰时段的数据WHERE dt 2025-01-06就够了。注意外部表和内部表的区别要想清楚。外部表删表不会删数据文件适合保存清洗后的基础数据内部表删表会连带删HDFS文件适合放中间结果。我当时把dwd层建成外部表把后面的聚合结果建成内部表逻辑上更安全。4. Hive分析模型设计决定了可视化能画什么4.1 建表与分区设计数据进了Hive接下来是分析模型的设计。这一步看着不像写SQL那么“技术”但恰恰决定了后面可视化能画出什么。我的指标表都是按“业务问题”而不是按“字段”来设计的。先问自己领导想看到哪几张图然后倒推SQL。最终我定了四个核心分析主题订单量按时段分布一天24小时哪些小时是高峰。热门口碑区域按起点经纬度聚类哪些区域叫车需求最旺。完单率与取消率不同时间段、不同区域的服务质量。客单价分布价格区间分布判断业务结构。分析结果按主题写到对应的MySQL表。比如订单时段分布CREATE TABLE hour_order_stat ( id INT AUTO_INCREMENT PRIMARY KEY, stat_date VARCHAR(20), hour_no INT, order_cnt INT );4.2 几个真正有用的分析SQLHive SQL本身不复杂但写的时候有几个习惯帮我省了很多事。第一个预约时段分组。Hive里取小时直接用hour函数INSERT OVERWRITE TABLE hour_order_stat SELECT dt, hour(start_time) AS hour_no, COUNT(*) AS order_cnt FROM dwd_order_detail WHERE status 1 GROUP BY dt, hour(start_time);这里只统计status为1的已完成订单避免把取消订单算进订单量里。如果临时要看取消率再单独跑一个带status分组的SQL。第二个去重和去极值。做客单价分析前先把价格区间限制在合理范围内比如2元到200元之间超出这个区间的要么是数据异常要么是特殊业务混在一起会把均价拉偏SELECT CASE WHEN price 10 THEN 0-10 WHEN price 20 THEN 10-20 WHEN price 30 THEN 30-50 ELSE 50 END AS price_range, COUNT(*) AS cnt FROM dwd_order_detail WHERE status 1 AND price BETWEEN 2 AND 200 GROUP BY CASE WHEN price 10 THEN 0-10 WHEN price 20 THEN 10-20 WHEN price 30 THEN 30-50 ELSE 50 END;第三热门区域我用了简单的网格聚合。把经纬度四舍五入到两位小数作为网格单元再按网格统计订单量。这个方法粗糙但对项目演示完全够用比引入GeoHash库更轻量SELECT ROUND(start_lat, 2) AS lat_grid, ROUND(start_lng, 2) AS lng_grid, COUNT(*) AS order_cnt FROM dwd_order_detail WHERE dt 2025-01-06 AND status 1 GROUP BY ROUND(start_lat, 2), ROUND(start_lng, 2) ORDER BY order_cnt DESC LIMIT 20;4.3 数据倾斜与小文件综合项目里最容易翻车的地方这两块是综合项目里最容易翻车、但教程里最不爱讲的内容。数据倾斜在真实网约车数据里非常明显热门区域订单量是冷门区域的几十倍GROUP BY热门网格时这个Reduce任务要处理的数据量远大于其他任务整个查询卡在最后一个Task上。我的处理办法是加盐把热点key先打散再聚合-- 第一步加盐打散二次聚合 SELECT lng_grid, lat_grid, SUM(partial_cnt) AS order_cnt FROM ( SELECT ROUND(start_lng, 2) AS lng_grid, ROUND(start_lat, 2) AS lat_grid, COUNT(*) AS partial_cnt FROM ( SELECT start_lng, start_lat, -- 加盐字段值为0-9随机数 (CAST(RAND() * 10 AS INT)) AS salt FROM dwd_order_detail WHERE dt 2025-01-06 AND status 1 ) t GROUP BY ROUND(start_lng, 2), ROUND(start_lat, 2), salt ) t2 GROUP BY lng_grid, lat_grid;加盐的原理很简单本来所有数据都挤在一个key上现在随机拆成10个key每个Reduce只处理十分之一的数据最后再合起来。这个技巧在综合项目里用一次就够展示功力了。小文件问题出在写入端。如果一次性写入的数据量不大Spark默认会生成很多小文件每个文件就几MB。Hive读这些文件时每个文件要起一个Map任务几十个小文件能把查询拖慢好几倍。解决办法是写入前做一次重分区df.repartition(4).write.mode(overwrite).insertInto(dwd_order_detail)把输出文件数控制在4个左右文件大小在100-200MB之间查询速度能明显提升。记住一个原则Hive查询性能不看数据量多大而看文件数量和质量。5. FlaskECharts把分析结果变成能看的图表5.1 Flask接口设计一个接口对应一个视图分析结果落在MySQL后Flask的作用是做一个薄薄的中间层把数据变成JSON喂给前端。我设计接口的原则很简单一个图表对应一个接口接口只做查询和格式化。以一个“订单时段分布”接口为例from flask import Flask, jsonify import pymysql app Flask(__name__) def query_mysql(sql, paramsNone): conn pymysql.connect( hostlocalhost, userroot, password123456, databaseanalysis, charsetutf8mb4 ) cursor conn.cursor() cursor.execute(sql, params) columns [desc[0] for desc in cursor.description] rows [dict(zip(columns, row)) for row in cursor.fetchall()] cursor.close() conn.close() return rows app.route(/api/hour_order) def hour_order(): sql SELECT hour_no, order_cnt FROM hour_order_stat WHERE stat_date%s ORDER BY hour_no rows query_mysql(sql, (2025-01-06,)) return jsonify({code: 0, data: rows}) if __name__ __main__: app.run(host0.0.0.0, port5000, debugTrue)这里有几个细节一是建立数据库连接的函数放在查询函数里不搞复杂ORM几十个指标的查询用轻量方式最省事二是返回格式统一是{code: 0, data: [...]}前端判断code再做渲染三是字符集明确指定utf8mb4不然中文内容会乱码。5.2 ECharts接数格式对齐是联调的关键ECharts的图表种类很多但接数逻辑都差不多。核心是把后端JSON里的字段映射到ECharts需要的x轴和y轴。折线图示例!DOCTYPE html html head meta charsetutf-8 title网约车订单时段分布/title script srchttps://cdn.jsdelivr.net/npm/echarts5/dist/echarts.min.js/script /head body div idhourChart stylewidth: 90%; height: 400px; margin: auto;/div script fetch(/api/hour_order) .then(res res.json()) .then(res { const hours res.data.map(i i.hour_no :00); const counts res.data.map(i i.order_cnt); const chart echarts.init(document.getElementById(hourChart)); chart.setOption({ title: { text: 2025-01-06 订单时段分布 }, tooltip: { trigger: axis }, xAxis: { type: category, data: hours }, yAxis: { type: value, name: 订单量 }, series: [{ name: 订单量, type: line, areaStyle: {}, data: counts }] }); }); /script /body /html我最开始联调时踩过一个格式坑后端返回的hour_no是数字ECharts x轴用来做类目数字当成类目时会自动排序导致下午的数据排到了上午前面。后来的处理是在前端拼成字符串比如把8变成“08:00”再放进x轴顺序问题就消失了。另一个需要注意的地方是图表类型选择。订单时段分布用折线图趋势直观区域热度用散点图或地图能看出空间聚集完单率比较用柱状图分类对比最清楚客单价构成用饼图或环形图。一张ECharts大肠图解决不了所有问题选对图表类型比堆砌交互效果重要得多。如果图表渲染的数据量比较大比如一次要画几万个点不要直接全量推到前端。ECharts渲染几千个点没问题几万个点会明显卡顿。更好的做法是在后端先做一次聚合降采样或者在前端用dataZoom只渲染可视区间的数据。6. 项目笔记要沉淀什么代码之外还有文档、部署和作品集6.1 目录结构与文档清单综合项目的最终交付物不只是能跑的代码还有一套能说服人的文档和演示。我吃了不少亏才明白代码写得再好如果别人打不开、看不懂、复现不了项目价值直接打五折。我的项目目录结构是这样组织的netcar-data-platform/ ├── README.md # 项目说明、架构图、快速启动指南 ├── docs/ │ ├── requirements.md # 需求分析和业务问题定义 │ ├── architecture.md # 架构图与选型说明 │ ├── data_dictionary.md # 数据字典 │ └── deployment.md # 部署步骤 ├── etl/ │ ├── clean.py # Spark清洗脚本 │ └── config.py # 路径、连接信息等配置 ├── analysis/ │ ├── hour_order.sql # 时段分析SQL │ ├── hotspot.sql # 热区分析SQL │ └── ... ├── web/ │ ├── app.py # Flask接口 │ ├── templates/ # HTML模板 │ └── static/ # JS、CSS └── data/ ├── raw/ # 原始数据 └── warehouse/ # 清洗后数据README里我写了两样东西架构图和快速启动命令。快速启动不是一句“按文档操作”而是逐条命令列清楚前提条件是什么、每一步大概多久、跑完怎么验证。这样无论是答辩评委、实习导师还是半年后的自己都能在30分钟内把项目跑起来。6.2 从“做完项目”到“留住成果”代码和文档之外综合项目还能往三个方向延展都是低成本高回报的事第一是数据可视化演示页打包。把Flask部署到服务器或者用gunicorn跑一个稳定的服务录一个3分钟的操作演示视频同步放到项目文档里。答辩或面试时直接放演示比口头解释高效得多。第二是软件著作权申请。如果你的项目代码是独立完成的完全可以整理成软著申请材料。我申请过几次程序鉴别材料和文档鉴别材料基本就是从项目代码和设计文档里提取再补一个软件说明书比想象中简单。软著对在校生的综测加分和评奖都有实际帮助。第三是备考软考软件设计师。中级里考的操作系统、数据库、软件工程这些基础在综合项目里都能找到对应场景。比如Hive的分区表涉及数据库设计理论Flask接口涉及软件工程里的模块化设计做项目的过程本身就是备考过程。反过来备考时学的规范化理论又能帮你审视项目里表设计是否合理。提示如果项目要求部署到多台机器或者考虑到演示环境的稳定性可以额外了解双机热备、数据库同步这类运维概念。综合项目做到最后拼的往往是这些“工程化”细节而不只是某个框架的API用得好。我个人的体会是综合项目最珍贵的不是技术栈本身而是“把一堆松散的技术点黏合成一个能回答问题系统”的能力。Spark清洗、Hive分析、Flask接口、ECharts图表任何一个单独拿出来都有海量教程但只有当它们在一个完整业务场景里协同工作时你才会真正理解每一层的边界在哪里、接口该怎么定义、数据格式为什么必须统一。这份笔记整理完后我又花了半天时间把README和架构图重新画了一遍把清洗脚本里每个判断的“为什么”注释清楚SQL文件统一加了版本号。过程很琐碎但当你三个月后再翻开这份笔记发现能在一小时内重新跑通整个项目时就会知道这些看似无关紧要的沉淀才是综合项目留给你的最大资产。
返回列表