Team Ai
Datasetpublic

codekingpro/portable-devtools

sourceHugging Faceupdated 5mo agoView on Hugging Face
1likes14kdownloads
kafka_manager.py67 linesDownload Raw Back to socketio
1import logging2import pickle3 4try:5    import kafka6except ImportError:7    kafka = None8 9from .pubsub_manager import PubSubManager10 11logger = logging.getLogger('socketio')12 13 14class KafkaManager(PubSubManager):  # pragma: no cover15    """Kafka based client manager.16 17    This class implements a Kafka backend for event sharing across multiple18    processes.19 20    To use a Kafka backend, initialize the :class:`Server` instance as21    follows::22 23        url = 'kafka://hostname:port'24        server = socketio.Server(client_manager=socketio.KafkaManager(url))25 26    :param url: The connection URL for the Kafka server. For a default Kafka27                store running on the same host, use ``kafka://``. For a highly28                available deployment of Kafka, pass a list with all the29                connection URLs available in your cluster.30    :param channel: The channel name (topic) on which the server sends and31                    receives notifications. Must be the same in all the32                    servers.33    :param write_only: If set to ``True``, only initialize to emit events. The34                       default of ``False`` initializes the class for emitting35                       and receiving.36    """37    name = 'kafka'38 39    def __init__(self, url='kafka://localhost:9092', channel='socketio',40                 write_only=False):41        if kafka is None:42            raise RuntimeError('kafka-python package is not installed '43                               '(Run "pip install kafka-python" in your '44                               'virtualenv).')45 46        super().__init__(channel=channel, write_only=write_only)47 48        urls = [url] if isinstance(url, str) else url49        self.kafka_urls = [url[8:] if url != 'kafka://' else 'localhost:9092'50                           for url in urls]51        self.producer = kafka.KafkaProducer(bootstrap_servers=self.kafka_urls)52        self.consumer = kafka.KafkaConsumer(self.channel,53                                            bootstrap_servers=self.kafka_urls)54 55    def _publish(self, data):56        self.producer.send(self.channel, value=pickle.dumps(data))57        self.producer.flush()58 59    def _kafka_listen(self):60        for message in self.consumer:61            yield message62 63    def _listen(self):64        for message in self._kafka_listen():65            if message.topic == self.channel:66                yield pickle.loads(message.value)67 
codekingpro/portable-devtools · Team Ai