-
Notifications
You must be signed in to change notification settings - Fork 0
/
kafka_pub.py
44 lines (37 loc) · 1021 Bytes
/
kafka_pub.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
from kafka import KafkaProducer
import time
import os
import json
import logging
QUOTE='quote'
TRADE='trade'
BAR='bar'
sequences = {BAR: {}, QUOTE: {}}
topics = {BAR: 'stock-bars', QUOTE: 'stock-quotes'}
producer = None
enabled = None
async def init(bootstrap, enable):
global producer, enabled
enabled = enable
if enabled:
producer = KafkaProducer(bootstrap_servers=bootstrap)
def incr(symbol, datatype):
global sequences
if symbol in sequences[datatype]:
sequences[datatype][symbol] += 1
else:
sequences[datatype][symbol] = 1
return sequences[datatype][symbol]
async def publish(symbol, datatype, data):
msg = {'symbol': symbol, datatype: data, 'seq': incr(symbol, datatype)}
if enabled:
producer.send(
topics[datatype],
key=bytearray(symbol, 'utf-8'),
value=bytearray(json.dumps(msg), 'utf-8')
)
else:
print(json.dumps(msg))
async def flush():
if enabled:
producer.flush()