【问题标题】:Not able to get newest data from kafka broker using kafka-python library无法使用 kafka-python 库从 kafka 代理获取最新数据
【发布时间】:2019-12-20 01:49:08
【问题描述】:

我正在尝试使用 kafka-python 库从 Kafka 代理消费数据,并且有多个代理以高频率生成数据,但在 Kafka 消费者端我需要大约 5 秒的处理时间,因此在处理第一条消息后在最后一次提交偏移之后,我应该得到最新消息而不是下一条消息。

我试过设置enable_auto_commit=Falseauto_offset_reset="latest"我也试过设置随机组id,我也试过设置group_id = None。这样做的唯一效果是我只在开始时获得最新信息,但之后每个数据都按偏移顺序出现,而不是队列末尾或最新数据。

consumer = KafkaConsumer(bootstrap_servers=kafka_brokers_address,
                      api_version=(2, 3, 0),
                      group_id='abcd',
                      value_deserializer=lambda v:json.loads(v.decode('utf-8')),
                      enable_auto_commit=False,
                      auto_offset_reset="latest")
    consumer_rpnl.assign([TopicPartition('topic', 0)])

c = next(consumer)
## also tried
for c in consumer:
     print(c.values)

【问题讨论】:

  • 你的问题我不清楚。您是否可以使用第一条消息,但不能使用之后的消息?

标签: python apache-kafka kafka-consumer-api kafka-producer-api kafka-python


【解决方案1】:

如何移至最后的示例:https://github.com/dpkp/kafka-python/issues/1405

def seek_to_last():
    consumer = KafkaConsumer(bootstrap_servers=config.kafka_bootstrap_server,
                             group_id=config.kafka_check_proxy_thread_group,
                             value_deserializer=lambda m: json.loads(m.decode('utf-8')), auto_offset_reset='latest',
                             enable_auto_commit=True)

    partitions = consumer.partitions_for_topic(config.kafka_raw_proxy_topic)

    if len(partitions) > config.TASK_OK_PROXY_SCAN_THREAD_N:
        logger.error("...................")

    for partation in partitions:
        p = TopicPartition(config.kafka_raw_proxy_topic, partation)
        mypartition = [p]
        consumer.assign(mypartition)
        # consumer.seek_to_end(p)
        last_pos = consumer.end_offsets(mypartition)
        pos = last_pos[p]
        logger.info("%s,  %s" % (partation, pos))
        # consumer.seek(p, pos)
        offset = OffsetAndMetadata(pos, b'')
        consumer.commit(offsets={p: offset})

【讨论】:

    猜你喜欢
    • 2016-02-13
    • 2018-02-01
    • 2018-05-29
    • 2018-06-12
    • 2018-06-21
    • 1970-01-01
    • 2017-11-11
    • 1970-01-01
    • 2015-11-05
    相关资源
    最近更新 更多