Loading src/nbi/service/vntm_recommend/VntRecommThread.py +30 −24 Original line number Diff line number Diff line Loading @@ -13,8 +13,10 @@ # limitations under the License. import logging, socketio, threading from typing import Dict, List from common.tools.kafka.Variables import KafkaConfig, KafkaTopic from kafka import KafkaConsumer from kafka import KafkaConsumer, TopicPartition from kafka.consumer.fetcher import ConsumerRecord from .Constants import SIO_NAMESPACE, SIO_ROOM logging.getLogger('kafka.client').setLevel(logging.WARNING) Loading Loading @@ -55,32 +57,36 @@ class VntRecommThread(threading.Thread): LOGGER.info('[run] Subscribed') while not self._terminate.is_set(): records = kafka_consumer.poll(timeout_ms=1000, max_records=1) if len(records) == 0: continue # no pending messages... continuing MSG = '[run] records={:s}' LOGGER.debug(MSG.format(str(records))) raise NotImplementedError('parse kafka records and extract recommendation') #if vntm_request.error(): # if vntm_request.error().code() == KafkaError._PARTITION_EOF: continue # MSG = '[run] Consumer error: {:s}' # LOGGER.error(MSG.format(str(vntm_request.error()))) # break #message_key = vntm_request.key().decode('utf-8') #message_value = vntm_request.value().decode('utf-8') #MSG = '[run] Recommendation: key={:s} value={:s}' #LOGGER.debug(MSG.format(str(message_key), str(message_value))) # #LOGGER.debug('[run] checking server namespace...') #server : socketio.Server = self._namespace.server #if server is None: continue #LOGGER.debug('[run] emitting recommendation...') #server.emit('recommendation', message_value, namespace=SIO_NAMESPACE, to=SIO_ROOM) #LOGGER.debug('[run] emitted') topic_records : Dict[TopicPartition, List[ConsumerRecord]] = \ kafka_consumer.poll(timeout_ms=1000, max_records=1) if len(topic_records) == 0: return # no pending records self.process_topic_records(topic_records) LOGGER.info('[run] Closing...') kafka_consumer.close() except: # pylint: disable=bare-except LOGGER.exception('[run] Unexpected Thread Exception') LOGGER.info('[run] Terminated') def process_topic_records( self, topic_records : Dict[TopicPartition, List[ConsumerRecord]] ) -> None: MSG = '[process_topic_records] topic_records={:s}' LOGGER.debug(MSG.format(str(topic_records))) for topic, records in topic_records.items(): if topic.topic == KafkaTopic.VNTMANAGER_REQUEST.value: for record in records: self.emit_recommendation(record) def emit_recommendation(self, record : ConsumerRecord) -> None: message_key = record.key.decode('utf-8') message_value = record.value.decode('utf-8') MSG = '[emit_recommendation] Recommendation: key={:s} value={:s}' LOGGER.debug(MSG.format(str(message_key), str(message_value))) LOGGER.debug('[emit_recommendation] checking server namespace...') server : socketio.Server = self._namespace.server if server is None: return LOGGER.debug('[emit_recommendation] emitting recommendation...') server.emit('recommendation', message_value, namespace=SIO_NAMESPACE, to=SIO_ROOM) LOGGER.debug('[emit_recommendation] emitted') Loading
src/nbi/service/vntm_recommend/VntRecommThread.py +30 −24 Original line number Diff line number Diff line Loading @@ -13,8 +13,10 @@ # limitations under the License. import logging, socketio, threading from typing import Dict, List from common.tools.kafka.Variables import KafkaConfig, KafkaTopic from kafka import KafkaConsumer from kafka import KafkaConsumer, TopicPartition from kafka.consumer.fetcher import ConsumerRecord from .Constants import SIO_NAMESPACE, SIO_ROOM logging.getLogger('kafka.client').setLevel(logging.WARNING) Loading Loading @@ -55,32 +57,36 @@ class VntRecommThread(threading.Thread): LOGGER.info('[run] Subscribed') while not self._terminate.is_set(): records = kafka_consumer.poll(timeout_ms=1000, max_records=1) if len(records) == 0: continue # no pending messages... continuing MSG = '[run] records={:s}' LOGGER.debug(MSG.format(str(records))) raise NotImplementedError('parse kafka records and extract recommendation') #if vntm_request.error(): # if vntm_request.error().code() == KafkaError._PARTITION_EOF: continue # MSG = '[run] Consumer error: {:s}' # LOGGER.error(MSG.format(str(vntm_request.error()))) # break #message_key = vntm_request.key().decode('utf-8') #message_value = vntm_request.value().decode('utf-8') #MSG = '[run] Recommendation: key={:s} value={:s}' #LOGGER.debug(MSG.format(str(message_key), str(message_value))) # #LOGGER.debug('[run] checking server namespace...') #server : socketio.Server = self._namespace.server #if server is None: continue #LOGGER.debug('[run] emitting recommendation...') #server.emit('recommendation', message_value, namespace=SIO_NAMESPACE, to=SIO_ROOM) #LOGGER.debug('[run] emitted') topic_records : Dict[TopicPartition, List[ConsumerRecord]] = \ kafka_consumer.poll(timeout_ms=1000, max_records=1) if len(topic_records) == 0: return # no pending records self.process_topic_records(topic_records) LOGGER.info('[run] Closing...') kafka_consumer.close() except: # pylint: disable=bare-except LOGGER.exception('[run] Unexpected Thread Exception') LOGGER.info('[run] Terminated') def process_topic_records( self, topic_records : Dict[TopicPartition, List[ConsumerRecord]] ) -> None: MSG = '[process_topic_records] topic_records={:s}' LOGGER.debug(MSG.format(str(topic_records))) for topic, records in topic_records.items(): if topic.topic == KafkaTopic.VNTMANAGER_REQUEST.value: for record in records: self.emit_recommendation(record) def emit_recommendation(self, record : ConsumerRecord) -> None: message_key = record.key.decode('utf-8') message_value = record.value.decode('utf-8') MSG = '[emit_recommendation] Recommendation: key={:s} value={:s}' LOGGER.debug(MSG.format(str(message_key), str(message_value))) LOGGER.debug('[emit_recommendation] checking server namespace...') server : socketio.Server = self._namespace.server if server is None: return LOGGER.debug('[emit_recommendation] emitting recommendation...') server.emit('recommendation', message_value, namespace=SIO_NAMESPACE, to=SIO_ROOM) LOGGER.debug('[emit_recommendation] emitted')