Pandas设置数据类型后转换为PySpark DataFrame时全部字段变为String类型的问题求助
解决Pandas DataFrame转Spark DataFrame时类型被强制转为String的问题
我太懂你这种明明设置了类型却不生效的挫败感了!咱们先揪出问题根源,再一步步解决:
核心问题1:Pandas的类型转换没真正生效
你这段代码里的astype调用犯了Pandas新手常踩的坑:
BDtable_FINAL.astype({ 'dateinsert' : 'int64', 'date': 'datetime64' })
astype是返回一个新的DataFrame对象,不会原地修改原数据,所以你必须把结果重新赋值给原变量。另外dateinsert是字符串格式生成的,直接转int64会报错,得先确保它是可转换的数值字符串,date列最好指定格式避免自动解析出错。修改后的类型转换代码如下:
# 先把dateinsert从字符串转成int64 BDtable_FINAL['dateinsert'] = BDtable_FINAL['dateinsert'].astype('int64') # 指定格式转换date列为datetime,避免解析异常 BDtable_FINAL['data'] = pd.to_datetime(BDtable_FINAL['data'], format='%Y-%m-%d %H:%M:%S')
核心问题2:转换函数的细节优化
你的pandas_to_spark思路没问题,但有两个小细节拖了后腿:
- Pandas的
float64对应Spark的DoubleType()(Spark的FloatType是单精度,和Pandas的双精度不匹配) - Databricks里更推荐用
spark上下文替代sqlContext
优化后的转换函数:
from pyspark.sql.types import ( StructType, StructField, StringType, LongType, IntegerType, DoubleType, TimestampType ) import pandas as pd def equivalent_type(f): if f == 'datetime64[ns]': return TimestampType() elif f == 'int64': return LongType() elif f == 'int32': return IntegerType() elif f == 'float64': return DoubleType() # 修正为DoubleType匹配Pandas双精度 else: return StringType() def define_structure(string, format_type): try: typo = equivalent_type(format_type) except: typo = StringType() return StructField(string, typo, nullable=True) # 加上nullable参数更灵活 def pandas_to_spark(pandas_df): columns = list(pandas_df.columns) types = list(pandas_df.dtypes) struct_list = [] print(f"Pandas原数据类型: {dict(zip(columns, types))}") # 打印清晰的类型日志 for column, typo in zip(columns, types): struct_list.append(define_structure(column, typo)) p_schema = StructType(struct_list) return spark.createDataFrame(pandas_df, p_schema) # Databricks用spark上下文
验证修正后的完整流程
现在重新跑一遍你的代码:
import pandas as pd from datetime import datetime BDtable_FINAL = pd.DataFrame({'data': ['0001-01-01 00:00:00', '2020-02-02 00:00:00', '2021-01-01 00:00:00']}) BDtable_FINAL = BDtable_FINAL[~BDtable_FINAL['data'].isin(['0001-01-01 00:00:00'])] datainsert = datetime.now().strftime('%Y%m%d%H%M') dateinsert = datainsert[:8] + '0000' BDtable_FINAL.insert(loc=0, column='dateinsert', value=dateinsert) # 修复后的类型转换 BDtable_FINAL['dateinsert'] = BDtable_FINAL['dateinsert'].astype('int64') BDtable_FINAL['data'] = pd.to_datetime(BDtable_FINAL['data'], format='%Y-%m-%d %H:%M:%S') # 转Spark DataFrame spark_df = pandas_to_spark(BDtable_FINAL) spark_df.printSchema()
运行后printSchema()应该会输出你想要的结果:
root |-- dateinsert: long (nullable = true) |-- data: timestamp (nullable = true)
为啥之前的方法没起效?
- Koalas的
to_spark()完全依赖Pandas的类型信息,如果Pandas本身类型没设置对,转过去自然还是String - 直接在Spark里设类型时,若源数据是String格式(比如你没修正Pandas的类型转换),Spark没法自动推断出正确类型,必须先清洗数据
另外悄悄说一句:在Databricks里其实不用自己写转换函数,只要Pandas的类型正确,直接用spark.createDataFrame(BDtable_FINAL)就能自动推断出对应类型啦!
内容的提问来源于stack exchange,提问作者Luís Gustavo
相关产品推荐
相关产品推荐

