【问题标题】:Kafka Python, How to track consumer started in a different processKafka Python,如何跟踪在不同进程中开始的消费者
【发布时间】:2019-08-29 06:38:33
【问题描述】:

我对 Python 还很陌生,刚开始接触 Kafka,如果我在某个地方错了,请原谅我的术语。

所以我有一个基于 Django 的 Web 应用程序,我在同一进程中通过 Kafka Producer 发送 json 消息。 然而,在务实地创建主题的同时,我还在针对该特定主题的单独流程中启动(订阅)新消费者。

#Consumer code snippet

 if topic_name is not None :
        #Create topic
        create_kafka_topic_instance(topic_name)
        #Initialize a consumer and subscribe to topic
        Process(target=init_kafka_consumer_instance, args=(topic_name))

def forgiving_json_deserializer(v):
    if v is None :
        return
    try:
        return json.loads(v.decode('utf-8'))
    except json.decoder.JSONDecodeError:
        import traceback
        print(traceback.format_exc())
        return None

def init_kafka_consumer_instance(topic, group_id=None):
    try:
        if topic is None:
            raise Exception("Invalid argument topic")
        comsumer = None
        comsumer = KafkaConsumer(topic, bootstrap_servers=[KAFKA_BROKER_URL], auto_offset_reset="earliest",
           urn comsumer
    except Exception as e:
        import traceback
        print(traceback.format_exc())
    return Noneurn comsumer
    except Exception as e:
        import traceback
        print(traceback.format_exc())
    return None

生产者代码片段

# assuming obj is a model instance
        serialized_obj = serializers.serialize('json', [ order, ])
        #send_message(topic_name,order)
        producer = KafkaProducer(bootstrap_servers=[KAFKA_BROKER_URL], value_serializer=lambda x: json.dumps(x).encode('utf-8'))
        x = producer.send("test", serialized_obj)
        producer.flush()

现在我有一些疑问,所以如果我的 Django 应用程序(服务器)以某种方式重新启动,我是否仍然让消费者收听该主题。

另外,我在消费者中有一些打印语句,我无法在我的服务器控制台中看到这些语句。

但是,在 python shell 中编写相同的代码 sn-p(初始化消费者),我可以在那里看到 print 语句中的消息,这意味着我的 Producer 工作正常。

【问题讨论】:

    标签: python python-3.x apache-kafka kafka-python


    【解决方案1】:

    Kafka 服务器不依赖于您的 Django 应用程序(服务器)。但是你的消费者是肯定的。

    所以您的主题在 Kafka 服务器中仍然存在(如果 kafka 服务器死了,那就是另一回事了),但是您的消费者会随着您的应用程序重新启动。

    因此,如果您希望您的消费者正常工作,请将其设置为与您的应用并行工作的 Worker,并且在您的应用关闭时不会重新启动

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2019-03-26
      • 1970-01-01
      • 2017-02-15
      • 2020-05-26
      • 2021-01-11
      • 1970-01-01
      • 2019-07-01
      • 2019-01-18
      相关资源
      最近更新 更多