【发布时间】:2019-12-20 01:49:08
【问题描述】:
我正在尝试使用 kafka-python 库从 Kafka 代理消费数据,并且有多个代理以高频率生成数据,但在 Kafka 消费者端我需要大约 5 秒的处理时间,因此在处理第一条消息后在最后一次提交偏移之后,我应该得到最新消息而不是下一条消息。
我试过设置enable_auto_commit=False和auto_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