如何使用PySpark提取文本中的指定列并构建DataFrame?
解决方案
你的文本文件结构是表头 + 多组数据在同一行,原始代码只提取了表头部分,没有处理后续的多组数据。以下是修正后的实现步骤和代码:
问题分析
你的文件内容是一行数据:Name|age|course|"A"|20|"Science"|"B"|23|"Math"|"C"|25|"Englsih",前3个元素是表头,后面每3个元素为一组用户数据。原始代码只取了数组前3位(表头),没有拆分并展开后续的多组数据。
修正代码(两种实现方式)
方式一:使用UDF分组数据
from pyspark.sql import functions as F from pyspark.sql.types import ArrayType, StringType # 1. 读取文本文件(用text方式保留整行内容) df = spark.read.text("test.txt") # 2. 拆分整行数据为数组 df_split = df.withColumn("all_data", F.split(F.col("value"), "\\|")) # 3. 分离表头和数据部分(从第4个元素开始是用户数据,slice是1-based索引) df_with_parts = df_split.select( F.slice(F.col("all_data"), 4, F.size(F.col("all_data")) - 3).alias("data_values") ) # 4. 定义UDF将数据数组按每3个元素分组 def group_data(arr): return [arr[i:i+3] for i in range(0, len(arr), 3)] group_udf = F.udf(group_data, ArrayType(ArrayType(StringType()))) # 5. 分组后展开为多行,再提取字段 df_final = df_with_parts.withColumn("grouped_rows", group_udf(F.col("data_values")))\ .withColumn("row_data", F.explode(F.col("grouped_rows")))\ .select( F.col("row_data")[0].alias("Name"), F.col("row_data")[1].cast("int").alias("Age"), F.col("row_data")[2].alias("Course") ) # 查看结果 df_final.show()
方式二:使用Spark内置函数(无需UDF)
from pyspark.sql import functions as F # 1. 读取文本文件 df = spark.read.text("test.txt") # 2. 拆分整行并分离数据部分 df_split = df.withColumn("all_data", F.split(F.col("value"), "\\|"))\ .select(F.slice(F.col("all_data"), 4, F.size(F.col("all_data")) - 3).alias("data_values")) # 3. 生成索引序列,按每3个元素拆分并展开 df_final = df_split.withColumn( "index", F.explode(F.sequence(0, F.size(F.col("data_values")) - 1, 3)) ).withColumn( "row_data", F.slice(F.col("data_values"), F.col("index") + 1, 3) # slice是1-based索引 ).select( F.col("row_data")[0].alias("Name"), F.col("row_data")[1].cast("int").alias("Age"), F.col("row_data")[2].alias("Course") ) df_final.show()
执行结果
两种方式都会输出正确的DataFrame:
+----+---+-------+ |Name|Age| Course| +----+---+-------+ | "A"| 20|Science| | "B"| 23| Math| | "C"| 25|Englsih| +----+---+-------+
内容的提问来源于stack exchange,提问作者marton mar suri
相关产品推荐
相关产品推荐

