PySpark中toDF()与createDataFrame()的异常行为问询
为什么PySpark中
parallelize(...).toDF()和createDataFrame()结果不一样? 这是刚接触Spark很容易踩的新手坑,我来给你把背后的逻辑掰明白:
核心差异:Schema推断的逻辑不同
Spark创建DataFrame时,Schema推断的规则会因为调用方式不同而完全不一样:
1. spark.parallelize(dic).toDF()的行为
当你把字典列表转成RDD再调用toDF()时,Spark会只拿RDD的第一个元素来推断Schema。你的RDD第一个元素是{"a":1},所以Spark直接认定这个DataFrame只有a这一列。后面的{"b":2}和{"c":3}没有a键,对应位置就只能填充null,这就是你看到只有a列且后两行是null的原因。
2. spark.createDataFrame(dic)的行为
而createDataFrame()方法会遍历整个输入数据集,收集所有出现过的键作为列名,生成包含所有列的Schema。所以它能识别出a、b、c三个列,每个字典只填充自己有的键对应的值,没有的就补null,这就是你看到三列完整结果的原因。
怎么让parallelize(...).toDF()得到相同的结果?
如果你一定要用RDD转DF的方式,有两种靠谱的解决办法:
方法一:手动指定Schema
最稳妥的方式是自己定义完整的Schema,彻底避免依赖自动推断:
from pyspark.sql.types import StructType, StructField, IntegerType # 定义包含所有列的Schema schema = StructType([ StructField("a", IntegerType(), nullable=True), StructField("b", IntegerType(), nullable=True), StructField("c", IntegerType(), nullable=True) ]) # 用指定的Schema转DF spark.parallelize(dic).toDF(schema).show()
方法二:让Spark全局推断Schema
你可以把RDD转成列表后传给createDataFrame(数据量小的时候适用),或者直接用createDataFrame包裹parallelize的结果,本质和直接调用createDataFrame(dic)完全一样:
spark.createDataFrame(spark.parallelize(dic)).show()
额外提醒
当你的数据集很大时,createDataFrame的全局Schema推断会有一定性能开销,因为要遍历所有数据。这时候手动指定Schema不仅能避免这种坑,还能提升创建DataFrame的效率。
内容的提问来源于stack exchange,提问作者Joel Lee
相关产品推荐
相关产品推荐

