| tierno | 5c01261 | 2018-04-19 16:01:59 +0200 | [diff] [blame] | 1 | import logging |
| 2 | import asyncio |
| 3 | import yaml |
| 4 | from aiokafka import AIOKafkaConsumer |
| 5 | from aiokafka import AIOKafkaProducer |
| 6 | from aiokafka.errors import KafkaError |
| tierno | 3054f78 | 2018-04-25 16:59:53 +0200 | [diff] [blame] | 7 | from osm_common.msgbase import MsgBase, MsgException |
| 8 | # import json |
| tierno | 5c01261 | 2018-04-19 16:01:59 +0200 | [diff] [blame] | 9 | |
| 10 | __author__ = "Alfonso Tierno <alfonso.tiernosepulveda@telefonica.com>, " \ |
| 11 | "Guillermo Calvino <guillermo.calvinosanchez@altran.com>" |
| tierno | 3054f78 | 2018-04-25 16:59:53 +0200 | [diff] [blame] | 12 | |
| 13 | |
| tierno | 5c01261 | 2018-04-19 16:01:59 +0200 | [diff] [blame] | 14 | class MsgKafka(MsgBase): |
| 15 | def __init__(self, logger_name='msg'): |
| 16 | self.logger = logging.getLogger(logger_name) |
| 17 | self.host = None |
| 18 | self.port = None |
| 19 | self.consumer = None |
| 20 | self.producer = None |
| 21 | self.loop = None |
| 22 | self.broker = None |
| tierno | 73da4fa | 2018-08-31 13:50:59 +0000 | [diff] [blame^] | 23 | self.group_id = None |
| tierno | 5c01261 | 2018-04-19 16:01:59 +0200 | [diff] [blame] | 24 | |
| 25 | def connect(self, config): |
| 26 | try: |
| 27 | if "logger_name" in config: |
| 28 | self.logger = logging.getLogger(config["logger_name"]) |
| 29 | self.host = config["host"] |
| 30 | self.port = config["port"] |
| 31 | self.loop = asyncio.get_event_loop() |
| 32 | self.broker = str(self.host) + ":" + str(self.port) |
| tierno | 73da4fa | 2018-08-31 13:50:59 +0000 | [diff] [blame^] | 33 | self.group_id = config.get("group_id") |
| tierno | 5c01261 | 2018-04-19 16:01:59 +0200 | [diff] [blame] | 34 | |
| 35 | except Exception as e: # TODO refine |
| 36 | raise MsgException(str(e)) |
| 37 | |
| 38 | def disconnect(self): |
| 39 | try: |
| tierno | ebbf353 | 2018-05-03 17:49:37 +0200 | [diff] [blame] | 40 | pass |
| 41 | # self.loop.close() |
| tierno | 5c01261 | 2018-04-19 16:01:59 +0200 | [diff] [blame] | 42 | except Exception as e: # TODO refine |
| 43 | raise MsgException(str(e)) |
| 44 | |
| 45 | def write(self, topic, key, msg): |
| tierno | 8657799 | 2018-05-10 16:51:17 +0200 | [diff] [blame] | 46 | """ |
| 47 | Write a message at kafka bus |
| 48 | :param topic: message topic, must be string |
| 49 | :param key: message key, must be string |
| 50 | :param msg: message content, can be string or dictionary |
| 51 | :return: None or raises MsgException on failing |
| 52 | """ |
| tierno | 5c01261 | 2018-04-19 16:01:59 +0200 | [diff] [blame] | 53 | try: |
| tierno | 8657799 | 2018-05-10 16:51:17 +0200 | [diff] [blame] | 54 | self.loop.run_until_complete(self.aiowrite(topic=topic, key=key, msg=msg)) |
| tierno | 5c01261 | 2018-04-19 16:01:59 +0200 | [diff] [blame] | 55 | |
| 56 | except Exception as e: |
| 57 | raise MsgException("Error writing {} topic: {}".format(topic, str(e))) |
| 58 | |
| 59 | def read(self, topic): |
| 60 | """ |
| tierno | 8657799 | 2018-05-10 16:51:17 +0200 | [diff] [blame] | 61 | Read from one or several topics. |
| tierno | 5c01261 | 2018-04-19 16:01:59 +0200 | [diff] [blame] | 62 | :param topic: can be str: single topic; or str list: several topics |
| 63 | :return: topic, key, message; or None |
| 64 | """ |
| 65 | try: |
| 66 | return self.loop.run_until_complete(self.aioread(topic, self.loop)) |
| 67 | except MsgException: |
| 68 | raise |
| 69 | except Exception as e: |
| 70 | raise MsgException("Error reading {} topic: {}".format(topic, str(e))) |
| 71 | |
| 72 | async def aiowrite(self, topic, key, msg, loop=None): |
| 73 | |
| 74 | if not loop: |
| 75 | loop = self.loop |
| 76 | try: |
| 77 | self.producer = AIOKafkaProducer(loop=loop, key_serializer=str.encode, value_serializer=str.encode, |
| 78 | bootstrap_servers=self.broker) |
| 79 | await self.producer.start() |
| tierno | 3054f78 | 2018-04-25 16:59:53 +0200 | [diff] [blame] | 80 | await self.producer.send(topic=topic, key=key, value=yaml.safe_dump(msg, default_flow_style=True)) |
| tierno | 5c01261 | 2018-04-19 16:01:59 +0200 | [diff] [blame] | 81 | except Exception as e: |
| tierno | 3054f78 | 2018-04-25 16:59:53 +0200 | [diff] [blame] | 82 | raise MsgException("Error publishing topic '{}', key '{}': {}".format(topic, key, e)) |
| tierno | 5c01261 | 2018-04-19 16:01:59 +0200 | [diff] [blame] | 83 | finally: |
| 84 | await self.producer.stop() |
| 85 | |
| 86 | async def aioread(self, topic, loop=None, callback=None, *args): |
| 87 | """ |
| 88 | Asyncio read from one or several topics. It blocks |
| 89 | :param topic: can be str: single topic; or str list: several topics |
| 90 | :param loop: asyncio loop |
| 91 | :callback: callback function that will handle the message in kafka bus |
| 92 | :*args: optional arguments for callback function |
| 93 | :return: topic, key, message |
| 94 | """ |
| 95 | |
| 96 | if not loop: |
| 97 | loop = self.loop |
| 98 | try: |
| 99 | if isinstance(topic, (list, tuple)): |
| 100 | topic_list = topic |
| 101 | else: |
| 102 | topic_list = (topic,) |
| 103 | |
| tierno | 73da4fa | 2018-08-31 13:50:59 +0000 | [diff] [blame^] | 104 | self.consumer = AIOKafkaConsumer(loop=loop, bootstrap_servers=self.broker, group_id=self.group_id) |
| tierno | 5c01261 | 2018-04-19 16:01:59 +0200 | [diff] [blame] | 105 | await self.consumer.start() |
| 106 | self.consumer.subscribe(topic_list) |
| 107 | |
| 108 | async for message in self.consumer: |
| 109 | if callback: |
| 110 | callback(message.topic, yaml.load(message.key), yaml.load(message.value), *args) |
| 111 | else: |
| 112 | return message.topic, yaml.load(message.key), yaml.load(message.value) |
| 113 | except KafkaError as e: |
| 114 | raise MsgException(str(e)) |
| 115 | finally: |
| 116 | await self.consumer.stop() |