-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathdata_generator.py
More file actions
150 lines (124 loc) · 7.01 KB
/
Copy pathdata_generator.py
File metadata and controls
150 lines (124 loc) · 7.01 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
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
import argparse
import random
import time
from data_generator.ping_percentagecpu import ping_sensor, percentagecpu_sensor
from data_generator.rand_data import data_generator as rand_data
from data_generator.blob_people_video import get_data as people_counter
from data_generator.blobs_factory_images import get_data as image_processing
def __check_data_generators(data_generators:str):
for data_gen in data_generators.split(","):
if data_gen not in ['rand', 'ping', 'percentagecpu', 'cars', 'people', 'images']:
raise argparse.ArgumentError(f"Invalid data type {data_gen}")
return data_generators
def __extract_conn(conn_info:str)->(str, tuple):
conns = {}
for conn in conn_info.split(","):
auth = ()
if '@' in conn:
auth, conn = conn.split('@')
auth = tuple(auth.split(':'))
conns[conn] = auth
return conns
def __generate_examples():
from data_generator.support import serialize_data
output = "Sample Values:"
output += f"\n\t{serialize_data(rand_data(db_name='test'))}"
output += f"\n\t{serialize_data(ping_sensor(db_name='test'))}"
output += f"\n\t{serialize_data(percentagecpu_sensor(db_name='test'))}"
output += "\n\nSample Calls"
output += "\n\tSending data to MQTT: python3 ~/Sample-Data-Generator/data_generator.py rand anyloguser:mqtt4AnyLog!@localhost:1883 mqtt --topic test --exception"
output += "\n\tSending data to Kafka: python3 ~/Sample-Data-Generator/data_generator.py rand 35.188.2.231:9092 kafka --topic test --exception"
output += "\n\tSending data via REST POST: python3 ~/Sample-Data-Generator/data_generator.py rand 127.0.0.1:32149 post --topic test --exception"
output += "\n\tSending data via REST PUT: python3 ~/Sample-Data-Generator/data_generator.py rand 127.0.0.1:32149 put --exception"
print(output)
def __generate_data(data_generator:str, db_name:str, last_blob:str=None, exception:bool=False):
payload = {}
if data_generator == 'ping':
payload = ping_sensor(db_name=db_name)
elif data_generator == 'percentagecpu':
payload = percentagecpu_sensor(db_name=db_name)
elif data_generator == 'rand':
payload = rand_data(db_name=db_name)
elif data_generator == 'cars':
from data_generator.blobs_car_video import car_counting
payload, last_blob = car_counting(db_name=db_name, last_blob=last_blob, exception=exception)
elif data_generator == 'people':
payload, last_blob = people_counter(db_name=db_name, last_blob=last_blob, exception=exception)
elif data_generator == 'images':
payload, last_blob = image_processing(db_name=db_name, last_blob=last_blob, exception=exception)
return payload, last_blob
def __publish_data(publisher:str, conn:str, payload:list, topic:str, qos:int=0, auth:tuple=(), timeout:float=30,
exception:bool=False):
if publisher == 'put':
from data_publisher.publisher_rest import publish_via_put
publish_via_put(conn=conn, payload=payload, auth=auth, timeout=timeout, exception=exception)
elif publisher == 'post':
from data_publisher.publisher_rest import publish_via_post
publish_via_post(conn=conn, payload=payload, topic=topic, auth=auth, timeout=timeout, exception=exception)
elif publisher == 'mqtt':
from data_publisher.publisher_mqtt import publish_mqtt
publish_mqtt(conn=conn, payload=payload, topic=topic, qos=qos, auth=auth, exception=exception)
elif publisher == 'kafka':
from data_publisher.publisher_kafka import publish_kafka
publish_kafka(conn=conn, payload=payload, topic=topic, auth=auth, exception=exception)
def main():
parser = argparse.ArgumentParser()
parser.add_argument('data_generator', type=__check_data_generators, default='rand',
help='data to generate')
parser.add_argument('conn', type=str, default='127.0.0.1:32149',
help='connection information (example: [user]:[passwd]@[ip]:[port])')
parser.add_argument('publisher', type=str, default='put',
choices=['put', 'post', 'mqtt', 'kafka'], help='format to publish data')
parser.add_argument('--batch-size', type=int, default=10, help='number of rows per insert batch')
parser.add_argument('--total-rows', type=int, default=10, help='total rows to insert - if set to 0 then run continuously')
parser.add_argument('--sleep', type=float, default=0.5, help='wait time between each row to insert')
parser.add_argument('--db-name', type=str, default='test', help='logical database name')
parser.add_argument('--topic', type=str, default='anylog-demo', help='topic name for POST, MQTT and Kafka')
parser.add_argument('--timeout', type=float, default=30, help='REST timeout')
parser.add_argument('--qos', type=int, choices=list(range(0, 4)), default=0, help='Quality of Service')
parser.add_argument('--exception', type=bool, nargs='?', const=True, default=False,
help='Whether to print exceptions')
parser.add_argument('--examples', type=str, nargs='?', const=True, default=False, help='print example calls and sample data')
args = parser.parse_args()
conns = __extract_conn(conn_info=args.conn)
data_generators = list(args.data_generator.split(","))
status = True
error_msg = ""
if len(data_generators) > 1:
for x in data_generators:
if x in ['cars', 'people', 'images']:
if f"blobs not supported with other data generators" not in error_msg:
error_msg += "Blobs not supported with other data generators\n"
status = False
if args.publisher != 'put':
error_msg += f"Multiple data generator types require put publishing type"
status = False
if status is False:
print(error_msg)
exit(1)
total_rows = 0
payloads = []
if args.data_generator in ['cars', 'people', 'images'] and args.publisher == 'put':
print(f"Data generator for blobs ({args.data_generator} cannot use PUT as a processing option")
exit(1)
if args.examples:
__generate_examples()
exit(1)
last_blob = None
while True:
conn = random.choice(list(conns.keys()))
data_generator = random.choice(data_generators)
auth = conns[conn]
payload, last_blob = __generate_data(data_generator=args.data_generator, db_name=args.db_name,
last_blob=last_blob, exception=args.exception)
payloads.append(payload)
if len(payloads) == args.batch_size or (args.total_rows <= len(payloads) + total_rows and args.total_rows != 0):
__publish_data(publisher=args.publisher, conn=conn, payload=payloads, topic=args.topic, qos=args.qos,
auth=auth, timeout=args.timeout, exception=args.exception)
total_rows += len(payloads)
payloads = []
if total_rows >= args.total_rows:
exit(1)
time.sleep(args.sleep)
if __name__ == '__main__':
main()