0
votes

I have a data frame which is ~49 Billion records. The data frame looks like -

id       transaction_no         amount
1            321                  100
1            100                   50
1            32                   200
2            54                    50
2            20                  1000
3            41                    44
4            78                   400
4            65                   200

My final output looks like -

id        count         amount
1           3            350
2           2           1050
3           1             44
4           2            600

I can do it in python, but how to do it in pyspark?

1

1 Answers

1
votes
import pyspark.sql.functions as f

df=spark.createDataFrame([(1,321,100),(1,100,50),(1,32,200),(2,54,50),(2,20,1000),(3,41,44),(4,78,400),(4,65,200)],['id','transaction_no','amount'])

df_req = df.groupby('id').agg(f.sum('amount').alias('amount'),f.count('id').alias('count'))
df_req.show()

+---+------+-----+
| id|amount|count|
+---+------+-----+
|  1|   350|    3|
|  3|    44|    1|
|  2|  1050|    2|
|  4|   600|    2|
+---+------+-----+