小编Ank*_*ain的帖子

如何在pyspark数据框中动态添加列

我正在尝试根据输入变量 vIssueCols 添加几列

from pyspark.sql import HiveContext
from pyspark.sql import functions as F
from pyspark.sql.window import Window
vIssueCols=['jobid','locid']
vQuery1 = 'vSrcData2= vSrcData'
vWindow1 = Window.partitionBy("vKey").orderBy("vOrderBy")
for x in vIssueCols:
Query1=vQuery1+'.withColumn("'+x+'_prev",F.lag(vSrcData.'+x+').over(vWindow1))'

exec(vQuery1)
Run Code Online (Sandbox Code Playgroud)

现在上面的查询将生成如下 vQuery1,并且它正在工作,但是

vSrcData2= vSrcData.withColumn("jobid_prev",F.lag(vSrcData.jobid).over(vWindow1)).withColumn("locid_prev",F.lag(vSrcData.locid).over(vWindow1))
Run Code Online (Sandbox Code Playgroud)

我不能写一个类似的查询吗

vSrcData2= vSrcData.withColumn(x+"_prev",F.lag(vSrcData.x).over(vWindow1))for x in vIssueCols
Run Code Online (Sandbox Code Playgroud)

并使用循环语句生成列。一些博客建议添加一个 udf 并调用它,但我将使用上面执行字符串方法来代替使用 udf。

window-functions pyspark

2
推荐指数
1
解决办法
6058
查看次数

标签 统计

pyspark ×1

window-functions ×1