Databricks中如何将字符串转为PySpark StructType供DataFrame读取
问题场景
业务方提供了一份CSV文件,其中记录了MATERIAL_MASTER(物料主数据)的字段详情,字段定义参考截图:
在Databricks中创建Notebook读取该字段定义文件并生成对应Schema,初始实现代码如下:
import pandas as pd rows=[] material_master_schema_df = spark.read.csv("/mnt/dentanalyticsdatalake/main-data/POC/RAW/MaterialMasterSchema.csv",header = True) material_master_schema_df= material_master_schema_df.toPandas() rows = material_master_schema_df.values.tolist() #print(rows) var= [ 'StructField'+str(tuple(i)) for i in rows] #print(var) fields = [] for i in var: fields.append(str(str(i.split(',')[0])+','+(i.split(',')[1].split("'")[1]+'Type()'+','+(i.split(',')[2].split("'")[1]+')')))) #print(fields) Current_Schema = 'StructType'+'('+str(fields)+')' material_master_schema = Current_Schema.replace('"','')
故障现象
需要在另一个Notebook中通过%run命令调用上述Schema生成Notebook,创建DataFrame读取数据时传入该schema,调用代码如下:
%run ./Schema_Creation_Notebook MATERIAL_MASTER_TXT_DF = spark.read.csv("/xxx/file.txt",header = True, sep='\t',schema = material_master_schema )
执行时抛出ParseException异常,排查确认根因:material_master_schema变量类型为字符串(str),但接口要求传入pyspark.sql.types.StructType类型对象。
解决方案
原方案通过配置文件自动生成Schema的思路完全可行,问题出在实现逻辑通过字符串拼接构造Schema内容,没有生成Spark可识别的StructType实例。
- 推荐实现:直接导入PySpark类型类,基于读取到的字段配置逐行生成StructField对象,最终组装为合法的StructType实例,完全避免字符串拼接的问题。修正后的Schema生成代码如下:
from pyspark.sql.types import StructType, StructField, StringType, IntegerType, DoubleType, BooleanType, DateType, TimestampType # 建立字段类型名到PySpark类型对象的映射,可根据业务实际使用的类型扩展 type_map = { "string": StringType(), "int": IntegerType(), "double": DoubleType(), "boolean": BooleanType(), "date": DateType(), "timestamp": TimestampType() } # 读取字段定义CSV,无需转pandas可直接遍历 schema_conf_df = spark.read.csv( "/mnt/dentanalyticsdatalake/main-data/POC/RAW/MaterialMasterSchema.csv", header=True ) field_list = [] for row in schema_conf_df.collect(): # 按CSV列顺序依次取:字段名、字段类型、是否允许为空 col_name = row[0] col_type = type_map[row[1].strip().lower()] is_nullable = True if row[2].strip().lower() == "true" else False field_list.append(StructField(col_name, col_type, is_nullable)) # 直接生成StructType实例 material_master_schema = StructType(field_list)
- 临时兼容方案:如果要沿用原有字符串拼接的逻辑,可以提前导入PySpark所有类型类,通过
eval()函数将拼接完成的Schema字符串转为StructType实例,但该方式存在代码注入风险,生产环境不建议使用:
from pyspark.sql.types import * # material_master_schema为原有逻辑生成的Schema字符串 material_master_schema = eval(material_master_schema)
修正完成后,通过%run调用该Notebook拿到的material_master_schema即为合法的StructType对象,直接传入读文件接口即可正常运行,不会再出现类型不匹配的异常。
内容的提问来源于stack exchange,提问作者sayan nandi
相关产品推荐
相关产品推荐

