PySpark中如何高效迭代解析df.collect()的特殊格式数据?
问题描述
我在Databricks中用PySpark读取S3上的非标准CSV文件,因为表头和数据位于文件中间,所以用文本格式读取:
df = (spark.read .format("text") .option("mode", "PERMISSIVE") .option("header", "false") .option("inferSchema","false") .load("s3://path/to/your/file") )
df.collect()返回结果如下:
[Row(value=':::::'), Row(value='./some_file_name.csv'), Row(value=':::::'), Row(value='col_one|col_two|col_three'), Row(value='1|2|3'), Row(value='./another_file.csv'), Row(value='more stuff|of|varying|length')]
我需要提取包含col_one|col_two|col_three的表头行,以及紧随其后的1|2|3这类数据行。我可以转成列表后用正则处理,但想知道更符合PySpark风格的实现方式,已知数据量不大,无需担心资源问题。
解决方案
以下是纯PySpark风格的实现方案,通过行号标记、筛选和字符串拆分完成需求:
- 添加行号:给每行分配自增行号,用于定位表头与数据行的关联关系
- 定位表头行:筛选出目标表头行,记录其行号
- 筛选目标行:保留表头行及紧随其后的一行数据
- 转换为结构化DataFrame:拆分字符串列,映射为对应表头的字段
具体代码
from pyspark.sql import Window import pyspark.sql.functions as F # 1. 为每行添加自增行号 df_with_row_num = df.withColumn("row_num", F.monotonically_increasing_id()) # 2. 获取表头行的行号(精确匹配表头内容) header_row_num = df_with_row_num.filter(F.col("value") == "col_one|col_two|col_three").select("row_num").first()[0] # 3. 筛选表头行和紧随的数据行 target_rows = df_with_row_num.filter( (F.col("row_num") == header_row_num) | (F.col("row_num") == header_row_num + 1) ) # 4. 提取表头列名 header_cols = target_rows.filter(F.col("row_num") == header_row_num).select(F.split("value", "\\|").alias("cols")).first()[0] # 5. 拆分数据行并映射为结构化DataFrame final_df = target_rows.filter(F.col("row_num") == header_row_num + 1) \ .select(F.split("value", "\\|").alias("data")) \ .select([F.col("data")[i].alias(header_cols[i]) for i in range(len(header_cols))]) # 查看结果 final_df.show()
适配动态表头的情况
如果表头内容不是固定字符串,而是符合特定正则格式(比如以col_开头的多字段行),可以用正则匹配定位表头行:
# 用正则匹配包含至少三个以col_开头字段的行 header_row_num = df_with_row_num.filter( F.regexp_extract("value", "^col_\w+\|col_\w+\|col_\w+$", 0) != "" ).select("row_num").first()[0]
这种方案完全基于PySpark API实现,无需将数据转成Python列表处理,同时因为数据量小,不会有性能瓶颈。
内容的提问来源于stack exchange,提问作者Chuck
相关产品推荐
相关产品推荐

