-
Notifications
You must be signed in to change notification settings - Fork 0
/
Copy pathreceive.py
executable file
·90 lines (71 loc) · 2.76 KB
/
receive.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
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
#!/usr/bin/python
import time, pika, sys, os, MySQLdb
credentials = pika.PlainCredentials(os.environ['MQ_USER'], os.environ['MQ_PASS'])
parameters = pika.ConnectionParameters(os.environ['MQ_SERVER'],int(os.environ['MQ_PORT']),'/',credentials)
######################
## Connect to rabbitmq
connection = pika.BlockingConnection(parameters)
channel = connection.channel()
## declare the que (idempotent)
channel.queue_declare(queue=os.environ['MQ_QUEUENAME'])
################
## Open the file
#file = open("receive_output", "a")
###################
## Connect to MySQL
db = MySQLdb.connect(host=os.environ['SQL_SERVER'], # your host, usually localhost
user=os.environ['SQL_USER'], # your username
port=int(os.environ['SQL_PORT']), # port
passwd=os.environ['SQL_PASS'], # your password
db=os.environ['SQL_DB']) # name of the data base
# you must create a Cursor object. It will let
# you execute all the queries you need
cursor = db.cursor()
## callback function for basic_cosume
def consumer_callback(ch, method, properties, body):
## print rabbitmq message on screen
print("%s" % body)
## write rabbitmq message to file
file = open("/var/log/receive_output", "a")
file.write("%s" % body + "\n")
file.close()
## split message into values
fields = body.split(",")
while len(fields) < 11:
fields.append('')
datetimestripped=fields[1].replace(" ", "").replace("-", "").replace(":", "").split(".")[0]
messagetype=fields[2].replace("$", "")
cursor.execute('INSERT INTO nmea_raw_data(datetime, messagetype, f1, f2, f3, f4, f5, f6, f7, f8 ) \
VALUES(%s, %s, %s, %s, %s, %s, %s, %s, %s, %s)',
(datetimestripped, messagetype, fields[3], fields[4], fields[5], fields[6], fields[7], fields[8], fields[9], fields[10]))
# print(type(properties))
# print(properties)
# if isinstance(properties, str) is True:
# if "noack" not in properties:
# db.commit()
# print("commit")
# else:
# print("noack is set")
db.commit()
## argument is given, process as log file
if len(sys.argv) > 1:
file = open(sys.argv[1], "r")
for line in file:
consumer_callback("","","noack",line)
db.commit()
file.close()
else:
# receive from queue
channel.basic_consume(consumer_callback,
queue=os.environ['MQ_QUEUENAME'],
no_ack=True,
exclusive=False, consumer_tag=None, arguments=None)
db.commit()
print(' [*] Waiting for messages. To exit press CTRL+C')
channel.start_consuming()
# close connection
connection.close()
#file.close()
#close the connection to the database.
cursor.close()
db.close()