【问题标题】:Python Kafka consumer reading already read messagesPython Kafka消费者阅读已经阅读的消息
【发布时间】:2020-06-20 19:23:42
【问题描述】:

Kafka 消费者代码 -

def test():
TOPIC = "file_data"
producer = KafkaProducer()
producer.send(TOPIC, "data")
consumer = KafkaConsumer(
    bootstrap_servers=['localhost:9092'],
    auto_offset_reset='latest',
    consumer_timeout_ms=1000,
    group_id="Group2",
    enable_auto_commit=False,
    auto_commit_interval_ms=1000
)
topic_partition = TopicPartition(TOPIC, 0)
assigned_topic = [topic_partition]
consumer.assign(assigned_topic)
consumer.seek_to_beginning(topic_partition)
for message in consumer:
    print("%s key=%s value=%s" % (message.topic, message.key, message.value))
consumer.commit()

预期行为 它应该只读取生产者写入的最后一条消息。它应该只打印:

file_data key=None value=b'data'

当前行为 运行代码后会打印:

file_data key=None value=b'data'
file_data key=None value=b'data'
file_data key=None value=b'data'
file_data key=None value=b'data'
file_data key=None value=b'data'
file_data key=None value=b'data'

【问题讨论】:

  • 不是 Kafka 专家,但 consumer.seek_to_beginning 调用听起来会导致这种情况?
  • @JohanSchiff 试过 consumer.seek_to_end 但它没有返回任何东西。没有消息被阅读。
  • 您的代码缩进不正确。你能解决这个问题吗?否则,我们必须对您的实际代码做出假设,而这种假设可能是错误的。如果您可以提供minimal reproducible example(一段代码,有人可以复制/粘贴以运行和测试它),那将会有所帮助。

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


【解决方案1】:
from kafka import KafkaConsumer
from kafka import TopicPartition
from kafka import KafkaProducer

def test():
    TOPIC = "file_data"
    producer = KafkaProducer()
    producer.send(TOPIC, b'data')
    consumer = KafkaConsumer(
        bootstrap_servers=['localhost:9092'],
        auto_offset_reset='latest',
        consumer_timeout_ms=1000,
        group_id="Group2",
        enable_auto_commit=False,
        auto_commit_interval_ms=1000
    )
    topic_partition = TopicPartition(TOPIC, 0)
    assigned_topic = [topic_partition]
    consumer.assign(assigned_topic)
    # consumer.seek_to_beginning(topic_partition)
    for message in consumer:
        print("%s key=%s value=%s" % (message.topic, message.key, message.value))
    consumer.commit()
test()

这符合您的期望。如果你想让它从头开始,那么只调用seekToBeginning

参考:seek_to_beginning

【讨论】:

    【解决方案2】:

    您需要使用consumer.seek_to_end(topic_partition),而不是consumer.seek_to_beginning(topic_partition)

    来自文档:

    Kafka 允许使用 seek(TopicPartition, long) 指定位置来指定新位置。还可以使用特殊方法来查找服务器维护的最早和最新偏移量(分别为 seekToBeginning(Collection) 和 seekToEnd(Collection))。

    简单的形式:

    consumer = KafkaConsumer('topicName', bootstrap_servers=[server], group_id=group_id, enable_auto_commit=True)
    consumer.poll()
    consumer.seek_to_end()
    
    for message in consumer:
        print(message)
    
    consumer.close()
    

    【讨论】:

    • consumer.seek_to_end(topic_partition) 没有显示任何结果。它不返回任何消息。
    • 输出保持不变。当 enable_auto_commit=true 和 consumer.seek_to_beginning(topic_partition) 时,将检索所有消息。也没有检索到消息 enable_auto_commit=false 和 consumer.seek_to_beginning(topic_partition)。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2020-04-13
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2019-02-16
    • 2017-11-09
    相关资源
    最近更新 更多