Loading src/common/tools/kafka/Variables.py +31 −26 Original line number Diff line number Diff line Loading @@ -16,6 +16,7 @@ import logging, time from enum import Enum #from confluent_kafka.admin import AdminClient, NewTopic from kafka.admin import KafkaAdminClient, NewTopic from kafka.errors import TopicAlreadyExistsError from common.Settings import get_setting Loading Loading @@ -88,6 +89,7 @@ class KafkaTopic(Enum): #create_topic_future_map = kafka_admin_client.create_topics(missing_topics, request_timeout=5*60) #LOGGER.debug('create_topic_future_map: {:s}'.format(str(create_topic_future_map))) try: topics_result = kafka_admin_client.create_topics( new_topics=missing_topics, timeout_ms=KAFKA_TOPIC_CREATE_REQUEST_TIMEOUT, validate_only=False Loading Loading @@ -115,6 +117,9 @@ class KafkaTopic(Enum): if len(failed_topic_creations) > 0: return False LOGGER.debug('All topics created.') except TopicAlreadyExistsError: LOGGER.debug('Some topics already exists.') # Wait until topics appear in metadata desired_topics = {topic.value for topic in KafkaTopic} missing_topics = set() Loading Loading
src/common/tools/kafka/Variables.py +31 −26 Original line number Diff line number Diff line Loading @@ -16,6 +16,7 @@ import logging, time from enum import Enum #from confluent_kafka.admin import AdminClient, NewTopic from kafka.admin import KafkaAdminClient, NewTopic from kafka.errors import TopicAlreadyExistsError from common.Settings import get_setting Loading Loading @@ -88,6 +89,7 @@ class KafkaTopic(Enum): #create_topic_future_map = kafka_admin_client.create_topics(missing_topics, request_timeout=5*60) #LOGGER.debug('create_topic_future_map: {:s}'.format(str(create_topic_future_map))) try: topics_result = kafka_admin_client.create_topics( new_topics=missing_topics, timeout_ms=KAFKA_TOPIC_CREATE_REQUEST_TIMEOUT, validate_only=False Loading Loading @@ -115,6 +117,9 @@ class KafkaTopic(Enum): if len(failed_topic_creations) > 0: return False LOGGER.debug('All topics created.') except TopicAlreadyExistsError: LOGGER.debug('Some topics already exists.') # Wait until topics appear in metadata desired_topics = {topic.value for topic in KafkaTopic} missing_topics = set() Loading