我想根據該組的上一行中該列的值設置列的值。然後這個更新的值將被用在下一行。如何遍歷pyspark中Dataframe/RDD的每一行。
我有以下數據幀
id | start_date|sort_date | A | B |
-----------------------------------
1 | 1/1/2017 | 31-01-2015 | 1 | 0 |
1 | 1/1/2017 | 28-02-2015 | 0 | 0 |
1 | 1/1/2017 | 31-03-2015 | 1 | 0 |
1 | 1/1/2017 | 30-04-2015 | 1 | 0 |
1 | 1/1/2017 | 31-05-2015 | 1 | 0 |
1 | 1/1/2017 | 30-06-2015 | 1 | 0 |
1 | 1/1/2017 | 31-07-2015 | 1 | 0 |
1 | 1/1/2017 | 31-08-2015 | 1 | 0 |
1 | 1/1/2017 | 30-09-2015 | 0 | 0 |
2 | 1/1/2017 | 31-10-2015 | 1 | 0 |
2 | 1/1/2017 | 30-11-2015 | 0 | 0 |
2 | 1/1/2017 | 31-12-2015 | 1 | 0 |
2 | 1/1/2017 | 31-01-2016 | 1 | 0 |
2 | 1/1/2017 | 28-02-2016 | 1 | 0 |
2 | 1/1/2017 | 31-03-2016 | 1 | 0 |
2 | 1/1/2017 | 30-04-2016 | 1 | 0 |
2 | 1/1/2017 | 31-05-2016 | 1 | 0 |
2 | 1/1/2017 | 30-06-2016 | 0 | 0 |
輸出:
id | start_date|sort_date | A | B | C
---------------------------------------
1 | 1/1/2017 | 31-01-2015 | 1 | 0 | 1
1 | 1/1/2017 | 28-02-2015 | 0 | 0 | 0
1 | 1/1/2017 | 31-03-2015 | 1 | 0 | 1
1 | 1/1/2017 | 30-04-2015 | 1 | 0 | 2
1 | 1/1/2017 | 31-05-2015 | 1 | 0 | 3
1 | 1/1/2017 | 30-06-2015 | 1 | 0 | 4
1 | 1/1/2017 | 31-07-2015 | 1 | 0 | 5
1 | 1/1/2017 | 31-08-2015 | 1 | 0 | 6
1 | 1/1/2017 | 30-09-2015 | 0 | 0 | 0
2 | 1/1/2017 | 31-10-2015 | 1 | 0 | 1
2 | 1/1/2017 | 30-11-2015 | 0 | 0 | 0
2 | 1/1/2017 | 31-12-2015 | 1 | 0 | 1
2 | 1/1/2017 | 31-01-2016 | 1 | 0 | 2
2 | 1/1/2017 | 28-02-2016 | 1 | 0 | 3
2 | 1/1/2017 | 31-03-2016 | 1 | 0 | 4
2 | 1/1/2017 | 30-04-2016 | 1 | 0 | 5
2 | 1/1/2017 | 31-05-2016 | 1 | 0 | 6
2 | 1/1/2017 | 30-06-2016 | 0 | 0 | 0
集團是ID和日期的
列C是衍生基於列A和B.
如果A == 1且B == 0,則C從前一行+ 1導出C.
還有其他一些條件,但我正在努力與這部分。
假設我們在數據框中有一個sort_date列。
我嘗試以下查詢:
SELECT
id,
date,
sort_date,
lag(A) OVER (PARTITION BY id, date ORDER BY sort_date) as prev,
CASE
WHEN A=1 AND B= 0 THEN 1
WHEN A=1 AND B> 0 THEN prev +1
ELSE 0
END AS A
FROM
Table
這是我做的UDAF
val myFunc = new MyUDAF
val w = Window.partitionBy(col("ID"), col("START_DATE")).orderBy(col("SORT_DATE"))
val df = df.withColumn("C", myFunc(col("START_DATE"), col("X"),
col("Y"), col("A"),
col("B")).over(w))
PS:我使用的Spark 1.6
您可以使用** Window函數**與Spark SQL。 – mrsrinivas
你可以添加你試過的代碼嗎? – mrsrinivas
請改善問題:你能否再解釋一下你試圖達到的目標,到目前爲止你做了什麼,你的輸入是什麼,你的期望輸出是什麼,你想在RDD中這樣做,就像標題所說或者在作爲專欄文字的數據框表明了什麼?你是什麼意思一個組?你的意思是一個groupby?你想如何分類? –