-
Notifications
You must be signed in to change notification settings - Fork 65
Expand file tree
/
Copy pathkafka_client.py
More file actions
58 lines (43 loc) · 2.1 KB
/
kafka_client.py
File metadata and controls
58 lines (43 loc) · 2.1 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
import time
from kafka import KafkaProducer
from locust import events
class KafkaClient:
def __init__(self, kafka_brokers=None):
print("creating message sender with params: " + str(locals()))
if kafka_brokers is None:
kafka_brokers = ['localhost:9092']
self.producer = KafkaProducer(bootstrap_servers=kafka_brokers)
def __handle_success(self, *arguments, **kwargs):
end_time = time.time()
elapsed_time = int((end_time - kwargs["start_time"]) * 1000)
try:
record_metadata = kwargs["future"].get(timeout=1)
request_data = dict(request_type="ENQUEUE",
name=record_metadata.topic,
response_time=elapsed_time,
response_length=record_metadata.serialized_value_size)
self.__fire_success(**request_data)
except Exception as ex:
print("Logging the exception : {0}".format(ex))
raise # ??
def __handle_failure(self, *arguments, **kwargs):
print("failure " + str(locals()))
end_time = time.time()
elapsed_time = int((end_time - kwargs["start_time"]) * 1000)
request_data = dict(request_type="ENQUEUE", name=kwargs["topic"], response_time=elapsed_time,
exception=arguments[0])
self.__fire_failure(**request_data)
def __fire_failure(self, **kwargs):
events.request_failure.fire(**kwargs)
def __fire_success(self, **kwargs):
events.request_success.fire(**kwargs)
def send(self, topic, key=None, message=None):
start_time = time.time()
future = self.producer.send(topic, key=key.encode() if key else None,
value=message.encode() if message else None)
future.add_callback(self.__handle_success, start_time=start_time, future=future)
future.add_errback(self.__handle_failure, start_time=start_time, topic=topic)
def finalize(self):
print("flushing the messages")
self.producer.flush(timeout=5)
print("flushing finished")