全字符串类型pandas dataframe按指定schema转spark dataframe的方法
报错原因
TypeError: an integer is required (got type str) 报错的核心逻辑是:直接调用
spark.createDataFrame()传入pandas DataFrame时,Spark不会自动对pandas端的原始数据做类型转换,只会严格校验pandas列的实际数据类型是否和指定Spark schema匹配。你的pandas列全为字符串类型,而目标schema指定为整数类型,类型不匹配因此直接抛出错误。
解决方法
方法1:Spark端显式类型转换(推荐,适合大数据量场景)
先将pandas DataFrame转为全字符串类型的Spark临时表,再按照目标schema批量做类型转换,转换过程可控,还可以自定义非法值的兜底逻辑:
import pandas as pd from pyspark.sql.types import * import pyspark.sql.functions as F # 原始pandas数据 d = {'col1': ['1', '2'], 'col2': ['3', '4']} df = pd.DataFrame(data=d) # 第一步:直接生成字符串类型的Spark临时表 spark_temp_df = spark.createDataFrame(df) # 第二步:定义目标schema target_schema = StructType([ StructField('col1', IntegerType(), True), StructField('col2', IntegerType(), True) ]) # 第三步:遍历目标schema批量做类型转换,无需手动逐个写字段 sparkDf = spark_temp_df.select( *[F.col(field.name).cast(field.dataType).alias(field.name) for field in target_schema.fields] ) display(sparkDf)
方法2:pandas端提前转换类型(适合小数据量场景)
先在pandas侧把列类型转为和目标Spark schema匹配的类型,再直接创建Spark DataFrame:
import pandas as pd from pyspark.sql.types import * d = {'col1': ['1', '2'], 'col2': ['3', '4']} df = pd.DataFrame(data=d) # pandas端提前完成类型转换 df['col1'] = df['col1'].astype(int) df['col2'] = df['col2'].astype(int) # 再指定schema创建Spark DF schema = StructType([ StructField('col1', IntegerType(), True), StructField('col2', IntegerType(), True) ]) sparkDf = spark.createDataFrame(df, schema = schema) display(sparkDf)
内容的提问来源于stack exchange,提问作者ARCrow
相关产品推荐
相关产品推荐

