PySpark如何将逗号分隔的字符串字段拆分为多行生成新RDD
实现方案
你可以根据当前使用的是DataFrame还是原生RDD选择对应实现方式,优先推荐DataFrame API,写法更简洁性能也更优。
方案1:基于DataFrame实现(推荐)
核心使用split函数将字符串字段拆分为数组,再用explode函数将数组元素展开为独立行:
from pyspark.sql import functions as F # 拆分new_column为数组,再炸开数组生成多行 df_result = df.withColumn("new_column", F.explode(F.split(F.col("new_column"), ","))) # 验证结果 df_result.show() # 如果需要转成RDD,调用.rdd方法即可 result_rdd = df_result.rdd
方案2:基于原生RDD实现
如果你的原始数据是RDD结构,使用flatMap算子实现一对多的行转换即可:
def process_row(row): # 拆分目标字段得到所有取值 split_values = row["new_column"].split(",") # 每个取值生成一条新数据 for val in split_values: new_row = row.copy() new_row["new_column"] = val yield new_row # flatMap会把每行生成的多个结果压平为一维RDD result_rdd = original_rdd.flatMap(process_row)
两种方案执行后得到的结果都和你要求的结构完全一致。
内容的提问来源于stack exchange,提问作者romi
相关产品推荐
相关产品推荐

