forked from oap-project/recdp
-
Notifications
You must be signed in to change notification settings - Fork 0
/
test_sortArrayByFrequency.py
51 lines (44 loc) · 1.78 KB
/
test_sortArrayByFrequency.py
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
#!/env/bin/python
import os
import pathlib
import sys
import numpy as np
import pandas as pd
import pyrecdp
import pyspark.sql.functions as f
import pyspark.sql.types as t
from pyrecdp.data_processor import *
from pyrecdp.utils import *
from pyspark import *
from pyspark.sql import *
def main():
path_prefix = "file://"
cur_folder = str(pathlib.Path(__file__).parent.absolute())
folder = cur_folder + "/data"
path = path_prefix + folder
recdp_path = pyrecdp.__path__[0]
scala_udf_jars = recdp_path + "/ScalaProcessUtils/target/recdp-scala-extensions-0.1.0-jar-with-dependencies.jar"
print(scala_udf_jars)
##### 1. Start spark and initialize data processor #####
spark = SparkSession.builder.master("local[1]")\
.config('spark.eventLog.enabled', False)\
.config('spark.driver.maxResultSize', '16G')\
.config('spark.driver.memory', '10g')\
.config('spark.worker.memory', '10g')\
.config('spark.executor.memory', '10g')\
.config("spark.driver.extraClassPath", f"{scala_udf_jars}")\
.config("spark.executor.extraClassPath", f"{scala_udf_jars}")\
.appName("test_sortArrayByFrequency")\
.getOrCreate()
proc = DataProcessor(spark)
print(f"DataSource path is {path}")
df = spark.read.parquet(f"{path}")
df = df.select("language", "tweet_timestamp")
df = df.withColumn("dt_hour", f.dayofweek(f.from_unixtime(f.col('tweet_timestamp'))).cast(t.IntegerType()))
df = df.groupby('dt_hour').agg(f.collect_list("language").alias("language_list"))
df = df.filter("size(language_list) > 3")
df = df.withColumn("sorted_langugage", f.expr(f"sortStringArrayByFrequency(language_list)"))
df.printSchema()
df.show(24, vertical = True, truncate = 100)
if __name__ == "__main__":
main()