【发布时间】: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