-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathspark_streaming_transactions_app.py
More file actions
94 lines (68 loc) · 3.09 KB
/
Copy pathspark_streaming_transactions_app.py
File metadata and controls
94 lines (68 loc) · 3.09 KB
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
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
from pyspark.streaming import StreamingContext
from pyspark.sql import SparkSession
from pyspark.sql import Row
import numpy as np
from plotly.graph_objs import *
from plot.transactions_map import maps_stream
from plot.transactions_pie import transactions_pie_stream
from plot.dashboard import upload_dashboard
from plot.color.color import convert_to_color
spark = SparkSession.builder.master("local[2]").appName("ScoringApp").getOrCreate()
sc = spark.sparkContext
sc.setLogLevel("ERROR")
batchIntervalSeconds = 5
hostname = 'localhost'
ip = 5900
columns = ['step', 'type', 'amount', 'nameOrig', 'oldbalanceOrg', 'newbalanceOrig', 'nameDest', 'oldbalanceDest',
'newbalanceDest', 'isFraud', 'isFlaggedFraud', 'gps_latitude', 'gps_longitude', "location", "id",
"entity_id"]
float_columns = ['step', 'amount', 'oldbalanceOrg', 'newbalanceOrig', 'oldbalanceDest', 'newbalanceDest',
'gps_latitude', 'gps_longitude']
transactions_pie_stream.open()
maps_stream.open()
upload_dashboard()
def parse_row(row):
new_row = []
for column in columns:
if column in float_columns:
new_row.append((column, float(row[columns.index(column)])))
else:
new_row.append((column, row[columns.index(column)]))
return Row(**dict(new_row))
def publish_transactions_to_map(rdd):
color = {}
for transaction in rdd.collect():
if transaction["entity_id"] not in color:
color[transaction["entity_id"]] = convert_to_color(transaction["entity_id"])
lat, lon = transaction["gps_latitude"], transaction["gps_longitude"]
maps_stream.write(
dict(lat=lat, lon=lon, type="Scattermapbox",
marker=Marker(size=15, color=color[transaction["entity_id"]]),
text="{}\n\tAmount: {}\n\tType: {}".format(transaction["id"] + "/" + transaction["entity_id"],
transaction["amount"],
transaction["type"])))
def publish_transactions_to_pie(rdd):
if not rdd.isEmpty():
array_data = np.transpose(np.array(rdd.collect())).tolist()
transactions_pie_stream.write(dict(labels=array_data[0],
values=array_data[1],
type='pie'))
def parse_rdd(rdd):
if not rdd.isEmpty():
return rdd.map(lambda line: line.split(",")).map(parse_row)
else:
return sc.parallelize([])
def create_stream(batch_interval):
ssc = StreamingContext(sc, batch_interval)
dstream_input = ssc.socketTextStream(hostname, ip)
dstream_parsed = dstream_input.transform(parse_rdd)
dstream_parsed.foreachRDD(publish_transactions_to_map)
dstream_total_by_location = dstream_parsed.map(lambda row: (row["location"], row["amount"])).reduceByKey(
lambda a, b: a + b)
# dstream_total_by_location.pprint()
dstream_total_by_location.foreachRDD(publish_transactions_to_pie)
return ssc
ssc = create_stream(batchIntervalSeconds)
ssc.start()
ssc.awaitTermination()
# maps_stream.close()