我想使用 pyspark 根據輸入創建新的資料框,它會列印出每個不同值列的第一次出現。rownumber() 作業還是 window()。不確定最好的方法是解決這個問題還是 sparksql 是最好的。基本上,第二個表是我希望輸出的內容,它僅從輸入中列印出第一次出現的值列。我只對“值”列的第一次出現感興趣。如果一個值重復,則只顯示第一個看到的值。
-------- -------- --------
| VALUE| DAY | Color
-------- -------- --------
|20 |MON | BLUE|
|20 |TUES | BLUE|
|30 |WED | BLUE|
-------- -------- --------
-------- -------- --------
| VALUE| DAY | Color
-------- -------- --------
|20 |MON | BLUE|
|30 |WED | BLUE|
-------- -------- --------
uj5u.com熱心網友回復:
這是我在不使用視窗的情況下執行此操作的方法。它可能會在大型資料集上表現更好,因為它可以使用更多的集群來完成作業。在您的情況下,您需要使用“價值”作為部門,使用“工資”作為“日期”。
from pyspark.sql import SparkSession,Row
spark = SparkSession.builder.appName('SparkByExamples.com').getOrCreate()
data = [("James","Sales",3000),("Michael","Sales",4600),
("Robert","Sales",4100),("Maria","Finance",3000),
("Raman","Finance",3000),("Scott","Finance",3300),
("Jen","Finance",3900),("Jeff","Marketing",3000),
("Kumar","Marketing",2000)]
df = spark.createDataFrame(data,["Name","Department","Salary"])
unGroupedDf = df.select( \
df["Department"], \
f.struct(*[\ # Make a struct with all the record elements.
df["Salary"].alias("Salary"),\ #will be sorted on Salary first
df["Department"].alias("Dept"),\
df["Name"].alias("Name")] )\
.alias("record") )
unGroupedDf.groupBy("Department")\ #group
.agg(f.collect_list("record")\ #Gather all the element in a group
.alias("record"))\
.select(\
f.reverse(\ #Make the sort Descending
f.array_sort(\ #Sort the array ascending
f.col("record")\ #the struct
)\
)[0].alias("record"))\ #grab the "Max element in the array
).select( f.col("record.*") ).show() # use struct as Columns
.show()
--------- ------ -------
| Dept|Salary| Name|
--------- ------ -------
| Sales| 4600|Michael|
| Finance| 3900| Jen|
|Marketing| 3000| Jeff|
--------- ------ -------
uj5u.com熱心網友回復:
在我看來,您想洗掉重復的專案VALUE。如果是這樣,請使用dropDuplicates
df.dropDuplicates(['VALUE']).show()
----- --- -----
|VALUE|DAY|Color|
----- --- -----
| 20|MON| BLUE|
| 30|WED| BLUE|
----- --- -----
uj5u.com熱心網友回復:
這是使用視窗的方法。在這個例子中,他們以薪水為例。在你的情況下,我認為你會為 orderBy 使用“DAY”,為 partitionBy 使用“Value”。
from pyspark.sql import SparkSession,Row spark = SparkSession.builder.appName('SparkByExamples.com').getOrCreate() data = [("James","Sales",3000),("Michael","Sales",4600), ("Robert","Sales",4100),("Maria","Finance",3000), ("Raman","Finance",3000),("Scott","Finance",3300), ("Jen","Finance",3900),("Jeff","Marketing",3000), ("Kumar","Marketing",2000)] df = spark.createDataFrame(data,["Name","Department","Salary"]) df.show() from pyspark.sql.window import Window from pyspark.sql.functions import col, row_number w2 = Window.partitionBy("department").orderBy(col("salary")) df.withColumn("row",row_number().over(w2)) \ .filter(col("row") == 1).drop("row") \ .show() ------------- ---------- ------ |employee_name|department|salary| ------------- ---------- ------ | James| Sales| 3000| | Maria| Finance| 3000| | Kumar| Marketing| 2000| ------------- ---------- ------
是的,您需要開發一種排序天數的方法,但我認為您明白這是可能的,并且您選擇了正確的工具。我總是喜歡警告人們,這使用一個視窗,他們將所有資料吸入 1 個執行者來完成作業。這不是特別有效。在小型資料集上,這可能是高性能的。在較大的資料集上,可能需要很長時間才能完成。
轉載請註明出處,本文鏈接:https://www.uj5u.com/caozuo/477866.html
標籤:阿帕奇火花 pyspark apache-spark-sql
上一篇:火花流中的偏移管理
