如何正确使用Spark的from_csv函数解析带表头的CSV字符串?
问题原因
from_csv函数是用来解析单条CSV格式字符串的,你的输入是包含表头+数据行的多行字符串,加上设置了header=True,它会把整个字符串的第一行当成表头,剩下的所有内容(包括换行后的真实数据)当成单个字段的值,所以才会出现uid字段取了表头的uid,datetime因为匹配不上数据而返回null,quality字段把剩下的所有内容都塞进去的错误结果。
解决步骤
核心是先把多行CSV字符串拆分成单独的行,再处理:
- 将单个多行字符串按换行符拆分成数组,再用
explode展开成单独的行 - 过滤掉拆分后产生的空行
- 区分表头行和数据行,只对数据行使用
from_csv解析,或者直接用Spark的CSV读取器处理拆分后的行
修改后的代码
from pyspark.sql import SparkSession from pyspark.sql.functions import from_csv, split, explode, trim from pyspark.sql.types import StructType, StructField, StringType, TimestampType # 初始化Spark spark = SparkSession.builder.appName("Extract CSV from string").master("local").getOrCreate() # 模拟Kafka过来的DataFrame csvString = [("uid,datetime,quality\nfcc5382b-201b-4c10-9db9-e5de7f8c1d32,2023-07-04 11:32:24,passed\n",)] df = spark.createDataFrame(csvString, ("value",)) # 1. 拆分多行字符串为单独的行并展开 df = df.withColumn("line", explode(split(df.value, "\n"))) # 2. 过滤空行和纯空格行 df = df.filter(trim(df.line) != "") # 定义CSV schema csvSchema = StructType([ StructField("uid", StringType()), StructField("datetime", TimestampType()), StructField("quality", StringType()) ]) # 3. 过滤掉表头行,只解析数据行 header = df.select("line").first()[0] df_data = df.filter(df.line != header) # 解析数据行 result = df_data.select( from_csv(df_data.line, csvSchema).alias("csv") ).collect() print(result)
另一种更简洁的方法(直接用Spark CSV读取器)
如果场景允许,也可以把拆分后的行转成RDD,直接用spark.read.csv读取:
# 接上面拆分后的df lines_rdd = df.select("line").rdd.map(lambda x: x[0]) df_result = spark.read.csv(lines_rdd, header=True, schema=csvSchema) print(df_result.collect())
这两种方法都能得到你预期的结果:
[Row(csv=Row(uid='fcc5382b-201b-4c10-9db9-e5de7f8c1d32', datetime=datetime.datetime(2023, 7, 4, 11, 32, 24), quality='passed'))]
内容的提问来源于stack exchange,提问作者Thomas DevAI
相关产品推荐
相关产品推荐

