| tierno | ae50192 | 2018-02-06 23:17:16 +0100 | [diff] [blame] | 1 | import logging |
| 2 | import asyncio |
| 3 | import yaml |
| gcalvino | 13fd64d | 2018-02-02 10:53:30 +0100 | [diff] [blame] | 4 | from aiokafka import AIOKafkaConsumer |
| 5 | from aiokafka import AIOKafkaProducer |
| 6 | from aiokafka.errors import KafkaError |
| 7 | from msgbase import MsgBase, MsgException |
| gcalvino | 13fd64d | 2018-02-02 10:53:30 +0100 | [diff] [blame] | 8 | #import json |
| 9 | |
| tierno | ae50192 | 2018-02-06 23:17:16 +0100 | [diff] [blame] | 10 | |
| tierno | f3c4dbc | 2018-02-05 14:53:28 +0100 | [diff] [blame] | 11 | class MsgKafka(MsgBase): |
| tierno | ae50192 | 2018-02-06 23:17:16 +0100 | [diff] [blame] | 12 | def __init__(self, logger_name='msg'): |
| 13 | self.logger = logging.getLogger(logger_name) |
| gcalvino | 13fd64d | 2018-02-02 10:53:30 +0100 | [diff] [blame] | 14 | self.host = None |
| 15 | self.port = None |
| 16 | self.consumer = None |
| 17 | self.producer = None |
| 18 | # create a different file for each topic |
| 19 | #self.files = {} |
| 20 | |
| 21 | def connect(self, config): |
| 22 | try: |
| tierno | ae50192 | 2018-02-06 23:17:16 +0100 | [diff] [blame] | 23 | if "logger_name" in config: |
| 24 | self.logger = logging.getLogger(config["logger_name"]) |
| gcalvino | 13fd64d | 2018-02-02 10:53:30 +0100 | [diff] [blame] | 25 | self.host = config["host"] |
| 26 | self.port = config["port"] |
| 27 | self.topic_lst = [] |
| 28 | self.loop = asyncio.get_event_loop() |
| 29 | self.broker = str(self.host) + ":" + str(self.port) |
| 30 | |
| 31 | except Exception as e: # TODO refine |
| 32 | raise MsgException(str(e)) |
| 33 | |
| tierno | f3c4dbc | 2018-02-05 14:53:28 +0100 | [diff] [blame] | 34 | def write(self, topic, key, msg): |
| gcalvino | 13fd64d | 2018-02-02 10:53:30 +0100 | [diff] [blame] | 35 | try: |
| tierno | f3c4dbc | 2018-02-05 14:53:28 +0100 | [diff] [blame] | 36 | self.loop.run_until_complete(self.aiowrite(topic=topic, key=key, msg=yaml.safe_dump(msg, default_flow_style=True))) |
| gcalvino | 13fd64d | 2018-02-02 10:53:30 +0100 | [diff] [blame] | 37 | |
| 38 | except Exception as e: |
| 39 | raise MsgException("Error writing {} topic: {}".format(topic, str(e))) |
| 40 | |
| 41 | def read(self, topic): |
| tierno | f3c4dbc | 2018-02-05 14:53:28 +0100 | [diff] [blame] | 42 | #self.topic_lst.append(topic) |
| gcalvino | 13fd64d | 2018-02-02 10:53:30 +0100 | [diff] [blame] | 43 | try: |
| tierno | f3c4dbc | 2018-02-05 14:53:28 +0100 | [diff] [blame] | 44 | return self.loop.run_until_complete(self.aioread(topic)) |
| gcalvino | 13fd64d | 2018-02-02 10:53:30 +0100 | [diff] [blame] | 45 | except Exception as e: |
| 46 | raise MsgException("Error reading {} topic: {}".format(topic, str(e))) |
| 47 | |
| tierno | f3c4dbc | 2018-02-05 14:53:28 +0100 | [diff] [blame] | 48 | async def aiowrite(self, topic, key, msg, loop=None): |
| gcalvino | 13fd64d | 2018-02-02 10:53:30 +0100 | [diff] [blame] | 49 | try: |
| tierno | f3c4dbc | 2018-02-05 14:53:28 +0100 | [diff] [blame] | 50 | if not loop: |
| 51 | loop = self.loop |
| 52 | self.producer = AIOKafkaProducer(loop=loop, key_serializer=str.encode, value_serializer=str.encode, |
| gcalvino | 13fd64d | 2018-02-02 10:53:30 +0100 | [diff] [blame] | 53 | bootstrap_servers=self.broker) |
| 54 | await self.producer.start() |
| 55 | await self.producer.send(topic=topic, key=key, value=msg) |
| 56 | except Exception as e: |
| 57 | raise MsgException("Error publishing to {} topic: {}".format(topic, str(e))) |
| 58 | finally: |
| 59 | await self.producer.stop() |
| 60 | |
| tierno | f3c4dbc | 2018-02-05 14:53:28 +0100 | [diff] [blame] | 61 | async def aioread(self, topic, loop=None): |
| 62 | if not loop: |
| 63 | loop = self.loop |
| 64 | self.consumer = AIOKafkaConsumer(loop=loop, bootstrap_servers=self.broker) |
| gcalvino | 13fd64d | 2018-02-02 10:53:30 +0100 | [diff] [blame] | 65 | await self.consumer.start() |
| tierno | f3c4dbc | 2018-02-05 14:53:28 +0100 | [diff] [blame] | 66 | self.consumer.subscribe([topic]) |
| gcalvino | 13fd64d | 2018-02-02 10:53:30 +0100 | [diff] [blame] | 67 | try: |
| 68 | async for message in self.consumer: |
| 69 | return yaml.load(message.key), yaml.load(message.value) |
| 70 | except KafkaError as e: |
| 71 | raise MsgException(str(e)) |
| 72 | finally: |
| 73 | await self.consumer.stop() |
| 74 | |
| 75 | |