3
ñ\ð]º  ã               @   s`   d dl Z d dlZyd dlZW n ek
r4   dZY nX ddlmZ e jdƒZG dd„ deƒZdS )é    Né   )ÚPubSubManagerÚsocketioc                   s>   e Zd ZdZdZd‡ fdd„	Zdd	„ Zd
d„ Zdd„ Z‡  Z	S )ÚKafkaManageraI  Kafka based client manager.

    This class implements a Kafka backend for event sharing across multiple
    processes.

    To use a Kafka backend, initialize the :class:`Server` instance as
    follows::

        url = 'kafka://hostname:port'
        server = socketio.Server(client_manager=socketio.KafkaManager(url))

    :param url: The connection URL for the Kafka server. For a default Kafka
                store running on the same host, use ``kafka://``.
    :param channel: The channel name (topic) on which the server sends and
                    receives notifications. Must be the same in all the
                    servers.
    :param write_only: If set ot ``True``, only initialize to emit events. The
                       default of ``False`` initializes the class for emitting
                       and receiving.
    Úkafkaúkafka://localhost:9092r   Fc                sf   t d krtdƒ‚tt| ƒj||d� |dkr8|dd … nd| _t j| jd�| _t j| j	| jd�| _
d S )NzZkafka-python package is not installed (Run "pip install kafka-python" in your virtualenv).)ÚchannelÚ
write_onlyzkafka://é   zlocalhost:9092)Zbootstrap_servers)r   ÚRuntimeErrorÚsuperr   Ú__init__Z	kafka_urlZKafkaProducerÚproducerZKafkaConsumerr   Úconsumer)ÚselfÚurlr   r	   )Ú	__class__© úA/tmp/pip-build-mqc4i71p/python-socketio/socketio/kafka_manager.pyr   %   s    zKafkaManager.__init__c             C   s&   | j j| jtj|ƒd� | j jƒ  d S )N)Úvalue)r   Úsendr   ÚpickleÚdumpsÚflush)r   Údatar   r   r   Ú_publish4   s    zKafkaManager._publishc             c   s   x| j D ]
}|V  qW d S )N)r   )r   Úmessager   r   r   Ú_kafka_listen8   s    zKafkaManager._kafka_listenc             c   s0   x*| j ƒ D ]}|j| jkr
tj|jƒV  q
W d S )N)r   Ztopicr   r   Úloadsr   )r   r   r   r   r   Ú_listen<   s    zKafkaManager._listen)r   r   F)
Ú__name__Ú
__module__Ú__qualname__Ú__doc__Únamer   r   r   r   Ú__classcell__r   r   )r   r   r      s    r   )	Úloggingr   r   ÚImportErrorZpubsub_managerr   Ú	getLoggerÚloggerr   r   r   r   r   Ú<module>   s   

