ARTICLE DETAIL

资讯详情

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

createDataFrame 结合 Spark SQL:从建表到数据清洗与慢 SQL 优化

createDataFrame 结合 Spark SQL:从建表到数据清洗与慢 SQL 优化 很多人第一次接触 Spark 的时候最先被敲门的往往是两个词DataFrame 和 SQL。原因很简单它们足够“像表”足够“像数据库”写起来既不陌生跑起来也快。而把这两者联系起来的就是createDataFrame。你可以把它理解成 Spark 里手工建表的那只手——无论是从 RDD、从代码里的集合、还是从文件读进来的数据最终想要让 SQL 直接认账基本都得过这一关。本篇就围绕createDataFrame结合 Spark SQL 的实际用法给你拆一拆背后的原理、步骤和踩过的坑顺便把相关热词里提到的集群搭建、数据清洗、窗口函数、慢 SQL 优化这些话题都串进去看完你能直接用在自己的数据分析或 ETL 场景里。1. 总体思路为什么用 DataFrame 而不是 RDD1.1 DataFrame 到底比 RDD 好在哪很多人学习 Spark 时会先接触 RDD因为教学例子大多从 RDD 开始。可实际工作中真正大量跑的业务代码几乎都是 DataFrame 或 SQL。本质上 DataFrame 是带 Schema 的 RDD但就这一层“字段定义”让整个计算逻辑发生了质变。RDD 在 Spark 眼里就是一堆 Java/Scala 对象它不知道你的数据里哪一列是数值、哪一列是字符串只能通过用户自己写的函数去处理。而 DataFrame 相当于把数据结构显式告诉了 SparkSpark 就能做类型检查、谓词下推、列裁剪以及 Catalyst 优化器那一整套逻辑改写。举个例子你用 RDD 写一个map函数去做过滤跑起来就是个 map但你用 SQL 写同样的条件Spark 可能把过滤条件直接推到底层数据源省掉读一整列数据的开销。这就是为什么业界普遍说 DataFrame 比 RDD 更适合数据分析性能差距可以轻松达到数倍甚至一个数量级。1.2 DataFrame 与 Spark SQL 的共生关系createDataFrame创建的 DataFrame并不只是给你当普通对象用的它天然可以和 Spark SQL 对接。核心操作是把它注册成一个临时视图然后spark.sql()里直接写 SQL 语句。意思就是你可以从一段 JSON、一个 CSV、甚至代码里的数组随意造出一个 DataFrame然后立刻用“SELECT、JOIN、GROUP BY”去查它。这让 Spark 在实际业务中非常灵活尤其是数据清洗阶段——从不同来源读进来的数据长相各异你完全可以用createDataFrame统一成标准的表格式结构再交给下游的 SQL 逻辑处理。从工作的角度看DataFrame 与 Spark SQL 的关系更像“编程接口”和“查询接口”的关系。DataFrame API 适合在代码里做复杂转换SQL 适合做即席查询和可读性强的分析逻辑。两者切换几乎没有成本因为底层都建立在 Catalyst 优化器之上。2.createDataFrame核心细节与实操要点2.1 从 RDD 创建 DataFrame最基础的路子先从最常用的场景说起。你手里有一个 RDD里头可能是Row对象、可能是元组、可能是 case class你想把它变成一张可以跑 SQL 的表。这时有三个选择隐式推断、显式指定 Schema、或者从 RDD of Row 加上动态构造的 Schema。隐式推断的写法Scala 里最常见val peopleRdd spark.sparkContext.textFile(people.txt) case class Person(name: String, age: Int) val peopleDF peopleRdd.map(line Person(line.split(,)(0), line.split(,)(1).toInt)).toDF()这里用到了 Scala 的 case classSpark 会自动把它反射成 Schema。Java 里没有这种天然特性你需要手动给StructType。Python 下面没有 case class但可以构造一个Row列表再传给createDataFrame。Python 中的等价钱from pyspark.sql import Row people_rdd sc.textFile(people.txt) people_rdd.map(lambda line: Row(line.split(,)[0], int(line.split(,)[1]))) people_df spark.createDataFrame(people_rdd, schema[name, age])显式指定 Schema 的好处在于可控。自动推断要扫描一次数据如果数据量很大这个扫描是有成本的。而且一线数据经常有空值、类型不统一自动推断特别容易把本该是数字的列推断成string。所以我的习惯是能显式写 Schema 就显式写哪怕多写几行后面省下的排查时间远远超过这几分几秒。2.2 从集合直接创建 DataFrame除了 RDD日常调试和写单测时最常用的还有直接从代码里的集合创建 DataFrame。这种写法非常直观特别适合给你自己快速验证一段逻辑或者写单元测试时造点测试数据。Python 示例data [(Alice, 34), (Bob, 45), (Cathy, 29)] df spark.createDataFrame(data, [name, age]) df.show()Scala 里也类似通过spark.createDataFrame(Seq((Alice, 34), ...)).toDF(name, age)。这里 Spark 会利用元组的隐式转换帮你编好 Schema非常省事。但有一个容易忽略的坑如果列表里的元素是 Python 的dictSpark 虽然能推断 Schema但字段顺序不一定按你想要的顺序输出。因为 dict 是乱序的Spark 会按 key 的哈希顺序来构造列从而导致df.select(...)的时候列名对不上预期。规避方法很简单使用Row或者先排序 keys。这种细节网上很少有人提但实测工作中特别容易在“从 JSON 转 DataFrame 再注册成临时视图”的场景里踩中。2.3 Schema 定义的正确姿势createDataFrame时候最重要的就是 Schema 设计。你可以在StructType里定义StructField每个字段同时指定类型和是否可空。想构建一个理想的 Schema特别是要用于 Spark SQL 的时候需要把nullable记得特别清楚。因为 SQL 中NULL的存在会影响聚合函数和过滤条件的结果如果你把该标成 nullable 的字段标成 false一旦数据里出现 null整个 job 直接抛异常排查的时候还要一条一条翻日志很痛苦。一个比较稳的做法是先在外部做一次数据探查搞清楚每个字段的最小最大值、空值率、类型范围之后再跑到StructType()里去定义。很多同学写到这里顺手就inferSchemaTrue让它自动推断看起来没问题但等到数据量一大或者是复杂嵌套格式自动推断经常会把嵌套字段搞成一堆struct你后续想用 SQL 里的get_json_object或者col(a.b.c)时容易处处受阻。2.4 从文件读取后转化为 DataFrame实际生产环境里很少有数据是直接在代码里构造好的。常见的是从 CSV、JSON、Parquet、JDBC 这些数据源读进来。这时不需要手动调用createDataFrame因为spark.read.csv()之类的接口返回的已经是 DataFrame 了。你可以把它理解成“Spark 内部帮我调了一次 createDataFrame”。但知道这条链路非常有用遇到奇奇怪怪的读文件问题说到底还是底层创建 DataFrame 的结构定义问题。比如读 CSV 时你可以不像默认那样由引擎猜 Schema而是手动传入from pyspark.sql.types import StructType, StructField, StringType, IntegerType schema StructType([ StructField(id, IntegerType(), True), StructField(name, StringType(), True), StructField(age, IntegerType(), True) ]) df spark.read.schema(schema).csv(path/to/file.csv)这样它读取的效率更高而且不会出现把00123这种前导零的 ID 误删掉的诡异状况。3. 实操过程从 DataFrame 到 SQL 查询的完整链路3.1 初始化 SparkSession选好版本本人现在主用的是 Spark 3.x。SparkSession是你所有操作的入口createDataFrame、sql、read全从它身上调用。Python 里通常这么写from pyspark.sql import SparkSession spark SparkSession.builder \ .appName(CreateDataFrameDemo) \ .master(local[*]) \ .config(spark.sql.shuffle.partitions, 4) \ .getOrCreate()这里config不是随便写的。spark.sql.shuffle.partitions控制 SQL JOIN 或聚合操作的分区数默认是 200对于本地小数据集、写单测、快速验证逻辑来说200 个分区纯属浪费本地跑会卡。我一般会在开发和调试时设小一点等真正提交到集群时再按集群核心数调整。有同学会在网上搜到SparkContext和SQLContext的老写法。老版本的 Spark 2.0 之前你需要单独 new 一个SQLContext现在完全不需要了。统一从 SparkSession 进这部分你在各种博客和社区帖子里会反复看到如果看的是旧笔记建议先对齐版本。3.2 模拟一份业务数据创建 DataFrame假设我们要处理一个网约车偏好的数据集。先将小批量样本写成 Python 里的数组data [ (2024-05-01, driver_001, passenger_888, 12.5, completed), (2024-05-01, driver_002, passenger_777, 8.0, cancelled), (2024-05-02, driver_001, passenger_999, 15.2, completed), (2024-05-02, driver_003, passenger_666, 30.9, completed), ] schema [order_date, driver_id, passenger_id, amount, status] trip_df spark.createDataFrame(data, schema) trip_df.printSchema() trip_df.show()printSchema()可以确认每个字段的类型。show()默认打印前 20 条对调试阶段太快了。这里需要注意如果传入的数据中有一个字段类型是整数你在 Python 里直接写过12.5Spark 会反推成DoubleType。如果你原本预期它是个 Decimal 类型来避免浮点误差就必须用StructType显式指定。不做这步后面做金额聚合可能出现 0.1 0.2 不等于 0.3 的经典浮点问题。3.3 注册临时视图与全局视图有了 DataFrame你就能让它变成一个 SQL 概念里的表。最常用的是createOrReplaceTempViewtrip_df.createOrReplaceTempView(trip_table)然后直接写spark.sql(SELECT order_date, COUNT(*) AS cnt, SUM(amount) AS total FROM trip_table GROUP BY order_date).show()这就完成了一次完整的“从 createDataFrame 到 Spark SQL”的链路造数据、定义表、跑查询。实际业务中你可以把这个临时视图当作中间层让多个 SQL 片段分别读取同一份数据避免来回搬运。你还会看到createGlobalTempView的写法。全局视图要在 SQL 里引用global_temp.trip_table并且它在一个 Spark 应用的所有会话之间可见。绝大多数数据管道里一个 SparkSession 就够用了所以临时视图更常用。全局视图带来的额外上下文切换和潜在命名冲突通常不值得。我建议不到万不得已只用createOrReplaceTempView。3.4 使用 DataFrame API 和 SQL 混合查询日常开发里也不是纯写 SQL 不可。DataFrame API 和 SQL 其实可以混着用。比如先用 DataFrame API 做过滤再注册成视图跑 SQLfiltered_df trip_df.filter(trip_df[amount] 10) filtered_df.createOrReplaceTempView(high_value_trips) spark.sql(SELECT driver_id, SUM(amount) FROM high_value_trips GROUP BY driver_id ORDER BY SUM(amount) DESC).show()这样的好处是既能用代码做精细的列控制又能享受 SQL 描述复杂聚合逻辑的简洁性。很多实际项目里最耗时的 JOIN 和去重逻辑用 SQL 一眼能看懂而一些朴素的条件拼接、随机采样这些代码里写也顺手。你不需要二选一两边都是 Spark 的一等公民。3.5 把查询结果保存成新视图供下游使用除了展示结果你还可以把某个 SQL 查询的结果注册成另一个视图形成多级管道。举个例子daily_report spark.sql( SELECT order_date, driver_id, COUNT(*) AS trip_count, SUM(amount) AS revenue FROM trip_table WHERE status completed GROUP BY order_date, driver_id ) daily_report.createOrReplaceTempView(daily_driver_report)然后下游再写一条 SQL 去查每个司机最高日流水。这种组织方式非常像数据仓库里的分层思路——先做明细汇总再做业务层统计。配合createDataFrame加载外部数据几乎可以无缝模拟一套小型数仓。4. 常见问题与排查技巧实录4.1createDataFrame报错 “TypeError: Can not merge type ...”最典型的错误是 Python 列表里字段类型不统一。比如data [(Alice, 30), (Bob, 50)]第一行 age 是 int第二行是 strSpark 无法合并推断一个通用类型直接报 TypeError。解决办法是把它们统一成字符串或者显式定义 Schema 并转成对应类型。实践中我从各种数据库导出的数据经常有这种形态所以写 ETL 的第一步永远都是做类型对齐。4.2 临时视图名称冲突createTempView会在同一个 SparkSession 里创建一个同名的视图如果重名且没有 REPLACE会直接报错。想省事就用createOrReplaceTempView。这个细节初学容易忽略总以为视图是自动覆盖的结果多次跑同一个 cell 就报诡异错误。4.3collect()到 Driver 造成 OOM注意用show()实际上是调用了 collect 的一部分。如果你写df.collect()那会把全部数据收集到 Driver 进程的内存里数据一多Driver 直接 OutOfMemory。这是初学者最喜欢踩的坑网上搜“spark 内存”也常常提到这个。正确姿势是使用take(n)、show(n)或者foreach在 Executor 端消费。4.4 SQL 里NULL与空字符串的纠缠createDataFrame时如果某个字段是 Python 的NoneSpark 会把它当作NULL。但如果你读 CSV 拿到的是空字符串Spark 不会自动视为NULL。这直接导致IS NULL查不出空字符串聚合函数如COUNT(column)也会把空字符串算进去。处理办法是读文件时设置nullValue参数或者在 DataFrame 里用when(col(x) , lit(None))做转换。很多 SQL 清洗项目里大量时间都花在“把空字符串变成真正的 NULL”上。4.5 建表时没显式 Schema 导致后续性能差刚才提过inferSchema会额外触发一次任务在数据量大时等于白读一遍。不要小看这个点。对于一天的日志量少一次全量扫描可能意味着几分钟到十几分钟的差距。我在实际调优中但凡数据在几百 MB 以上就一定手动写 Schema。宁可代码长一点也不靠自动推断。4.6 慢 SQL 优化先看执行计划热词里有“慢 sql 优化”。Spark SQL 的查询慢很多时候不是 SQL 写得不对而是数据倾斜、扫描太多、没有分区。你可以在执行前spark.sql(EXPLAIN SELECT ...).show(truncateFalse)查看 Spark 物理执行计划确认你的过滤条件是不是下推到数据源了。如果发现明明在过滤大表却还是全量扫描那要考虑列裁剪和分区裁剪是否生效。createDataFrame创建的小 DataFrame 没这种问题但生产环境从 Parquet 读进来时就不一样了。5. 工具选型与实践建议5.1 在 PySpark 和 Scala 之间怎么选关于createDataFrame和 SQL 结合写当前大环境里 PySpark 的社区热度明显更高。原因很简单Python 开发的便利性和数据分析生态太强了。Scala 的性能有细微优势尤其是复杂类型操作和 UDF 上但它语法相对硬核迭代速度慢。我个人的建议是数据工程师用 PySpark 作为主力必要时用 Scala 写底层 UDF 或性能敏感模块两边可以无缝混用DataFrame 和 SQL 的语义完全一致。5.2 小集群自建与测试时的配置建议热词里提到 spark 集群搭建。如果你是自建小集群练手节点数是 3 到 5 台建议配置上重点考虑spark.executor.memory、spark.executor.cores还有spark.sql.shuffle.partitions。比如三台 16G 机器每台跑 2 个 executor每个 executor 4Gshuffle 分区按总核心数去设置别盲目用默认的 200。本地写代码验证时可以把 master 设成local[*]方便调试。真正提交到集群时建议用spark-submit指定 master避免代码里硬编码 local 模式否则上了集群还在本地跑等于没充分利用分布式资源。5.3 关于热词里的其他话题热词中还有“spark 数据分析案例”“农产品价格数据清洗”等其实都和createDataFrame有关。数据分析案例的核心就是读数据、清洗成 DataFrame、注册视图、SQL 聚合。农产品价格数据清洗面对的往往是缺失值、异常值、单位不一致这些通通可以通过先createDataFrame自定义一张标准表再优雅地用 SQL 或者 DataFrame API 清洗。理解了createDataFrame这一层就能理解为什么 Spark SQL 能成为这类大数据项目的利器。用一张表来描述createDataFrame模式实操时的选择依据对新手挺友好场景推荐方式原因调试/单测造少量数据Python 列表 spark.createDataFrame(data, schema)清楚直观无多余读取读取 CSV/JSON 大批量spark.read.schema(schema).csv()避免二次扫描和类型误判从 RDD 转换rdd.toDF()或createDataFrame(rdd, schema)兼容旧逻辑适合迁移需要精确 Decimal/金额显式StructType定义字段防止 Double 浮点误差需要跨 Session 共享createGlobalTempView保持全局可查5.4 一个加速 SQL 的小细节提前 cache()当你多次针对同一个 DataFrame 做不同 SQL 查询时如果不做缓存Spark 每次都会重新读数据源或重算上游转换链。比较快的处理方式是trip_df.cache() trip_df.createOrReplaceTempView(trip_table)这样第一次触发计算之后数据会留在 Executor 内存里后续不管跑多少条 SQL 都会命中缓存速度提升明显。需要注意如果只做了createOrReplaceTempView而没有真正触发一次 action比如马上执行spark.sql(...).show()那也会触发一次完整链路计算然后缓存生效。cache()只影响后续 action并不会在你调用那一刻主动去 cache。初学时容易以为缓存是立即的其实它还是 lazy 的。6. 一段完整示例代码供直接抄作业最后给你一个可以直接跑的 PySpark 脚本。它综合了创建 DataFrame、临时视图、SQL 查询和缓存优化。如果你手头有 Spark 环境本地直接跑一遍全过程不到十秒。from pyspark.sql import SparkSession from pyspark.sql.types import StructType, StructField, StringType, DoubleType spark SparkSession.builder \ .appName(CreateDataFrame2SQL) \ .master(local[*]) \ .config(spark.sql.shuffle.partitions, 4) \ .getOrCreate() # 显式定义 schema schema StructType([ StructField(order_date, StringType(), True), StructField(driver_id, StringType(), True), StructField(amount, DoubleType(), True), StructField(status, StringType(), True), ]) data [ (2024-05-01, driver_001, 12.5, completed), (2024-05-01, driver_002, 8.0, cancelled), (2024-05-02, driver_001, 15.2, completed), (2024-05-02, driver_003, 30.9, completed), ] df spark.createDataFrame(data, schema) df.cache() df.createOrReplaceTempView(trip_table) result spark.sql( SELECT order_date, COUNT(*) AS total_orders, SUM(CASE WHEN status completed THEN amount ELSE 0 END) AS completed_amount FROM trip_table GROUP BY order_date ORDER BY order_date ) result.show()跑出来的结果应该是一张带日期、总订单数、完成金额的汇总表。这已经是一个最基础的数据分析师视角的产出。如果你要继续深入不妨把订单数据换成 CSV 或者 JSON 文件将createDataFrame换成spark.read.schema(schema).csv(...)再配合同样的createOrReplaceTempView步骤。这样你在生产环境里处理网约车订单、农产品价格、日志清洗的逻辑基本都能直接套用。我个人在实际操作中最大的体会是createDataFrame虽然只是 Spark 很多 API 中的一个但它是理解整个 Spark SQL 的钥匙。把这个 API 用熟了你自然就理解了 Schema 的含义、临时视图的机制、SQL 和 DataFrame 的转换关系再去读那些“Spark 数据分析案例”“数据清洗”的复杂项目都会顺畅许多。所以别嫌它基础基础的地方往往正是所有高阶用法的出发点。
返回列表