ARTICLE DETAIL

资讯详情

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

大数据数据清洗实战:Spark、Great Expectations与规则引擎选型指南

大数据数据清洗实战:Spark、Great Expectations与规则引擎选型指南 1. 大数据场景下的数据清洗难点和普通ETL不一样在哪接手过几个数据平台项目之后我越来越确信一件事很多人对数据清洗的理解还停留在写脚本把空值填上、把重复行删掉这个层面。这个认知在单机小数据集上没大问题但一旦进入大数据领域数据清洗的难度会呈指数级上升它不再是一个纯粹的数据修正问题而是一个集存储、计算、调度、质量度量于一体的系统工程。先说说我最近一次踩坑经历。当时负责一个用户行为日志的清洗任务数据量大概是每天几亿条存储在HDFS上。最初版本我直接用Python脚本逐文件读取、清洗、再写回。本地测试小样本一切正常但一到集群全量跑就暴露出几个问题单机脚本处理几亿条数据动辄跑十几个小时中间某个节点内存溢出导致整个任务失败清洗逻辑里有个字段解析错误跑到一半才发现只能杀掉任务从头再来。那时候我才真正意识到大数据场景下的清洗核心难点根本不是怎么清洗而是用什么架构去承载清洗。具体来说大数据清洗跟传统数据清洗有几处本质差异数据量级不同。传统工具处理的是几十万到几百万行大数据清洗面对的是几十亿到上百亿行存储和计算必须分布式。数据来源复杂。日志文件、数据库binlog、消息队列、第三方API都有可能格式从结构化表到嵌套JSON、Avro、Parquet都有清洗需要先完成格式适配和schema对齐。质量问题的形态更多样。除了常见的缺失值、重复值、异常值还会遇到因为上游升级导致的字段含义变化、时间戳时区混乱、半结构化数据嵌套层级不一致这类数据血缘层面的脏。时效性要求不同。离线清洗可以容忍小时级延迟但实时链路的数据清洗要在秒级甚至毫秒级完成这直接决定了工具选型。明白了这几点再去选工具才会有方向感。市面上号称能做数据清洗的工具很多但真正扛得住大数据场景的其实就那么几类下面我按使用场景逐个说清楚。2. 通用批处理清洗Spark与Pandas谁更顺手如果只能选一个工具承担大数据环境下的主力清洗任务我会毫不犹豫推荐Apache Spark。原因很简单它天生就是为分布式数据处理的场景设计的而且对Python和SQL都支持得很好能覆盖从数据读取、清洗转换到落库的全链路。2.1 Spark洗数据的核心玩法用Spark做数据清洗最常用的无非三种姿势第一种是用Spark SQL直接写清洗逻辑。对于熟悉SQL的团队来说这是最顺手的。比如日志表里有个user_id字段因为上游埋点问题经常出现空字符串清洗时直接INSERT OVERWRITE TABLE cleaned_user_log SELECT user_id, COALESCE(NULLIF(user_id, ), unknown) AS user_id_clean, event_time, ... FROM raw_user_log WHERE dt 2024-06-01这种写法最大的好处是逻辑透明、好维护团队成员看一眼就能懂。而且Spark SQL的优化器会自动做谓词下推、列裁剪这些优化不需要手动调优就能跑得不错。第二种是DataFrame API。适合清洗逻辑比较复杂、用SQL不太好表达的场景。举个例子处理嵌套JSON时需要展开多层结构来提取字段from pyspark.sql import functions as F df spark.read.json(hdfs:///data/raw/events/) df_clean ( df .withColumn(event_id, F.col(event.id)) .withColumn(user_id, F.col(context.user.id)) .withColumn(event_time, F.to_timestamp(F.col(event_time), yyyy-MM-dd HH:mm:ss)) .withColumn(platform, F.when(F.col(platform).isin(iOS, Android, Web), F.col(platform)).otherwise(unknown)) .dropDuplicates([event_id, user_id, event_time]) )第三种是Spark 3.x引入的Delta Lake能力。我特别喜欢用Delta Lake做清洗的中间层因为它支持ACID事务、时间旅行和upsert。这意味着清洗任务如果跑挂了不需要全量重跑可以基于之前的快照继续处理而且数据回滚非常方便。在实际项目中我把这个能力用在先清洗、后校验、校验不过就回滚的流程里省掉了大量重跑成本。2.2 Pandas在什么场景下依然不可替代看到这里有人会问那Pandas是不是就没用了也不是。Pandas在数据量在几千万行以内、单机内存能扛住的场景下依然有它不可替代的优势生态成熟、API丰富、做探索性分析极其顺手。我常用的一个组合是先用Spark做粗清洗和抽样把数据量降下来再导出小样本到Pandas里做深度分析和规则探索。清洗规则在Pandas里验证通过后再翻译成Spark SQL或DataFrame代码跑全量。这种方式既能利用Pandas的灵活性又能保证全量任务的分布式能力。2.3 为什么不用Hive或MapReduce其实很多老项目里是用Hive做清洗的Hive本身也能胜任但它的痛点在于第一Hive的延迟较高不适合快速迭代第二Hive在复杂嵌套数据类型的处理上远不如Spark顺手第三Hive对机器学习类清洗逻辑比如用聚类算法识别异常值基本无能为力。MapReduce就更不用提了开发效率太低现在几乎没有团队新写MapReduce任务。所以综合来看Spark是当前大数据离线清洗场景下的最优解之一。2.4 Dask作为轻量替代如果你的数据量没有大到必须上Spark但Pandas又撑不住Dask是个很好的中间选择。Dask的接口跟Pandas几乎一致但能利用多核和分布式调度来扩容。我记得有个项目单机Pandas处理3000万行数据要跑40分钟迁移到Dask的distributed模式配了四台机器时间压缩到10分钟以内。关键是改动量极小基本就是换个import的事。当然Dask的生态和稳定性跟Spark比还是有差距所以我一般把它定位在Pandas的分布式扩展而不是Spark的替代品。3. 数据质量规则引擎从被动救火到大开箱的Great Expectations清洗做得多了就会面临一个更头疼的问题数据质量问题不是在清洗时发现的而是在下游报表出来之后才暴露的。这种事后救火的模式非常被动因为脏数据已经流通出去了。所以我在项目里会主动引入数据质量规则引擎把质量校验前置到清洗管道里。3.1 Great Expectations实战体验提到数据质量工具首推Great Expectations以下简称GX。它不是一个清洗工具而是一个数据质量断言框架核心思想是把你对数据的期望写成规则然后自动去校验。它的关键概念有三个Expectation、Checkpoint、Data Doc。Expectation是规则本身。比如user_id列不能为空金额字段必须大于0这些都可以表达成Expectationimport great_expectations as gx context gx.get_context() batch context.get_batch( batch_request{ datasource_name: my_spark_datasource, data_connector_name: default_inferred_data_connector_name, data_asset_name: clean_user_log, } ) expectation_suite context.suite.get(nameuser_log_quality) expectation_suite.add_expectation( gx.expectations.ExpectColumnValuesToNotBeNull(columnuser_id) ) expectation_suite.add_expectation( gx.expectations.ExpectColumnValuesToBeBetween(columnamount, min_value0, max_value100000) ) validation_result batch.validate(expectation_suite)Checkpoint负责把规则跟具体的数据批次绑定起来以调度方式定期执行。Data Doc则自动生成可视化的数据质量报告不需要自己写报表。这个工具的实用价值在于它能让你建立起数据质量基线每次清洗管道跑完自动生成一份质量报告规则挂了就报警。我还遇到过一种情况上游改造导致某个字段的含义变了旧规则瞬间报警如果靠人去肉眼审视根本发现不了。3.2 其他规则引擎的横向对比除了GX业界还有几个不错的开源方案我整理了一张表方便对比工具核心特点适合场景注意点Great Expectations规则丰富、Data Doc可视化、支持Spark/Pandas/SQL需要对数据质量建立基线的中大型团队规则一多之后维护成本不低DeequAmazon开源基于Spark主打数据质量约束的自动验证数据湖上的大规模自动化质量校验只在Spark上运行跟AWS生态绑定较深Soda Core轻量级、配置即代码、跟CI/CD集成友好敏捷团队做数据管道质量门禁生态相对年轻社区没有GX大dbt testdbt内置的测试功能基于SQL写断言已经用dbt做数仓建模的团队只覆盖dbt能访问到的数据不适合做全链路质量校验我给团队的建议是如果技术栈是Spark优先考虑Deequ或GXGX对Spark有原生支持如果清洗管道主要是SQL形态dbt test是最低成本的接入方式。3.3 规则引擎背后隐藏的清洗规则沉淀价值规则引擎还有个容易被忽略的价值它倒逼你把什么是脏数据这件事显性化。很多团队的数据清洗规则是散落在脚本里的别人没法知道这个数据集到底校验了什么、清洗了什么。而通过规则引擎里的Expectation列表新人能快速了解这个数据集的质量契约。我在实际项目中把每一个Expectation都跟具体的数据集文档做了关联下游要消费数据时先看质量报告再决定是否使用这比任何口口相传都有用。4. 数据比对与抽样老牌工具照样值得留一手前面说的都是工程化的清洗框架但真实业务里还有一个高频场景是它们覆盖不到的数据清洗到底改了什么、准确率如何这需要对比原始数据和清洗后数据的差异还需要抽样验证清洗规则的合理性。4.1 用SQL做数据质量巡检在大数据平台上我见过不少人把数据比对做成一个笨重的人工流程但真正高效的做法是把比对逻辑写成定时SQL任务每天自动巡检。举个例子检查清洗前后的记录数是否一致、主键是否重复、关键字段分布是否发生漂移-- 对比raw和clean两个表的记录数 SELECT record_count AS metric, (SELECT COUNT(*) FROM raw_log WHERE dt2024-06-01) AS raw_count, (SELECT COUNT(*) FROM cleaned_log WHERE dt2024-06-01) AS clean_count; -- 检查清洗后是否仍有重复主键 SELECT event_id, COUNT(*) AS cnt FROM cleaned_log WHERE dt2024-06-01 GROUP BY event_id HAVING COUNT(*) 1;这种方式的优点是直观、可以挂在调度平台上每天自动跑有异常就触发告警。相比在一堆日志里翻问题自动巡检能节省大量时间。4.2 抽样是清洗规则验证的加速器关于抽样很多人的第一反应是随机抽样。但在数据清洗场景下纯随机抽样往往效率很低因为脏数据在全集里的比例可能只有万分之一随机抽一万条可能一条脏数据都碰不到。我实际用的是分层抽样和异常聚焦抽样的组合。分层抽样的逻辑是先按业务维度比如渠道、日期、用户等级分层再从每个层里随机抽样本保证样本对全集的代表性。异常聚焦抽样则是先用简单的统计规则比如字段为空、字段值超出正常范围、时间戳不在合法区间把可疑数据筛出来再对这些可疑数据做高比例抽查用来验证清洗规则是否cover住了真实的脏数据形态。抽样的核心价值是在小样本上快速迭代清洗规则而不需要每次都在全量数据上反复试错。我通常会把抽出来的样本存成独立的数据集专门供规则探索和模型验证使用。4.3 OpenRefine与DataCleaner这两个老牌桌面工具虽然大数据工具链很发达但我依然会在某些特定时刻用回OpenRefine和DataCleaner这类桌面级工具。OpenRefine最擅长做的是探索式清洗它提供了很好的交互界面能看到每一列的数据分布、聚类结果、重复组还能通过GREL表达式写复杂的转换逻辑。我在早期做数据探查时经常用它对小样本做快速处理比如合并不同写法但含义相同的分类值北京和北京市这种OpenRefine的聚类功能对这种场景很有用它会自动把相近的文本归为一组省去大量手写规则的时间。DataCleaner更偏数据质量管理方向内置了完整性、唯一性、有效性等分析指标还有数据血缘追踪能力。对于想快速给一个数据集体检的场景DataCleaner的开箱体验很好。这两个工具不适合处理大数据量但作为规则探索台和数据质量体检器放在整个流程里非常有用。我的习惯是先用这些工具在小样本上探索出数据质量问题形成初步的清洗规则然后再把规则迁移到Spark或规则引擎上去跑全量。5. 为什么要留一个“脏数据仓库”在做数据清洗这个领域有个原则可能反直觉但极其重要清洗前必须保留原始数据清洗后的结果也要能往回追溯。我通常建议团队在数仓里专门留一层脏数据仓库。5.1 分层架构里的ODS层别急着做数据清洗标准的数据仓库分层里ODS操作数据存储层存放原始数据。很多团队为了省事在数据接入ODS时就顺手做了一堆清洗和过滤这其实是个隐患。因为一旦后续发现清洗规则有问题或者需要追溯历史数据的某个细节原始数据已经被污染了没法恢复。正确的做法是ODS层尽量保持原样哪怕字段名不规范、有缺失值、有异常数据都先原封不动落下来。清洗工作放到DWD或DWS层去完成。这样做还有一个好处ODS层的数据可以支持问题回溯——当业务方质疑某个数据不对时你可以直接对比ODS原始数据和清洗后数据快速定位问题出在规则设计还是上游数据源头。5.2 清洗规则版本化比代码版本化更复杂谈到追溯就不得不提清洗规则的版本管理。我自己在这个问题上吃过亏有一次修改了清洗规则结果下游报表数据发生了跳动但因为没有保存旧版本的规则根本无法解释数据变化的原因。后来我才意识到清洗规则其实也是数据资产的一部分规则版本化必须跟代码版本化同步。我现在的方法是每一条清洗规则都记录它的生效时间、版本号、修改人、修改原因。这些元数据存到数据资产目录里这样任何一次数据变化都能追溯到具体的规则变更。这看起来像是加重了流程负担但在数据量越大、下游受众越多的场景下回报会非常明显。5.3 清洗失败时怎么办幂等性设计不可忽视数据清洗任务在分布式环境下并不是每次都能成功的网络抖动、内存溢出、上游数据延迟都可能导致任务失败。因此清洗任务必须具备幂等性也就是说不管跑多少遍同样的输入拿到同样的输出不会产生重复数据或脏数据。我在设计清洗管道时会遵循几条原则清洗任务一律使用主键分区做去重约束保证重复执行不会产生重复记录。输出目标表使用覆盖写的模式INSERT OVERWRITE而不是追加写。每次清洗任务带上批次号batch_id方便按批次回滚和处理。任务失败时自动进入重试队列如果重试仍然失败必须发出告警并暂停下游调度而不是让脏数据继续流动。这些设计可能看起来属于工程基本功但我在实际项目里确实见过不少团队因为忽略幂等性导致数据翻倍、报表错乱、最后整个平台数据可信度崩塌的案例。数据清洗这个环节做得好不好往往不在清洗本身而在这些细节是否扎实。6. 我日常选型的心得和避坑经验写到这里工具层面的东西基本说完了但我还想分享一些在实际项目里反复踩坑后总结出来的心得希望帮你少走弯路。首先工具选型一定要先想清楚自己的数据规模和技术栈。我曾经见过一个团队数据量就几百万行非要去搭一套Spark集群做数据清洗结果集群运维成本比清洗本身还高。反过来说数据量已经到亿级了还坚持用Pandas在单机上硬扛最终只会耽误业务进度。我的建议是做一个简单的判断矩阵数据规模技术栈/偏好推荐方案百万行以内Python为主Pandas / OpenRefine千万行级别Python为主Dask / Pandas抽样并行千万到亿级有Spark环境Spark SQL / PySpark十亿级以上构建数据湖/数仓Spark Delta Lake Great Expectations/Deequ已经用dbt建模SQL优先dbt test 自建的SQL巡检其次不要试图用一个工具解决所有问题。很多人会问到底哪个工具最好但真实世界里从来没有万能的清洗工具。Spark很强大但它的强项是批处理GX擅长质量校验但不适合做复杂的转换逻辑OpenRefine交互体验好但处理不了大数据。把不同的工具组合起来让每个工具做自己最擅长的事才是工程上最务实的做法。最后持续提醒自己数据清洗不是一个一次性的任务而是一个持续演化的过程。业务在变上游数据源在变数据质量问题的形态也在变。今天认为正确的清洗规则三个月后可能就成了错误规则。所以清洗管道必须留出扩展性规则引擎的规则要能灵活增删清洗逻辑要模块化这样才能避免每次需求变化都推倒重来。再分享一个我个人的小习惯每写一条清洗规则都会在旁边注明为什么需要这条规则、它解决的是什么类型的数据质量问题。这个习惯在团队协作中帮我解决了很多问题——新同事接手时不需要反复来问这个字段为什么这样处理业务方找过来时也能快速给出解释。把规则背后的业务逻辑沉淀下来这比工具本身更重要因为工具可以随时换但对业务数据的理解才是清洗工作的核心价值。
返回列表