| tierno | 87858ca | 2018-10-08 16:30:15 +0200 | [diff] [blame] | 1 | # -*- coding: utf-8 -*- |
| 2 | |
| 3 | # Licensed under the Apache License, Version 2.0 (the "License"); |
| 4 | # you may not use this file except in compliance with the License. |
| 5 | # You may obtain a copy of the License at |
| 6 | # |
| 7 | # http://www.apache.org/licenses/LICENSE-2.0 |
| 8 | # |
| 9 | # Unless required by applicable law or agreed to in writing, software |
| 10 | # distributed under the License is distributed on an "AS IS" BASIS, |
| 11 | # WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or |
| 12 | # implied. |
| 13 | # See the License for the specific language governing permissions and |
| 14 | # limitations under the License. |
| 15 | |
| tierno | 5c01261 | 2018-04-19 16:01:59 +0200 | [diff] [blame] | 16 | import asyncio |
| aticig | 3dd0db6 | 2022-03-04 19:35:45 +0300 | [diff] [blame] | 17 | import logging |
| 18 | |
| tierno | 5c01261 | 2018-04-19 16:01:59 +0200 | [diff] [blame] | 19 | from aiokafka import AIOKafkaConsumer |
| 20 | from aiokafka import AIOKafkaProducer |
| 21 | from aiokafka.errors import KafkaError |
| tierno | 3054f78 | 2018-04-25 16:59:53 +0200 | [diff] [blame] | 22 | from osm_common.msgbase import MsgBase, MsgException |
| aticig | 3dd0db6 | 2022-03-04 19:35:45 +0300 | [diff] [blame] | 23 | import yaml |
| tierno | 5c01261 | 2018-04-19 16:01:59 +0200 | [diff] [blame] | 24 | |
| garciadeblas | 2644b76 | 2021-03-24 09:21:01 +0100 | [diff] [blame] | 25 | __author__ = ( |
| 26 | "Alfonso Tierno <alfonso.tiernosepulveda@telefonica.com>, " |
| 27 | "Guillermo Calvino <guillermo.calvinosanchez@altran.com>" |
| 28 | ) |
| tierno | 3054f78 | 2018-04-25 16:59:53 +0200 | [diff] [blame] | 29 | |
| 30 | |
| tierno | 5c01261 | 2018-04-19 16:01:59 +0200 | [diff] [blame] | 31 | class MsgKafka(MsgBase): |
| garciadeblas | 2644b76 | 2021-03-24 09:21:01 +0100 | [diff] [blame] | 32 | def __init__(self, logger_name="msg", lock=False): |
| tierno | 1e9a329 | 2018-11-05 18:18:45 +0100 | [diff] [blame] | 33 | super().__init__(logger_name, lock) |
| tierno | 5c01261 | 2018-04-19 16:01:59 +0200 | [diff] [blame] | 34 | self.host = None |
| 35 | self.port = None |
| 36 | self.consumer = None |
| 37 | self.producer = None |
| tierno | 5c01261 | 2018-04-19 16:01:59 +0200 | [diff] [blame] | 38 | self.broker = None |
| tierno | 73da4fa | 2018-08-31 13:50:59 +0000 | [diff] [blame] | 39 | self.group_id = None |
| tierno | 5c01261 | 2018-04-19 16:01:59 +0200 | [diff] [blame] | 40 | |
| 41 | def connect(self, config): |
| 42 | try: |
| 43 | if "logger_name" in config: |
| 44 | self.logger = logging.getLogger(config["logger_name"]) |
| 45 | self.host = config["host"] |
| 46 | self.port = config["port"] |
| tierno | 5c01261 | 2018-04-19 16:01:59 +0200 | [diff] [blame] | 47 | self.broker = str(self.host) + ":" + str(self.port) |
| tierno | 73da4fa | 2018-08-31 13:50:59 +0000 | [diff] [blame] | 48 | self.group_id = config.get("group_id") |
| tierno | 5c01261 | 2018-04-19 16:01:59 +0200 | [diff] [blame] | 49 | |
| 50 | except Exception as e: # TODO refine |
| 51 | raise MsgException(str(e)) |
| 52 | |
| 53 | def disconnect(self): |
| 54 | try: |
| tierno | ebbf353 | 2018-05-03 17:49:37 +0200 | [diff] [blame] | 55 | pass |
| tierno | 5c01261 | 2018-04-19 16:01:59 +0200 | [diff] [blame] | 56 | except Exception as e: # TODO refine |
| 57 | raise MsgException(str(e)) |
| 58 | |
| 59 | def write(self, topic, key, msg): |
| tierno | 8657799 | 2018-05-10 16:51:17 +0200 | [diff] [blame] | 60 | """ |
| 61 | Write a message at kafka bus |
| 62 | :param topic: message topic, must be string |
| 63 | :param key: message key, must be string |
| 64 | :param msg: message content, can be string or dictionary |
| 65 | :return: None or raises MsgException on failing |
| 66 | """ |
| garciadeblas | 2644b76 | 2021-03-24 09:21:01 +0100 | [diff] [blame] | 67 | retry = 2 # Try two times |
| delacruzramo | 562435a | 2019-12-10 12:06:01 +0100 | [diff] [blame] | 68 | while retry: |
| 69 | try: |
| Gulsum Atici | a06b854 | 2023-05-09 13:42:13 +0300 | [diff] [blame^] | 70 | asyncio.run(self.aiowrite(topic=topic, key=key, msg=msg)) |
| delacruzramo | 562435a | 2019-12-10 12:06:01 +0100 | [diff] [blame] | 71 | break |
| 72 | except Exception as e: |
| 73 | retry -= 1 |
| 74 | if retry == 0: |
| garciadeblas | 2644b76 | 2021-03-24 09:21:01 +0100 | [diff] [blame] | 75 | raise MsgException( |
| 76 | "Error writing {} topic: {}".format(topic, str(e)) |
| 77 | ) |
| tierno | 5c01261 | 2018-04-19 16:01:59 +0200 | [diff] [blame] | 78 | |
| 79 | def read(self, topic): |
| 80 | """ |
| tierno | 8657799 | 2018-05-10 16:51:17 +0200 | [diff] [blame] | 81 | Read from one or several topics. |
| tierno | 5c01261 | 2018-04-19 16:01:59 +0200 | [diff] [blame] | 82 | :param topic: can be str: single topic; or str list: several topics |
| 83 | :return: topic, key, message; or None |
| 84 | """ |
| 85 | try: |
| Gulsum Atici | a06b854 | 2023-05-09 13:42:13 +0300 | [diff] [blame^] | 86 | return asyncio.run(self.aioread(topic)) |
| tierno | 5c01261 | 2018-04-19 16:01:59 +0200 | [diff] [blame] | 87 | except MsgException: |
| 88 | raise |
| 89 | except Exception as e: |
| 90 | raise MsgException("Error reading {} topic: {}".format(topic, str(e))) |
| 91 | |
| Gulsum Atici | a06b854 | 2023-05-09 13:42:13 +0300 | [diff] [blame^] | 92 | async def aiowrite(self, topic, key, msg): |
| tierno | 05ede8f | 2019-01-28 16:20:18 +0000 | [diff] [blame] | 93 | """ |
| 94 | Asyncio write |
| 95 | :param topic: str kafka topic |
| 96 | :param key: str kafka key |
| 97 | :param msg: str or dictionary kafka message |
| tierno | 05ede8f | 2019-01-28 16:20:18 +0000 | [diff] [blame] | 98 | :return: None |
| 99 | """ |
| tierno | 5c01261 | 2018-04-19 16:01:59 +0200 | [diff] [blame] | 100 | try: |
| garciadeblas | 2644b76 | 2021-03-24 09:21:01 +0100 | [diff] [blame] | 101 | self.producer = AIOKafkaProducer( |
| garciadeblas | 2644b76 | 2021-03-24 09:21:01 +0100 | [diff] [blame] | 102 | key_serializer=str.encode, |
| 103 | value_serializer=str.encode, |
| 104 | bootstrap_servers=self.broker, |
| 105 | ) |
| tierno | 5c01261 | 2018-04-19 16:01:59 +0200 | [diff] [blame] | 106 | await self.producer.start() |
| garciadeblas | 2644b76 | 2021-03-24 09:21:01 +0100 | [diff] [blame] | 107 | await self.producer.send( |
| 108 | topic=topic, key=key, value=yaml.safe_dump(msg, default_flow_style=True) |
| 109 | ) |
| tierno | 5c01261 | 2018-04-19 16:01:59 +0200 | [diff] [blame] | 110 | except Exception as e: |
| garciadeblas | 2644b76 | 2021-03-24 09:21:01 +0100 | [diff] [blame] | 111 | raise MsgException( |
| 112 | "Error publishing topic '{}', key '{}': {}".format(topic, key, e) |
| 113 | ) |
| tierno | 5c01261 | 2018-04-19 16:01:59 +0200 | [diff] [blame] | 114 | finally: |
| 115 | await self.producer.stop() |
| 116 | |
| garciadeblas | 2644b76 | 2021-03-24 09:21:01 +0100 | [diff] [blame] | 117 | async def aioread( |
| 118 | self, |
| 119 | topic, |
| garciadeblas | 2644b76 | 2021-03-24 09:21:01 +0100 | [diff] [blame] | 120 | callback=None, |
| 121 | aiocallback=None, |
| 122 | group_id=None, |
| 123 | from_beginning=None, |
| 124 | **kwargs |
| 125 | ): |
| tierno | 5c01261 | 2018-04-19 16:01:59 +0200 | [diff] [blame] | 126 | """ |
| tierno | 05ede8f | 2019-01-28 16:20:18 +0000 | [diff] [blame] | 127 | Asyncio read from one or several topics. |
| tierno | 5c01261 | 2018-04-19 16:01:59 +0200 | [diff] [blame] | 128 | :param topic: can be str: single topic; or str list: several topics |
| Benjamin Diaz | 48b78e1 | 2018-10-18 17:55:12 -0300 | [diff] [blame] | 129 | :param callback: synchronous callback function that will handle the message in kafka bus |
| 130 | :param aiocallback: async callback function that will handle the message in kafka bus |
| tierno | 10602af | 2019-02-18 14:53:54 +0000 | [diff] [blame] | 131 | :param group_id: kafka group_id to use. Can be False (set group_id to None), None (use general group_id provided |
| 132 | at connect inside config), or a group_id string |
| tierno | 41ca4d0 | 2020-07-16 11:22:12 +0000 | [diff] [blame] | 133 | :param from_beginning: if True, messages will be obtained from beginning instead of only new ones. |
| 134 | If group_id is supplied, only the not processed messages by other worker are obtained. |
| 135 | If group_id is None, all messages stored at kafka are obtained. |
| Benjamin Diaz | 48b78e1 | 2018-10-18 17:55:12 -0300 | [diff] [blame] | 136 | :param kwargs: optional keyword arguments for callback function |
| 137 | :return: If no callback defined, it returns (topic, key, message) |
| tierno | 5c01261 | 2018-04-19 16:01:59 +0200 | [diff] [blame] | 138 | """ |
| tierno | 10602af | 2019-02-18 14:53:54 +0000 | [diff] [blame] | 139 | if group_id is False: |
| 140 | group_id = None |
| 141 | elif group_id is None: |
| 142 | group_id = self.group_id |
| tierno | 5c01261 | 2018-04-19 16:01:59 +0200 | [diff] [blame] | 143 | try: |
| 144 | if isinstance(topic, (list, tuple)): |
| 145 | topic_list = topic |
| 146 | else: |
| 147 | topic_list = (topic,) |
| garciadeblas | 2644b76 | 2021-03-24 09:21:01 +0100 | [diff] [blame] | 148 | self.consumer = AIOKafkaConsumer( |
| garciadeblas | 2644b76 | 2021-03-24 09:21:01 +0100 | [diff] [blame] | 149 | bootstrap_servers=self.broker, |
| 150 | group_id=group_id, |
| 151 | auto_offset_reset="earliest" if from_beginning else "latest", |
| 152 | ) |
| tierno | 5c01261 | 2018-04-19 16:01:59 +0200 | [diff] [blame] | 153 | await self.consumer.start() |
| 154 | self.consumer.subscribe(topic_list) |
| 155 | |
| 156 | async for message in self.consumer: |
| 157 | if callback: |
| garciadeblas | 2644b76 | 2021-03-24 09:21:01 +0100 | [diff] [blame] | 158 | callback( |
| 159 | message.topic, |
| 160 | yaml.safe_load(message.key), |
| 161 | yaml.safe_load(message.value), |
| 162 | **kwargs |
| 163 | ) |
| Benjamin Diaz | 48b78e1 | 2018-10-18 17:55:12 -0300 | [diff] [blame] | 164 | elif aiocallback: |
| garciadeblas | 2644b76 | 2021-03-24 09:21:01 +0100 | [diff] [blame] | 165 | await aiocallback( |
| 166 | message.topic, |
| 167 | yaml.safe_load(message.key), |
| 168 | yaml.safe_load(message.value), |
| 169 | **kwargs |
| 170 | ) |
| tierno | 5c01261 | 2018-04-19 16:01:59 +0200 | [diff] [blame] | 171 | else: |
| garciadeblas | 2644b76 | 2021-03-24 09:21:01 +0100 | [diff] [blame] | 172 | return ( |
| 173 | message.topic, |
| 174 | yaml.safe_load(message.key), |
| 175 | yaml.safe_load(message.value), |
| 176 | ) |
| tierno | 5c01261 | 2018-04-19 16:01:59 +0200 | [diff] [blame] | 177 | except KafkaError as e: |
| 178 | raise MsgException(str(e)) |
| 179 | finally: |
| 180 | await self.consumer.stop() |