Loading src/nbi/service/sse_telemetry/StreamSubscription.py +16 −16 Original line number Diff line number Diff line Loading @@ -55,7 +55,7 @@ KAFKA_BOOT_SERVERS = KafkaConfig.get_kafka_address() class StreamSubscription(Resource): # @HTTP_AUTH.login_required def get(self, subscription_id : int): LOGGER.warning('[get] begin') LOGGER.debug('[get] begin') #db = Engine.get_engine() #if db is None: Loading @@ -71,56 +71,56 @@ class StreamSubscription(Resource): # raise NotFound(description=msg) def event_stream(): LOGGER.warning('[stream:event_stream] begin') LOGGER.debug('[stream:event_stream] begin') topic = 'subscription.{:s}'.format(str(subscription_id)) LOGGER.warning('[stream:event_stream] Checking Topics...') LOGGER.info('[stream:event_stream] Checking Topics...') kafka_admin = KafkaAdminClient(bootstrap_servers=KAFKA_BOOT_SERVERS) existing_topics = set(kafka_admin.list_topics()) LOGGER.warning('[stream:event_stream] existing_topics={:s}'.format(str(existing_topics))) LOGGER.info('[stream:event_stream] existing_topics={:s}'.format(str(existing_topics))) if topic not in existing_topics: LOGGER.warning('[stream:event_stream] Creating Topic...') LOGGER.info('[stream:event_stream] Creating Topic...') to_create = [NewTopic(topic, num_partitions=3, replication_factor=1)] try: kafka_admin.create_topics(to_create, validate_only=False) LOGGER.warning('[stream:event_stream] Topic Created') LOGGER.info('[stream:event_stream] Topic Created') except TopicAlreadyExistsError: pass LOGGER.warning('[stream:event_stream] Connecting Consumer...') LOGGER.info('[stream:event_stream] Connecting Consumer...') kafka_consumer = KafkaConsumer( bootstrap_servers = KAFKA_BOOT_SERVERS, group_id = None, # consumer dispatch all messages sent to subscribed topics auto_offset_reset = 'latest', ) LOGGER.warning('[stream:event_stream] Subscribing topic={:s}...'.format(str(topic))) LOGGER.info('[stream:event_stream] Subscribing topic={:s}...'.format(str(topic))) kafka_consumer.subscribe(topics=[topic]) LOGGER.warning('[stream:event_stream] Subscribed') LOGGER.info('[stream:event_stream] Subscribed') while True: LOGGER.warning('[stream:event_stream] Waiting...') #LOGGER.debug('[stream:event_stream] Waiting...') topic_records : Dict[TopicPartition, List[ConsumerRecord]] = \ kafka_consumer.poll(timeout_ms=1000, max_records=1) if len(topic_records) == 0: time.sleep(0.5) continue # no pending records LOGGER.warning('[stream:event_stream] topic_records={:s}'.format(str(topic_records))) #LOGGER.info('[stream:event_stream] topic_records={:s}'.format(str(topic_records))) for _topic, records in topic_records.items(): if _topic.topic != topic: continue for record in records: message_key = record.key.decode('utf-8') #message_key = record.key.decode('utf-8') message_value = record.value.decode('utf-8') MSG = '[stream:event_stream] message_key={:s} message_value={:s}' LOGGER.warning(MSG.format(str(message_key), str(message_value))) #MSG = '[stream:event_stream] message_key={:s} message_value={:s}' #LOGGER.debug(MSG.format(str(message_key), str(message_value))) yield message_value LOGGER.warning('[stream:event_stream] sent') #LOGGER.debug('[stream:event_stream] Sent') LOGGER.info('[stream:event_stream] Closing...') kafka_consumer.close() LOGGER.warning('[stream] ready to stream...') LOGGER.info('[stream] Ready to stream...') return Response(event_stream(), mimetype='text/event-stream') #update_counter = 1 Loading Loading
src/nbi/service/sse_telemetry/StreamSubscription.py +16 −16 Original line number Diff line number Diff line Loading @@ -55,7 +55,7 @@ KAFKA_BOOT_SERVERS = KafkaConfig.get_kafka_address() class StreamSubscription(Resource): # @HTTP_AUTH.login_required def get(self, subscription_id : int): LOGGER.warning('[get] begin') LOGGER.debug('[get] begin') #db = Engine.get_engine() #if db is None: Loading @@ -71,56 +71,56 @@ class StreamSubscription(Resource): # raise NotFound(description=msg) def event_stream(): LOGGER.warning('[stream:event_stream] begin') LOGGER.debug('[stream:event_stream] begin') topic = 'subscription.{:s}'.format(str(subscription_id)) LOGGER.warning('[stream:event_stream] Checking Topics...') LOGGER.info('[stream:event_stream] Checking Topics...') kafka_admin = KafkaAdminClient(bootstrap_servers=KAFKA_BOOT_SERVERS) existing_topics = set(kafka_admin.list_topics()) LOGGER.warning('[stream:event_stream] existing_topics={:s}'.format(str(existing_topics))) LOGGER.info('[stream:event_stream] existing_topics={:s}'.format(str(existing_topics))) if topic not in existing_topics: LOGGER.warning('[stream:event_stream] Creating Topic...') LOGGER.info('[stream:event_stream] Creating Topic...') to_create = [NewTopic(topic, num_partitions=3, replication_factor=1)] try: kafka_admin.create_topics(to_create, validate_only=False) LOGGER.warning('[stream:event_stream] Topic Created') LOGGER.info('[stream:event_stream] Topic Created') except TopicAlreadyExistsError: pass LOGGER.warning('[stream:event_stream] Connecting Consumer...') LOGGER.info('[stream:event_stream] Connecting Consumer...') kafka_consumer = KafkaConsumer( bootstrap_servers = KAFKA_BOOT_SERVERS, group_id = None, # consumer dispatch all messages sent to subscribed topics auto_offset_reset = 'latest', ) LOGGER.warning('[stream:event_stream] Subscribing topic={:s}...'.format(str(topic))) LOGGER.info('[stream:event_stream] Subscribing topic={:s}...'.format(str(topic))) kafka_consumer.subscribe(topics=[topic]) LOGGER.warning('[stream:event_stream] Subscribed') LOGGER.info('[stream:event_stream] Subscribed') while True: LOGGER.warning('[stream:event_stream] Waiting...') #LOGGER.debug('[stream:event_stream] Waiting...') topic_records : Dict[TopicPartition, List[ConsumerRecord]] = \ kafka_consumer.poll(timeout_ms=1000, max_records=1) if len(topic_records) == 0: time.sleep(0.5) continue # no pending records LOGGER.warning('[stream:event_stream] topic_records={:s}'.format(str(topic_records))) #LOGGER.info('[stream:event_stream] topic_records={:s}'.format(str(topic_records))) for _topic, records in topic_records.items(): if _topic.topic != topic: continue for record in records: message_key = record.key.decode('utf-8') #message_key = record.key.decode('utf-8') message_value = record.value.decode('utf-8') MSG = '[stream:event_stream] message_key={:s} message_value={:s}' LOGGER.warning(MSG.format(str(message_key), str(message_value))) #MSG = '[stream:event_stream] message_key={:s} message_value={:s}' #LOGGER.debug(MSG.format(str(message_key), str(message_value))) yield message_value LOGGER.warning('[stream:event_stream] sent') #LOGGER.debug('[stream:event_stream] Sent') LOGGER.info('[stream:event_stream] Closing...') kafka_consumer.close() LOGGER.warning('[stream] ready to stream...') LOGGER.info('[stream] Ready to stream...') return Response(event_stream(), mimetype='text/event-stream') #update_counter = 1 Loading