ARTICLE DETAIL

资讯详情

深耕网站视觉设计与运营推广的一线实战洞察。

大数据复杂场景下数据科学实战:从数据清洗到治理全链路掌控

大数据复杂场景下数据科学实战:从数据清洗到治理全链路掌控 先从一个我最近反复被问到的事情说起。很多朋友看到数据科学四个字第一反应是机器学习模型、神经网络、调参炼丹。但真正到了大数据领域的复杂场景里比如热搜里反复出现的网约车大数据综合项目、校园大数据分析、MathorCup大数据挑战赛你会发现事情完全不是这么回事。你拿到手的是一份几百GB的CSV里面订单号和手机号混在一起时间格式乱成麻花经纬度有落在海里的金额还有负的。这时候你脑子里那些模型根本派不上用场真正决定项目成败的是你有没有一套能从头到尾把数据捋顺的系统能力。这篇文章我想从实战视角聊聊数据科学在大数据复杂场景下到底该怎么打。内容会围绕数据清洗、分析链路、可视化落地、集群部署、数据治理这几条主线展开穿插我在网约车数据项目、校园数据项目里真实踩过的坑。不管你是准备大数据毕业设计、参加竞赛还是刚入行做数据开发这篇文章的目标是让你少走弯路知道真正的功夫该下在哪里。1. 大数据复杂场景的真实面貌不是模型难是数据难1.1 热搜词背后暴露的共同痛点如果你去翻大数据相关的热搜词会发现一个很有意思的现象大家搜索最多的不是神经网络也不是深度学习而是数据清洗数据分析数据可视化大数据集群部署策略数据质量检查框架。这说明什么说明绝大多数人在真实项目里卡住的环节恰恰是那些听起来最没技术含量的活。网约车大数据综合项目里基于Spark的数据清洗、基于MapReduce的数据清洗、数据分析Hive、数据可视化FlaskECharts这些词被反复搜索本质上是因为它们构成了一个完整项目的地基和骨架。算法模型只是最后锦上添花的一步而前面这些环节任何一个出问题整个项目都会崩。我见过太多人把精力花在调一个看似高深的模型上结果连数据里有多少重复订单、多少空值都没搞清楚。最后模型跑出来的结果根本没法解释因为数据本身是脏的、乱的、不可信的。在大数据场景里数据科学的第一性问题永远是数据质量而不是模型复杂度。1.2 复杂度的四个来源规模、维度、质量、时效大数据场景之所以复杂我总结下来主要来自四个维度你可以对照自己手里的项目感受一下第一是规模。单机内存装不下这是最直接的冲击。一张订单表几亿行过去在Excel里拉个透视表的操作习惯全部失效。你连这个字段到底有多少种取值这种最基础的问题都得通过分布式计算引擎才能回答。操作习惯变了思维方式也得变。第二是维度。数据来源杂接口日志、业务库、埋点、第三方数据每张表都有自己的字段口径和命名习惯。用户ID在这张表叫user_id在那张表叫uid合并的时候口径对不上分析结论就全是错的。第三是质量。这个最磨人。缺失值、重复值、异常值、格式不统一、前后不一致这些问题不是零星出现的而是成片存在的。你以为是脏数据是少数实际上在真实场景里脏数据往往是常态干净数据才是少数。第四是时效。业务方不会等你慢慢跑数。今天的数据最好今天出结果最晚明天早上要看到报表。这逼着你必须在存储结构、计算引擎、调度策略上做取舍而不是只写一条SQL拉倒。1.3 数据科学家的角色早已不只是建模在这样复杂的场景里数据科学家的角色发生了很大的变化。你不再是那个只负责训练模型的人你更像是一个数据管家要懂数据怎么采集、怎么清洗、怎么建模、怎么调度还要懂怎么把分析结果用业务听得懂的话解释出来。这也是我想强调的一个观点在大数据领域数据科学能力的核心不是算法储备而是对整个数据链路的掌控力。你不需要会写最复杂的模型但你必须知道一条数据从产生到被业务使用中间要经过哪些环节每个环节可能出什么问题出了问题怎么排查。这种能力恰恰是学校和培训班最不教你、但项目里最要命的东西。2. 数据清洗复杂场景的第一主战场2.1 从网约车订单数据看脏数据的常见形态先看一个典型的网约车订单数据集。字段大概有三十多个订单号、乘客ID、司机ID、上车时间、下车时间、上车经纬度、下车经纬度、里程、金额、状态等等。表面上看结构完整但真正开始分析的时候你会陆续发现这些问题重复记录同一个订单号出现两遍甚至三遍。多是由于系统重试写入或者接口重复推送导致的。时间格式混乱有的存的是2019-08-01 12:30:00有的是2019/08/01 12:30还有的是2019年8月1日如果不统一排序和区间统计全是错的。坐标漂移GPS上报偶发异常经纬度跑到城市范围之外甚至出现(0,0)这种坐标。做热力图的时候这些漂移点会直接把渲染范围撑爆。金额异常负数、0元、或者上千元的极端值不一定是错误数据但必须在分析时单独处理否则平均客单价会被拉得离谱。关键字段缺失乘客ID、司机ID为空导致无法做关联分析。这些脏数据形态几乎在每个真实项目里都会碰到。校园大数据项目也一样学生选课记录里有重复选课、成绩有负数、时间格式不统一。所以说会识别脏数据的常见形态是数据科学在大数据场景下的基本功。2.2 清洗方案的选型逻辑MapReduce还是Spark确定了要清洗的问题接下来是选工具。热搜词里基于MapReduce的数据清洗和基于Spark的数据清洗同时出现很多朋友就纠结了到底学哪个、用哪个我的观点很直接如果是新项目、新写的代码直接用Spark如果是为了理解分布式计算原理MapReduce值得学但别用它写复杂ETL。MapReduce的问题是一个简单的过滤去重逻辑得写Map、Reduce两个类还要处理序列化、Partitioner这些细节中间结果大量落盘调试一次要等很久。而Spark的DataFrame API把这些东西全部封装掉一个链式调用就完成了多步清洗逻辑而且基于内存计算跑起来快得多。打个比方MapReduce像是手写汇编Spark像用高级语言写业务。汇编能让你理解CPU是怎么工作的但没人会用汇编去写一个Web系统。学习阶段理解原理用MapReduce没问题工程落地直接用Spark。下面是我在网约车项目里常用的一套Spark清洗骨架你可以直接套用自己的数据from pyspark.sql import SparkSession from pyspark.sql.functions import col, when, to_date, regexp_replace spark SparkSession.builder.appName(order_clean).getOrCreate() # 读取原始数据 df spark.read.option(header, True).csv(hdfs:///data/raw/order.csv) # 第一步去重同一订单号保留最新一条 df df.dropDuplicates([order_id]) # 第二步统一时间格式 df df.withColumn( pickup_time, when(col(pickup_time).contains(/), to_date(regexp_replace(col(pickup_time), /, -), yyyy-MM-dd) ).otherwise(to_date(col(pickup_time), yyyy-MM-dd HH:mm:ss)) ) # 第三步过滤异常坐标 df df.filter( col(pickup_lat).between(22.5, 25.5) col(pickup_lng).between(105.0, 115.0) ) # 第四步过滤异常金额 df df.filter(col(amount) 0) # 第五步缺失值填充 df df.na.fill({passenger_id: unknown, driver_id: unknown}) # 写入清洗后的分区表 df.write.mode(overwrite).partitionBy(dt).parquet(hdfs:///data/clean/order)注意几个细节。去重前先想清楚按什么去重、保留哪一条网约车场景里同一订单因为多次推送会有多条记录但业务上应该保留最后状态所以要先按时间倒序排再去重或者直接用dropDuplicates配合排序后的DataFrame。时间格式化一定要统一到一种标准否则后续按小时聚合会错得莫名其妙。坐标过滤的边界值要按城市实际范围设定别用全国范围的经纬度否则还是会有漂移点漏网。2.3 一个必踩的坑导出数据变成科学计数法清洗完数据往往需要把结果导出看一眼。这时候有个经典坑相信很多人经历过用DBeaver导出数据打开CSV发现订单号、手机号、身份证全变成了科学计数法后面的位数变成了0看起来像数据丢了。先别慌数据本身没丢是Excel的显示问题。Excel对超过11位的数字自动转科学计数法而且精度只有15位超过的部分全部归零。你在DBeaver里看到的是完整的订单号导出后类型变成了数值Excel就自作主张帮你截断了。但这不是小事如果你直接把导出的CSV拿去关联分析或者导入另一个数据库就会因为精度丢失而匹配失败。解决的办法有几个在SQL里直接对ID类字段做CAST(order_id AS STRING)让导出结果本身就是文本。DBeaver导出时在导出设置里勾选将字符串括在引号中这样CSV里的ID会带着引号Excel会识别为文本。用文本编辑器比如Notepad、VS Code打开CSV确认原始值别直接用Excel看。导入目标数据库时确保目标字段类型是VARCHAR而不是BIGINT或DOUBLE。这个小坑背后是一个通用的教训跨系统搬运数据时类型转换是事故高发区。你以为你在看数据其实你看到的是工具帮你处理过的数据。做数据科学你得养成用原始格式确认数据的习惯。3. 分析链路Hive与Spark SQL的分工协作3.1 数仓四层架构中分析发生在哪清洗好的数据会进入数仓。传统数仓一般分四层ODS原始数据层、DWD明细数据层、DWS汇总数据层、ADS应用数据层。热搜词里大数据架构包括四个层次说的也是类似的东西只不过不同公司叫法略有差异。在这个体系里数据科学分析通常发生在DWD到ADS之间。DWD是清洗后的明细数据字段完整、口径统一DWS是按业务主题汇总后的数据比如订单日汇总表用户维度表ADS是面向具体应用的结果表比如每小时订单量趋势Top10司机榜单。理解分层的意义在于它让不同团队可以在同一个数据体系里协同时不会互相踩脚。你做的分析不应该直接读ODS的原始垃圾表而是应该基于DWD的干净数据再产出ADS的结果表。这个习惯一旦养成了后续接任何可视化系统、报表系统都会很顺畅。3.2 一套可复用的分析SQL拆解在网约车项目里一个典型的分析需求是统计每个时段按小时的订单量、完单率、平均客单价。对应到Hive或Spark SQL上核心SQL大概是这样的SELECT hour(pickup_time) AS hour_no, COUNT(*) AS order_cnt, SUM(CASE WHEN status completed THEN 1 ELSE 0 END) / COUNT(*) AS complete_rate, AVG(amount) AS avg_amount FROM dwd_order_info WHERE dt 2024-01-01 GROUP BY hour(pickup_time) ORDER BY hour_no;这段SQL看着不难但里面有几个容易被忽略的细节第一hour()函数提取小时的前提是pickup_time必须是标准Timestamp类型如果你在清洗阶段没统一时间格式这里直接就报错或者返回NULL。第二AVG(amount)在数据里有极端值(比如几千元的豪华车订单)时平均值会被拉高。更稳妥的做法是同时算PERCENTILE_APPROX(amount, 0.5)作为中位数或者做缩尾处理。第三WHERE dt 2024-01-01利用分区裁剪只扫描当天的分区。如果没建分区表或者查询条件不带分区字段几亿行数据全表扫描跑一次要十几分钟体验极差。大数据SQL的第一原则就是尽量把分区条件带上。3.3 大数据SQL面试题到底想考什么很多准备面试的朋友会狂刷大数据SQL面试题比如求每个用户当天第一单的时间统计订单量Top10的司机找出连续三天有订单的用户。这些题刷起来很上头但你要理解它背后真正考察的东西。以每个用户当天第一单时间为例核心是窗口函数SELECT passenger_id, pickup_time FROM ( SELECT passenger_id, pickup_time, ROW_NUMBER() OVER (PARTITION BY passenger_id, dt ORDER BY pickup_time) AS rn FROM dwd_order_info ) t WHERE rn 1;ROW_NUMBER() OVER (PARTITION BY ... ORDER BY ...)这一句考的是你有没有真正理解分组内排序这种逻辑以及窗口函数分布在不同计算引擎里的执行计划差异。面试官很难通过背题看出你的水平但他能从你写出的SQL里判断你有没有真实处理过大数据的经验——比如你有没有主动加上分区条件有没有考虑数据倾斜有没有注意到某个字段可能存在NULL。我的建议是刷题不如刷项目。你把一个真实项目里所有的分析SQL自己从头写一遍再设计几个如果数据量和数据质量变差怎么办的问题比刷一百道题都管用。面试官问SQL的目的不是为了考语法而是想确认你能不能在大规模、低质量的数据环境里把活干成。4. 可视化落地FlaskECharts打通最后一公里4.1 为什么可视化是数据科学的最后一公里分析做完SQL汇报那页PPT写得再漂亮业务方也感知不强。真正让他们哇出来的是可视化大屏或者交互图表。热搜词里网约车大数据综合项目——数据可视化flaskecharts校园大数据—数据可视化被反复搜索说明大家都意识到可视化是项目交付中不可或缺的一环。这里想多说一句可视化的本质不是画图是降低理解门槛。好的可视化能让人一眼看出哪个时段订单最多哪个区域叫车需求最旺盛而不是丢给业务方一张几百行的Excel表格让他们自己看。我做项目的习惯是先做一版最核心的图表再根据实际使用反馈逐步增加交互而不是一上来就把所有图表堆在页面上。4.2 后端接口怎么写才不拖后腿Flask作为轻量级Web框架用来做数据可视化的后端非常合适。核心思路是后端从数仓结果表读取数据封装成JSON接口前端用ECharts发起请求并渲染图表。后端接口的骨架我用网约车项目的每小时订单量接口举例from flask import Flask, jsonify from pyspark.sql import SparkSession app Flask(__name__) spark SparkSession.builder.appName(dashboard).enableHiveSupport().getOrCreate() app.route(/api/hourly_orders) def hourly_orders(): df spark.sql( SELECT hour_no, order_cnt FROM ads_hourly_order_stats WHERE dt 2024-01-01 ORDER BY hour_no ) rows df.collect() hours [r[hour_no] for r in rows] orders [r[order_cnt] for r in rows] return jsonify({hours: hours, orders: orders}) if __name__ __main__: app.run(host0.0.0.0, port5000, debugFalse)注意这里有几个实际坑SparkSession不能每个请求都创建。它的启动和初始化非常重几十秒甚至几分钟。正确做法是全局初始化一次后面所有接口复用。如果内存紧张甚至可以考虑把SparkSQL查询结果缓存成临时文件再让Flask直接读文件避免长时间占用计算资源。接口返回的数据量要控制。图表的点位数一般几百个以内没问题但如果一次返回几十万行前端渲染会卡死。在接口层面做聚合或采样是后端工程师的自觉。debugFalse是必须的。Flask的debug模式会在代码变更时自动重启但在生产环境会有安全风险而且会和SparkSession的初始化流程冲突导致重复创建Session。4.3 前端图表联调中的真实坑前端用ECharts画折线图的代码不算复杂fetch(/api/hourly_orders) .then(response response.json()) .then(data { const chart echarts.init(document.getElementById(chart)); chart.setOption({ title: { text: 每小时订单量趋势 }, tooltip: { trigger: axis }, xAxis: { type: category, data: data.hours }, yAxis: { type: value }, series: [{ name: 订单量, type: line, data: data.orders, smooth: true }] }); }) .catch(error console.error(加载数据失败:, error));联调阶段最容易遇到的是跨域问题。前端项目跑在3000端口Flask后端跑在5000端口浏览器会拦截跨域请求。解决方式有几种最简单的是在Flask里加CORS头或者用Nginx做反向代理让前后端同源。另一个我踩过的坑是时区问题。服务器部署在国内但如果测试时连的是本地Hive集群时间字段会有UTC和北京时间的差异最终图表里会出现曲线整体平移8小时这种诡异现象。排查思路是核对数仓里存储的时间时区是什么再确认Flask接口返回的JSON里时间字符串本身是否带时区标识最后看浏览器渲染时用的时区。这三者如果不一致图表就漂了。还有一个很隐蔽的坑图表数据存在NULL值。订单量某个小时可能没数据Spark SQL返回的JSON该字段是null。ECharts默认会在null处断开线条如果业务方不理解会以为数据出了问题。更好的做法是后端默认填充为0或者在SQL里用COALESCE(order_cnt, 0)处理。5. 集群部署策略小团队撑起大数据项目的现实方案5.1 组件选型和硬件预估说到部署很多人上来就想搭建一套CDH或者HDP全家桶各种组件装了一大堆最后发现根本没有数据量去支撑这些组件。我自己的经验是小团队、小项目优先保证核心链路能跑通别贪全。一套我反复推荐的中小型大数据部署方案是组件作用选型建议HDFS分布式存储3节点起步副本数设为2Spark计算引擎Standalone模式即可不必强上YARNHive数仓SQL用MySQL作为MetastoreFlink实时计算非必须可根据是否需要实时需求加装Flink CDC增量同步如果有业务库实时同步需求再加硬件上3台8核16G内存、1T硬盘的服务器是性价比比较高的起点。如果你的数据量在100G到1T之间这个配置够用了。单机伪分布式模式也能跑教学项目但至少要16G内存否则Spark任务一跑就OOM。磁盘方面如果只是学习单块大容量SSD就够了如果数据量上百G建议至少2块盘做RAID。这里特别提醒一下集群部署不要一上来就追求三节点高可用。先在一台机器上跑通全流程再用克隆或新装的方式扩充节点。很多初学者第一次就搭三节点结果网络配置、时钟同步、SSH免密这些问题连环爆排查到最后连自己都懵了。一步一个脚印先把一条链路走通比什么都强。5.2 部署中容易被忽略的四个细节部署文档网上很多我重点说几个容易被忽略、但坑人无数的细节第一内存分配要预留系统余量。Spark的Executor内存不是越大越好。假设一台机器16G内存系统本身要占4GHDFS的DataNode要占2G剩下的10G才是给执行任务的。如果spark.executor.memory12G跑任务时系统直接卡死。我一直遵循的规则是物理内存减去4~6G系统保留再在剩余范围内分配。第二Hive和Spark共用Metastore时配置必须一致。最典型的坑是Hive里建好的表SparkSQL里show tables看不到。原因是两者连的Metastore地址不一致。解决办法是在hive-site.xml、spark-defaults.conf中明确指定同一个javax.jdo.option.ConnectionURL指向同一个MySQL库。第三时钟同步和免密登录要提前配置。多节点集群里如果节点时间差太多任务提交后会出现执行计划分配异常。建议直接用ntpdate同步或者用chrony服务。SSH免密不只是为了登录方便Spark Standalone模式启动Worker、提交任务时都依赖SSH。第四临时目录要及时清理。Spark任务和MapReduce任务会在/tmp下产生大量临时文件HDFS的/tmp目录也容易被写满。我见过一个项目因为/tmp满了所有任务都失败最后排查了半天才发现是磁盘满了。建议在部署时配置定时清理脚本比如每周清一次超过三天的临时文件。6. 数据治理与质量检查让分析结论站得住脚6.1 轻量级质量检查框架应该包含什么数据治理听上去像大厂才需要做的事其实不然。哪怕你只是做一个毕业设计质量检查这套思路也值得贯穿始终。不然你的分析结果可能因为一个字段的数据错误得出完全相反的结论到时候发现问题就得从头返工。我做项目时一般会在写完清洗流程后加一个自动化的质量检查模块。它至少包含五类检查完整性必填字段的NULL率是否异常。比如订单表的order_id不应该有NULL如果NULL率超过0.1%就要检查清洗规则是否漏了。唯一性主键字段有没有重复。用COUNT(*)和COUNT(DISTINCT order_id)对比不一致就是有重复或NULL。有效性字段值域是否符合业务规则。比如金额是否在0到10000之间、状态字段是否在预期的枚举值内。及时性数据时间是否满足SLA。比如今天的分区数据必须在次日凌晨2点前产出监控脚本发现晚于这个时间就要告警。一致性同一指标在不同表中是否对得上。比如ADS层的总订单量和DWS层的汇总值如果对不上说明链路里某一步算错了。对应的质量检查SQL长这样SELECT COUNT(*) AS total_cnt, COUNT(DISTINCT order_id) AS distinct_order_cnt, SUM(CASE WHEN order_id IS NULL THEN 1 ELSE 0 END) AS null_order_cnt, SUM(CASE WHEN amount 0 THEN 1 ELSE 0 END) AS invalid_amount_cnt FROM dwd_order_info WHERE dt 2024-01-01;跑完之后把几个指标和前一天对比如果差异超过阈值比如今天的NULL率突然从0.1%涨到5%就触发告警。告警方式很简单写个Python脚本发现问题就发邮件或者钉钉机器人通知。6.2 治理不用一步到位但要从小做起完整的数据治理体系还包括数据血缘、元数据管理、权限管控、数据字典等等。这些东西对小项目来说太重了但我建议从三个最轻量的动作开始第一建一份数据字典。把每个表的字段含义、类型、枚举值、负责人写在一份文档里。听起来很土但项目做到一半你一定会感谢自己当时写清楚了。第二记录任务依赖关系。哪怕只是用一张图表示Hive清洗任务→分析任务→可视化接口的顺序也能在某个任务失败时快速定位影响范围。有条件的话用调度平台的任务流功能没条件就写清楚依赖清单。第三区分权限和角色。开发人员有读写权限业务查看人员只有读权限。小项目可能用不上特别复杂的权限系统但至少在数据库账号上做区分不要所有人都用同一个root账号操作。治理这件事的价值不是让你去通过某个认证而是让你的数据科学结论经得起追问。业务方问你这个数准不准、怎么来的时你能用清晰的口径和文档回答而不是支支吾吾说大概跑的没问题吧。这种可信度才是数据科学家最值钱的东西。说到底数据科学在大数据领域里拼的从来不是模型有多高级而是你对每一个环节的掌控力。我自己的习惯是每次接手一个新数据集先花半天时间做数据体检——看总量、看唯一性、看空值、看分布、看几个典型字段的样例再花一天时间把清洗流程和数仓模型定下来最后才开始写分析代码。前期多投入一小时后期能省下十小时。每个项目都像是把一堆乱麻捋成线头的过程。你手里的工具可以换Spark也好、Hive也好、Flask也好它们都是手段。真正重要的是你始终记得数据科学的目的是从混乱中建立秩序从噪音中提取信号。下一次再面对一份几百GB的脏数据时不要慌先从它到底脏在哪开始然后一层层去清理、去建模、去表达。那些看起来繁琐的清洗和排查工作恰恰是把你和只会跑模型的人区分开来的地方。
返回列表