PySpark:使用For循环向DataFrame添加行的实现疑问
解决PySpark中给DataFrame每行添加衍生行的问题
我来帮你搞定这个需求!你想要给原DataFrame的每一行都新增一条符合规则的衍生行,同时保留原数据对吧?咱们一步步来实现:
需求分析
你期望的逻辑很明确:
- 对原DataFrame的每一行,生成一条对应的新行
- 新行的
ColNum是原行的ColNum + 1 - 新行的滞后列依次前移:
ColB_lag2= 原行的ColB_lag1,ColB_lag1= 原行的ColB - 新行的
ColB由自定义函数someFunc()生成 - 最终结果要保留原行和衍生行,保持每组数据的顺序
解决方案
关键是不要直接修改原DataFrame,而是先复制原数据生成衍生行,再把原数据和衍生行合并:
1. 导入必要模块(如果还没导的话)
from pyspark.sql.types import IntegerType from pyspark.sql.functions import col
2. 生成衍生行DataFrame
基于原DataFrame创建衍生行的数据集,应用你需要的字段变换:
# 基于原df生成衍生行,这一步不会修改原数据 derived_df = df.withColumn("ColNum", (col("ColNum") + 1).cast(IntegerType())) \ .withColumn("ColB_lag2", col("ColB_lag1")) \ .withColumn("ColB_lag1", col("ColB")) \ .withColumn("ColB", someFunc()) # 替换成你实际的ColB生成函数
3. 合并原数据和衍生行
用union方法把原DataFrame和衍生行DataFrame合并:
# 合并两个数据集 result_df = df.union(derived_df)
4. 排序保证顺序
最后按ColA分组、ColNum升序排序,就能得到你想要的规整结果:
# 排序后输出 result_df = result_df.orderBy("ColA", "ColNum")
验证结果
执行上述代码后,result_df的输出就会和你期望的完全一致:
ColA ColNum ColB ColB_lag1 ColB_lag2
Xyz 25 123 234 345
Xyz 26 789 123 234
Abc 40 456 567 678
Abc 41 890 456 567
小提示
- 要确保
someFunc()返回的类型和原ColB的类型一致,避免出现类型不匹配的错误 - 如果你的DataFrame还有其他字段,要保证衍生行的字段和原DataFrame完全一致,否则
union操作会报错
内容的提问来源于stack exchange,提问作者Apoorv Agarwal
相关产品推荐
相关产品推荐

