我正在尝试根据输入变量 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。