ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

SparkSQL实战指南:从建表读写到性能调优的全面梳理

SparkSQL实战指南:从建表读写到性能调优的全面梳理 从事大数据开发这行有个东西你早晚得正面刚——SparkSQL。不管你是做离线数仓、实时数仓还是数据湖湖仓一体跑批、清洗、统计、关联十有八九都要落到SparkSQL上。今天不聊那些理论满天飞的架构解读就实实在在把SparkSQL的常用操作顺一遍从建表读写到过滤聚合从多表关联到性能调优穿插一些我自己实战中踩过的坑和验证过的经验。无论你是刚入门想找个学习路线还是写了几个月Hive想试试Spark的威力这篇都能让你直接上手去用而不是看完还懵懵的。1. SparkSQL的角色定位为什么现在大数据处理绕不开它1.1 从RDD到DataFrame一张表带来的思维转变刚接触Spark的人通常都是从RDD开始的。map、flatMap、filter、reduceByKey这些算子是RDD的核心写起来确实灵活但有个很现实的问题——你写的代码Debug起来费劲优化起来更费劲因为RDD只是一个分布式数据集它不关心数据里面装的是什么。后来SparkSQL的出现从根本上改变了这个局面。它把数据抽象成了DataFrame和Dataset对标关系型数据库里的表。DataFrame有Schema有列名有类型这意味着Spark能够知道每一列是什么数据类型可以做基于列的优化。这个变化用一个生活类比来说就是——RDD像是你从菜市场买回来的一袋乱糟糟的菜你自己清洗、切配、下锅而DataFrame则是按荤素分好的净菜还贴了标签厨师拿到就能直接炒甚至还能提前规划需要用几个锅。这就引出了SparkSQL最核心的价值让分析师和工程师用同一种语言对话。业务分析师懂SQL不太会写Java/Scala/Python工程师写代码很利索但是跟业务对口径的时候效率不高。SparkSQL提供了一条路——分析师写SQL工程师用DataFrame API大家操作的是同一份数据、同一个执行引擎最终逻辑还能互相转换。1.2 Catalyst优化器你写的SQL是怎么被偷偷优化的很多人觉得SparkSQL就是把SQL翻译成RDD操作然后跑就完事了。实际不是中间有一套非常关键的组件——Catalyst优化器。Catalyst是一个基于树的查询优化框架它做的事情简单说就是四步分析Analyzer、逻辑优化Logical Optimization、物理规划Physical Planning、代码生成Code Generation。第一阶段把SQL字符串解析成语法树然后绑定元数据搞清楚每一列的类型、表属于哪个数据库第二阶段做逻辑优化比如谓词下推、列剪枝、常量折叠把没必要的计算直接砍掉第三阶段把逻辑计划转化成物理计划选择一个最优执行策略最后通过WSCGWholeStageCodeGen把整个Stage变成一段Java代码来执行。我举个谓词下推的例子。你写了一个查询SELECT user_id, order_time FROM orders WHERE amount 100在没有优化的情况下引擎会把整张表的数据都读进来然后过滤最后投影出两列。有了谓词下推之后Spark会先把过滤条件下推到数据源——如果是Parquet文件直接在读取时就跳过大量不满足条件的行然后做列剪枝只读取三列中需要的两列。这一下一张亿级大表可能扫描量直接少了一个数量级。所以你会遇到一个现象一样的SQL在Hive里跑30分钟在SparkSQL里跑5分钟。除了内存计算之外Catalyst的优化起了很大作用。这也是为什么SparkSQL的适用场景远不止离线批处理——只要你的数据是结构化或半结构化的用SparkSQL处理效率、开发速度、资源利用率都会比纯RDD开发好很多。2. 环境准备与前置认知先把手头的工具调好2.1 版本选型与依赖引入别小看这一步写了这么多年Spark我最大的忠告就是先定好版本再写代码。版本不一致导致的坑比业务逻辑本身的坑还要多。目前主流的生产环境版本是Spark 3.x其中3.3、3.5使用率较高。Spark 3.2开始对JDK版本要求是JDK8/11/17均可3.4的时候JDK17可以正常跑但如果你还挂在JDK8上也不必着急升其实大多数场景都没问题。关键是要注意Scala版本——Spark 3.x对应Scala 2.12/2.13。如果你写的是Java代码或者Scala代码引入Maven依赖时一定要把scala版本带上比如dependency groupIdorg.apache.spark/groupId artifactIdspark-sql_2.12/artifactId version3.3.2/version /dependencygrouldId里的_2.12指的就是Scala编译版本如果跟你本地安装的Scala不一致运行时会报出让人抓狂的兼容性问题。如果你是在Python环境用PySpark也需要检查PySpark版本跟Scala版本配套——PySpark 3.3对应Spark 3.3这个不会出错但也要注意Python 3.8。还有一种常见场景是公司内部有Hadoop集群你提交Spark任务时要跟HDFS版本兼容。目前Spark 3.3默认的Hadoop是3.3.4如果你的集群是Hadoop 3.x问题不大如果还在2.x就需要引入对应版本的spark-hadoop-cloud或者手动指定hadoop-client版本否则会出现RPC协议不兼容一类的问题。2.2 两种编程入口SparkSession与Hive集成从Spark 2.0开始SparkSession就是统一入口把SQLContext和HiveContext合并了。现在写代码第一步永远是from pyspark.sql import SparkSession spark SparkSession.builder \ .appName(spark_sql_demo) \ .config(spark.sql.shuffle.partitions, 200) \ .enableHiveSupport() \ .getOrCreate()注意这个enableHiveSupport()。如果你的集群配置了Hive Metastore开启这个支持之后SparkSQL就可以直接操作Hive里的库表读取元数据并复用现有的Hive表这是很多企业从Hive迁移到Spark的关键一步——表结构不用重新建SQL不用大改直接切引擎。但是如果你只是本地测试没有配Hive就不要随意开启否则Spark会尝试连接Hive Metastore然后报错。本地开发时可以直接这样spark SparkSession.builder \ .appName(local_demo) \ .master(local[*]) \ .getOrCreate()另外还要注意一个坑SparkSQL读Hive表的时候默认会把文件内容读成InternalRow如果表里有一些自定义SerDe或者特殊字段类型可能会解析失败。多数情况下Parquet/ORC格式的表都没问题TextFile格式的表最好先做一次预处理。3. 高频实操一张真实数据表带你过完常用操作3.1 造一张网约车订单表先学会建表与加载实操是最有说服力的。我用一个网约车订单数据的场景来串联整个章节——这也是我在实际项目中反复操练过的模板覆盖了SparkSQL最常用到的所有操作。假设我们需要分析一个网约车平台的司机接单行为数据从Kafka实时写入HDFS落成Parquet分区表按天分区。那么建表语句可以是CREATE TABLE IF NOT EXISTS dwd_trip_order ( order_id STRING, driver_id STRING, passenger_id STRING, start_lng DOUBLE, start_lat DOUBLE, end_lng DOUBLE, end_lat DOUBLE, order_amount DECIMAL(10,2), order_status INT, city_id STRING, duration_min INT, distance_km DOUBLE, eta_min INT, load_ts TIMESTAMP ) PARTITIONED BY (dt STRING) STORED AS PARQUET这里建了分区字段dt按天管理数据后续查询如果带上分区过滤扫描的数据量会大大减少。这是一个被说烂但依然极其重要的优化点——很多SQL跑得慢根本不是计算慢而是扫描的数据太多。建完表之后从HDFS加载数据可以这样df spark.read.parquet(hdfs://namenode:8020/data/dwd/trip_order/dt2024-06-01) df.write.mode(overwrite).insertInto(dwd_trip_order)这里有个细节要注意insertInto是按照表的列名顺序插入的而saveAsTable会使用DataFrame的Schema来建表或重写表结构。两者行为不一样用错容易导致数据错位。如果你不确定当前DataFrame列和表字段顺序一致最稳妥的办法是先注册成临时视图再走SQLdf.createOrReplaceTempView(tmp_order) spark.sql( INSERT OVERWRITE TABLE dwd_trip_order PARTITION (dt2024-06-01) SELECT order_id, driver_id, passenger_id, start_lng, start_lat, end_lng, end_lat, order_amount, order_status, city_id, duration_min, distance_km, eta_min, load_ts FROM tmp_order )这么写的好处是字段顺序一目了然出问题也好排查。3.2 日常读写与转换操作filter、select、withColumn表建好之后日常用得最多的就是过滤、选列、派生列。先说过滤SparkSQL有两种写法——SQL风格的WHERE和DataFrame API风格的filter。# DataFrame API picked_orders df.filter(df.order_status 1) # SQL风格 picked_orders df.where(order_status 1)这两种写法到底有没有区别底层结果一样where本身就是filter的别名。但有一点值得注意如果你在filter里用了字符串表达式比如order_status 1Spark需要解析这个字符串为表达式有一定解析开销虽然很小如果用Column表达式df.order_status 1编译器能更早发现问题。性能上差距微乎其微我建议代码可读性优先——团队里如果SQL熟人多就直接用字符串写法大家看得懂。接下来是select和withColumn。很多人刚接触时会把这两个搞混。select是从现有DataFrame选出若干列生成一个全新的DataFrame不会增加新的列withColumn则是在现有DataFrame基础上新增加一列或者替换一列。# select 使用示例 driver_city df.select(driver_id, city_id, order_amount) # withColumn 新增一列计算每单的客单价元/公里 df2 df.withColumn(amount_per_km, func.round(df.order_amount / df.distance_km, 2))这里有个经验之谈withColumn虽然名叫with但它返回的是新DataFrame原来的df没有被修改。如果你在循环里反复调df.withColumn(...)实际上是不断在原有DataFrame上构建新对象如果循环几百次会生成一个超长的血缘关系链Lineage导致执行计划优化困难甚至StackOverflow。这在做特征工程时经常遇到——几十个withColumn串在一起任务直接崩。解法是尽量合并表达式或者把一个一个的字段计算放在同一个select里完成。df3 df.select( *, func.round(df.order_amount / df.distance_km, 2).alias(amount_per_km), func.unix_timestamp(load_ts).alias(load_ts_epoch) )3.3 分组聚合与窗口函数统计口径都在这里没有聚合就没有数据仓库。SparkSQL里的分组聚合最常用的就是groupByagg组合。比如统计每个城市每天的完单数、总流水、平均实付金额city_daily_stats df.filter(df.order_status 1) \ .groupBy(city_id, dt) \ .agg( func.count(order_id).alias(finish_orders), func.sum(order_amount).alias(total_amount), func.round(func.avg(order_amount), 2).alias(avg_amount) )对应的SQL写法是SELECT city_id, dt, COUNT(order_id) AS finish_orders, SUM(order_amount) AS total_amount, ROUND(AVG(order_amount), 2) AS avg_amount FROM dwd_trip_order WHERE order_status 1 GROUP BY city_id, dt这里有几个聚合时的细节count(列名)不会统计NULL值count(1)或count(*)会统计行数。如果你要统计订单数最好用count(order_id)并确保order_id无NULL如果你想统计记录条数直接用count(*)。SUM遇上NULL结果是NULL不是0。如果某天没有订单所有聚合函数的结果除了count(*)外可能都是NULL需要COALESCE处理。窗口函数也是使用频率极高的功能。比如我想给每个城市、每一天的司机按流水排名只保留该城市当天Top 10的司机SELECT * FROM ( SELECT driver_id, city_id, dt, SUM(order_amount) AS day_amount, ROW_NUMBER() OVER (PARTITION BY city_id, dt ORDER BY SUM(order_amount) DESC) AS rn FROM dwd_trip_order WHERE order_status 1 GROUP BY driver_id, city_id, dt ) t WHERE rn 10窗口函数在SparkSQL里支持得很完整ROW_NUMBER、RANK、DENSE_RANK、LAG、LEAD、SUM OVER这些都可用。但有个性能上的提醒——ROW_NUMBER() OVER (PARTITION BY ... ORDER BY ...)在数据量巨大时分区内排序需要做全局排序如果分区键的基数特别高比如几百万个driver_id那么每个driver_id一个分区反而增加了调度开销。通常建议分区粒度不要太细能用city_id dt就不要用单条的order_id。4. 多表关联与复杂场景从单表到SQL思维4.1 Join操作三种关联方式的取舍数据分析真正复杂的地方在于多表关联。网约车项目后端一般有订单表、司机表、城市维度表、评价表等分析时通常需要把它们join到一起。SparkSQL里最常用的是JOINorder_df spark.sql(SELECT * FROM dwd_trip_order WHERE dt2024-06-01) driver_df spark.sql(SELECT * FROM dim_driver WHERE dt2024-06-01) result order_df.join(driver_df, order_df.driver_id driver_df.driver_id, left)第四种——left其实是left_outer的简写表示左外连接保留左表全部记录。其他的还有inner、left_outer、right_outer、full_outer、left_semi、left_anti。这里我说一个容易被忽视但实际影响巨大的操作——left_semi和left_anti。它们不常用但应对取订单表中哪些司机在今天有订单这类判断非常高效-- 有订单的司机 SELECT * FROM dim_driver WHERE driver_id IN (SELECT DISTINCT driver_id FROM dwd_trip_order WHERE order_status 1) -- 用 left_semi 实现同样的逻辑 SELECT * FROM dim_driver d LEFT SEMI JOIN dwd_trip_order o ON d.driver_id o.driver_id AND o.order_status 1left_semi与普通join最大的不同在于它不会把右表的任何列带出来只返回左表满足关联条件的记录并且在语义上等价于EXISTS。Spark对left_semi的优化通常比IN (SELECT ...)走得更顺因为优化器知道只需要判断是否存在不需要维护关联列之外的状态。4.2 子查询与临时视图写完即用的SQL片段SparkSQL的SQL语法是完整的HiveQL超集所以子查询、CTE公用表表达式都用得上。在代码里操作时用得更多的模式是先用createOrReplaceTempView把多个DataFrame注册成临时视图然后用一段完整的SQL把它们串起来。order_df.createOrReplaceTempView(tmp_order) driver_df.createOrReplaceTempView(tmp_driver) city_df.createOrReplaceTempView(tmp_city) spark.sql( WITH city_stats AS ( SELECT o.city_id, COUNT(DISTINCT o.driver_id) AS active_drivers, COUNT(o.order_id) AS total_orders, SUM(o.order_amount) AS gmv FROM tmp_order o WHERE o.order_status 1 GROUP BY o.city_id ) SELECT c.city_name, cs.active_drivers, cs.total_orders, cs.gmv FROM city_stats cs LEFT JOIN tmp_city c ON cs.city_id c.city_id ORDER BY cs.gmv DESC ).show()把多个DataFrame注册成临时视图再统一写SQL这个习惯我会推荐给团队里的每个人。对比在代码里把join条件写成层层的lambda表达式以SQL为胶水语言的写法维护成本低太多——面试造火箭、工作拧螺丝拧螺丝也要拧得干净利落让后来接手的人能快速看懂。有一点需要注意临时视图的生命周期跟随SparkSession。如果你在Jupyter Notebook里跑了createOrReplaceTempView只要同一个SparkSession没停后面的cell都可以继续查。但在生产任务里每个SparkApplication独立任何临时视图都只在本Application内可见不要指望跨任务共享。4.3 写出SparkSQL可以听得懂的条件表达式多表关联处理中条件写得好不好直接影响能否触发优化。最常见的一个低级错误是把过滤条件写在ON子句里或者反过来把关联条件写在WHERE子句里导致执行计划低效。标准做法是关联条件只负责关联过滤条件放WHERE。比如-- 推荐写法 SELECT o.order_id, d.driver_name FROM dwd_trip_order o JOIN dim_driver d ON o.driver_id d.driver_id WHERE o.dt 2024-06-01 AND d.dt 2024-06-01如果你不小心在ON后面加了AND o.dt 2024-06-01虽然结果不一定错inner join情况下等价但Spark执行优化时对谓词下推的判断会变复杂尤其在left join场景下右表的过滤条件写在ON里和WHERE里结果完全不一样这是新手最容易翻车的地方。再补一个跟条件判断有关的细节宁用等值条件少用非等值。SparkSQL的join优化如Broadcast Join、SortMerge Join都是基于等值join来设计的。如果你写的是o.order_amount d.threshold这种非等值join优化器基本无法做各种加速只能退化为NestedLoopJoin大表之间直接完蛋。遇到这种情况要么改业务口径要么提前把阈值条件用CASE WHEN转换成等值条件再造一个辅助字段。5. 性能调优一样的SQL为什么别人跑得比你快5.1 分区裁剪与谓词下推先让扫描量降下来性能调优的第一个原则是减少扫描量不要让Spark读它不需要的数据。在建表时使用PARTITIONED BY并查询时带上分区条件这是一个分区裁剪的动作-- 好只读一天数据 SELECT * FROM dwd_trip_order WHERE dt 2024-06-01 -- 差全表扫描再过滤 SELECT * FROM dwd_trip_order WHERE load_ts 2024-06-01 00:00:00 AND load_ts 2024-06-02 00:00:00数据量稍微上亿这两种写法查询时间能差几十倍。在实际运维中我见过很多任务因为业务方只传了一个时间范围没有落到分区字段上导致每天凌晨的任务扫了整张历史表集群被打到OOM。如果确实需要按load_ts过滤但不是分区字段建议建一个映射表或者把时间用于推导分区后再查询。另一个跟谓词下推相关的点是文件格式的选择。Parquet和ORC都是列式存储天然支持谓词下推。如果你用TextFile格式存储比如CSV即使写了WHERE条件Spark也得把整行数据读出来再过滤。所以数仓建模时底层明细表存储格式强烈建议Parquet Snappy压缩既有列式存储的扫描优势又有不错的压缩比。5.2 缓存与复用不要重复计算同一份数据在同一个Spark任务中如果一份数据会被多次使用可以把它缓存起来。SparkSQL的缓存有两个级别DataFrame的cache()和临时视图的CACHE TABLE。# 方式一DataFrame 缓存 filtered_df df.filter(df.dt 2024-06-01).cache() filtered_df.count() # 触发实际计算并缓存 # 方式二SQL缓存 spark.sql(CACHE TABLE cached_order AS SELECT * FROM dwd_trip_order WHERE dt2024-06-01)缓存的坑也不少。默认的缓存存储级别是MEMORY_AND_DISK意思是先放内存放不下就溢写到磁盘。如果你在缓存前没有对DataFrame做列剪枝把几十列全缓存了占用的内存会异常大。所以缓存的黄金法则是先select出需要的列、过滤掉不需要的行再cache。还有一个我特别想提醒的操作细节cache()是一个lazy操作真正触发缓存的是后面的action比如count()、show()。如果你写完df.cache()之后又跑了一个df.write看起来好像缓存生效了但实际第一次读取时并没有复用缓存——因为写入也读了一次这个过程才把缓存填充上。要确认缓存是否生效可以看一下Spark UI的Storage页签或者执行spark.sparkContext.getPersistentRDDs排错时挺管用的。5.3 Shuffle优化数据倾斜的识别与处理数据倾斜是SparkSQL性能问题里最顽固的一个而且很多问题在跑小数据量时完全看不出来一上生产就爆。它的本质是某些key的数据量远超其他key导致某个Task处理的数据量特别大其他Task早就跑完了还在等它。典型的场景就是GROUP BY city_id某个城市的订单量特别大比如一线城市占了40%的订单。体现在Spark UI上就是某个Stage的一个Task运行时间远大于其他Task甚至会OOM。解决数据倾斜的方式有很多种我按实操中使用的优先级列一下如果倾斜的key是NULL可以把NULL值赋一个随机前缀让它们分散到多个Task处理最后再过滤掉。如果是少数热点key就加随机前缀打散再做两次聚合。第一次聚合时给热点key加随机数让数据分散第二次去掉前缀再做最终聚合。提高shuffle分区数spark.sql.shuffle.partitions把单个Task的数据量降下来。但要小心——分区数不是越大越好过多分区增加调度开销小文件也会变多。一般建议分区的目标大小在100MB到200MB之间。在join场景如果小表可以完全放进内存可以用Broadcast提示SELECT /* BROADCAST(d) */ o.order_id, d.driver_name FROM dwd_trip_order o JOIN dim_driver d ON o.driver_id d.driver_id这样Spark会把小表广播到每个Executor避免shuffle。尤其在城市维度表、司机维度表这种不过几万行的表效果立竿见影。但要注意广播表默认阈值是10MBspark.sql.autoBroadcastJoinThreshold如果表特别大即使加了hint也可能不会广播需要手动调大阈值或重新考虑方案。6. 常见问题与排错实录6.1 明明SQL没错却报语法错误SparkSQL虽然支持SQL语法但它不是完全等同于MySQL或Hive。最典型的差异包括字符串用单引号双引号在某些版本里会被当成列标识符解析导致报错。关键字冲突比如date、timestamp、window这类词如果字段名叫date查询时要加反引号否则报错。隐式类型转换严格Spark的cast规则比MySQL严格比如字符串2024-06-01和DATE类型比较有时候需要显式to_date。遇到语法错误不要慌先从最简单的查询开始缩小范围。先SELECT * FROM table LIMIT 10验证表能读再逐步加字段、加条件定位到具体出错的那一段。6.2 数据倾斜典型场景定位如果任务卡住或者极慢第一件事就是看Spark UI里哪个Stage耗时最长进入Stage详情看Task运行时间分布。如果看到某个Task运行时间比中位数高出几个数量级或者某个Task处理的数据明显大于其他Task基本就是数据倾斜。常见解决手段前面提过加随机前缀和Broadcast还有一个实用技巧是给倾斜的key加上盐值后在内层关联里用范围条件过滤减少无效关联。比如join driver表和order表时driver_idD001的司机订单特别多那么可以先把D001摘出来单独算再和其他key的结果union回去。虽然代码丑一点但好在可控。6.3 小文件问题与写出优化SparkSQL写数据时默认每个Task写一个文件。如果你的shuffle分区数是200即使最后结果只有几KB也会生成200个小文件。大量小文件带来的问题有两个一是HDFS NameNode压力大二是下次读取时扫描效率低下因为要频繁打开文件。解决小文件的思路是合并输出df.coalesce(10).write.mode(overwrite).partitionBy(dt).parquet(path)coalesce只做窄依赖合并不触发shuffle适合源数据本身已经分区清晰的情况如果需要重新分布数据得用repartition它会做一次shuffle。写出后再配合Hive的ALTER TABLE ... RECOVER PARTITIONS把新分区注册上这一套流程就完整了。还有一个经验写完数据后用一个小任务做文件数巡检列一下每个分区下的文件数量和大小超过阈值的就压缩合并一次。跑批任务最怕的就是数据越写越碎定期治理能让整个数仓的健康度保持住。6.4 表数据与元数据不一致怎么办msck repair table table_name这招我用了无数次场景是外部分区表直接往HDFS路径丢数据后SparkSQL查不到新分区。因为Metastore里没有分区信息Spark/Hive根本不知道有这些新数据。MSCK REPAIR TABLE dwd_trip_order;这个命令会扫描表目录下的所有分区目录把新的分区元数据注册进去。大数据量表跑这个命令可能较慢因为它要列目录。不过在数据入仓后执行一次能避免很多明明文件在那为什么查询查不到的诡异问题。回到最开始说的那个点SparkSQL这套技术栈已经是当前大数据处理的事实标准。不管你是从Hive迁移过来的还是直接上手Spark的理解它的常用操作背后那些执行机制和优化思路比背几个API要重要得多。我在实际项目中最大的体会是先保证逻辑对再考虑性能先能用起来再慢慢调优。SparkSQL的生态足够成熟网上案例也很多真正把常用操作吃透绝大多数业务场景都能覆盖了。最后再分享一个小技巧——如果你在写一段特别长的SQL最好先在客户端用EXPLAIN命令看一下执行计划确认过滤条件下推、join类型是否符合预期这十几秒的检查比任务跑挂了再排查要省一个小时。
返回列表