使用PySpark读取CSV时能否仅覆盖单个列的推断数据类型
解决方案
需求完全可行,你当前代码的问题是:手动传入仅包含单列的Schema时,Spark会直接使用该Schema作为全表读取规则,自动忽略其余未声明的列,同时inferschema参数会被显式传入的Schema覆盖失效。
最优实现方式
方式一:直接转换已读取的DataFrame(绝大多数场景优先选)
先开启inferschema读取全量数据,再单独对类型错误的列做类型转换,不需要二次读取文件,代码最简洁:
from pyspark.sql.types import StringType df = spark.read.format('com.databricks.spark.csv') \ .option('delimiter', ',') # 注意你原代码这里参数名写错了,是delimiter不是delimited .option('header', 'true') \ .option('inferschema', 'true') \ .load('dbfs:/FileStore/some.csv') # 仅修改指定列的类型 df_final = df.withColumn("column_one_of_many", df["column_one_of_many"].cast(StringType()))
方式二:预生成修改后Schema读取(大文件场景优先选)
如果CSV文件体积极大,不想在类型转换阶段做全量计算,可以先预推断Schema、修改指定列类型后,再用完整Schema一次性读取正确类型的数据:
from pyspark.sql.types import StringType # 采样少量数据推断Schema,避免全量扫描浪费资源 df_temp = spark.read.format('com.databricks.spark.csv') \ .option('delimiter', ',') \ .option('header', 'true') \ .option('inferschema', 'true') \ .load('dbfs:/FileStore/some.csv').limit(100) # 仅采样前100行足够推断大部分列类型 # 修改目标列的类型定义 modified_schema = df_temp.schema for field in modified_schema.fields: if field.name == "column_one_of_many": field.dataType = StringType() break # 用修改后的完整Schema读取全量数据 df_final = spark.read.format('com.databricks.spark.csv') \ .option('delimiter', ',') \ .option('header', 'true') \ .schema(modified_schema) \ .load('dbfs:/FileStore/some.csv')
内容的提问来源于stack exchange,提问作者Zhao Li
相关产品推荐
相关产品推荐

