自定义Schema实现pandas转PySpark DataFrame报错排查咨询
问题原因分析
- 第一个错误:
StructField初始化要求传入3个位置参数,分别为字段名、数据类型、是否允许为空,你在列表推导式里直接写StructField(i),只把三元组i作为第一个参数传入,没有拆分三个元素对应传参,自然会提示缺少dataType参数。 - 第二个错误:你把数据类型存成了字符串格式的
"IntegerType()"、"FloatType()",不是PySpark原生的类型对象,且eval的使用位置完全错误,套在StructType外层不会生效,反而会引发额外的语法错误。 - 第三个逻辑错误:循环内判断类型时固定取了
df.y.dtype,没有读取当前循环列的类型,会导致所有字段的类型都和y字段保持一致,不符合自定义Schema的需求。
修正后代码
from pyspark.sql.types import StructType, StructField, IntegerType, FloatType, StringType dtype_l, name_l, true_l = [],[],[] for col in df.columns: name_l.append(col) true_l.append(True) # 读取当前循环列的类型做判断 current_dtype = df[col].dtype if current_dtype == 'int64': dtype_l.append(IntegerType()) # 直接传入类型对象,不要加引号 elif current_dtype == 'float64': dtype_l.append(FloatType()) # 可按需补充其他类型映射规则,比如字符串类映射为StringType else: dtype_l.append(StringType()) # 拆分三元组元素分别传给StructField对应参数 mySchema = StructType([StructField(name, dtype, nullable) for name, dtype, nullable in zip(name_l, dtype_l, true_l)]) mySchema
如果不需要做特殊的类型规则映射,也可以直接用PySpark内置方法完成转换,代码更简洁:
spark_df = spark.createDataFrame(df)
内容的提问来源于stack exchange,提问作者Ekat Sim
相关产品推荐
相关产品推荐

