PySpark无法将数据映射到指定Schema的问题求助
问题根源分析
你踩了一个常见的格式匹配坑——你的数据是每行一个完整JSON对象的格式,但你一直在用csv格式的读取器加载数据。CSV读取器是按指定分隔符(你设置的是\t)拆分列的,可你的数据里根本没有制表符,所以整个JSON字符串被当作了单独一列,不管怎么指定schema都没法正确拆分字段。
解决方案
根据你的需求,有几种可行的处理方式:
1. 直接使用JSON读取器加载(最简单)
Spark原生支持读取每行一个JSON对象的数据集,直接用json格式读取器即可,它会自动推断嵌套的schema:
df = spark.read.json(path) df.show(1, False)
执行后你会看到data是一个嵌套结构体列,uploadedDate是顶层列,之后可以用下面的代码把嵌套字段展开成单独列:
df_expanded = df.select("data.probability", "data.customerId", "data.region", "uploadedDate") df_expanded.show(1, False)
2. 手动指定严格的Schema(生产环境推荐)
如果你需要严格控制字段类型(避免自动推断出错),得构建包含嵌套结构体的Schema,完全匹配原始JSON的层级结构:
from pyspark.sql.types import StructType, StructField, DoubleType, LongType, StringType SCHEMA = StructType([ StructField('data', StructType([ StructField('probability', DoubleType(), True), # 原始数据是小数,用DoubleType而非LongType StructField('customerId', LongType(), True), StructField('region', StringType(), True) ]), True), StructField('uploadedDate', LongType(), True) # 原始数据是毫秒级时间戳,用LongType而非StringType ]) df = spark.read.schema(SCHEMA).json(path) # 展开嵌套字段 df_expanded = df.select("data.*", "uploadedDate") df_expanded.show(1, False)
注意:你之前定义的Schema完全不匹配原始数据的层级(没考虑data的嵌套结构,还错误地把probability命名为probabilityMale、把时间戳设为StringType),这也是schema失效的关键原因之一。
3. 若必须用CSV读取器(特殊场景)
如果你的数据确实是存储为CSV格式(每行只有一个JSON字符串列),可以先读取成单列,再用from_json函数解析JSON字符串:
# 先读取成单列DataFrame df_raw = spark.read.format('csv').option('header','false').load(path) # 定义正确的JSON Schema(和方案2一致) SCHEMA = StructType([ StructField('data', StructType([ StructField('probability', DoubleType(), True), StructField('customerId', LongType(), True), StructField('region', StringType(), True) ]), True), StructField('uploadedDate', LongType(), True) ]) # 解析JSON字符串 from pyspark.sql.functions import from_json df_parsed = df_raw.select(from_json(df_raw._c0, SCHEMA).alias("parsed_data")) # 展开字段得到最终结果 df_final = df_parsed.select("parsed_data.data.*", "parsed_data.uploadedDate") df_final.show(1, False)
内容的提问来源于stack exchange,提问作者Saurabh Deshpande
相关产品推荐
相关产品推荐

