如何在PySpark中拆分含CSV字符串的列并解析为结构化DataFrame?
问题描述
我有如下DataFrame:
+---+---------------------+ | id| csv| +---+---------------------+ | 1|a,b,c\n1,2,3\n2,3,4\n| | 2|a,b,c\n3,4,5\n4,5,6\n| | 3|a,b,c\n5,6,7\n6,7,8\n| +---+---------------------+
希望拆分字符串类型的csv列,最终得到如下DataFrame:
+--+--+--+ | a| b| c| +--+--+--+ | 1| 2| 3| | 2| 3| 4| | 3| 4| 5| | 4| 5| 6| | 5| 6| 7| | 6| 7| 8| +--+--+--+
查阅from_csv的文档可知,它仅能处理单行CSV字符串,因此该方法不可用。
我尝试循环遍历DataFrame的每一行,提取并解析CSV字符串后合并结果,代码如下:
rows = df.collect() for (i, row) in enumerate(rows): data = row['csv'] data = data.split('\\n') rdd = spark.sparkContext.parallelize(data) df_row = (spark.read .option('header', 'true') .schema('a int, b int, c int') .csv(rdd)) if i == 0: df_new = df_row else: df_new = df_new.union(df_row) df_new.show()
但这种方法效率极低,有没有更优的实现方式?
高效实现方案
绝对不要用collect()将数据拉取到Driver端循环处理,这会严重影响性能,尤其是大数据量场景。推荐使用Spark原生的分布式字符串处理函数完成,全程在集群节点并行执行:
方法一:纯分布式处理(推荐)
通过拆分、展开、过滤、类型转换等步骤完成,完全不需要拉取数据到本地:
from pyspark.sql import functions as F from pyspark.sql.types import IntegerType df_result = (df # 将csv列按换行符拆分成数组,再展开为单行记录 .select(F.explode(F.split(F.col("csv"), "\n")).alias("line")) # 过滤空行和表头行 .filter(F.col("line") != "").filter(F.col("line") != "a,b,c") # 拆分每行字符串并转换为整数类型 .select( F.split(F.col("line"), ",").getItem(0).cast(IntegerType()).alias("a"), F.split(F.col("line"), ",").getItem(1).cast(IntegerType()).alias("b"), F.split(F.col("line"), ",").getItem(2).cast(IntegerType()).alias("c") ) ) df_result.show()
方法二:拼接完整CSV后读取(适合小数据量)
如果数据量不大,可以先拼接所有CSV内容再一次性读取,比循环处理效率高:
from pyspark.sql import functions as F # 拼接所有csv字符串为完整的CSV内容 full_csv_content = "\n".join([row.csv for row in df.select("csv").collect()]) # 并行化后读取CSV df_result = spark.read.csv( spark.sparkContext.parallelize([full_csv_content]), header=True, schema="a int, b int, c int" ) df_result.show()
两种方法中,方法一完全分布式执行,性能最优,适合大数据量场景;方法二仅适合数据量较小的情况,避免频繁创建小DataFrame和union操作带来的开销。
内容的提问来源于stack exchange,提问作者Kai Roesner
相关产品推荐
相关产品推荐

