forked from snowplow/snowplow-python-tracker
-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathredis_worker.py
More file actions
74 lines (59 loc) · 1.79 KB
/
redis_worker.py
File metadata and controls
74 lines (59 loc) · 1.79 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
import sys
from snowplow_tracker import Emitter
from typing import Any
from snowplow_tracker.typing import PayloadDict
import json
import redis
import signal
import gevent
from gevent.pool import Pool
def get_url_from_args():
if len(sys.argv) != 2:
raise ValueError("Collector Endpoint is required")
return sys.argv[1]
class RedisWorker:
def __init__(self, emitter: Emitter, key) -> None:
self.pool = Pool(5)
self.emitter = emitter
self.rdb = redis.StrictRedis()
self.key = key
signal.signal(signal.SIGTERM, self.request_shutdown)
signal.signal(signal.SIGINT, self.request_shutdown)
signal.signal(signal.SIGQUIT, self.request_shutdown)
def send(self, payload: PayloadDict) -> None:
"""
Send an event to an emitter
"""
self.emitter.input(payload)
def pop_payload(self) -> None:
"""
Get a single event from Redis and send it
If the Redis queue is empty, sleep to avoid making continual requests
"""
payload = self.rdb.lpop(self.key)
if payload:
self.pool.spawn(self.send, json.loads(payload.decode("utf-8")))
else:
gevent.sleep(5)
def run(self) -> None:
"""
Run indefinitely
"""
self._shutdown = False
while not self._shutdown:
self.pop_payload()
self.pool.join(timeout=20)
def request_shutdown(self, *args: Any) -> None:
"""
Halt the worker
"""
self._shutdown = True
def main():
collector_url = get_url_from_args()
# Configure Emitter
emitter = Emitter(collector_url, batch_size=1)
# Setup worker
worker = RedisWorker(emitter=emitter, key="redis_key")
worker.run()
if __name__ == "__main__":
main()