You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Kafka Consumer定位指定Offset触发ValueError的原因及解决疑问

Kafka Consumer指定Offset时触发ValueError的问题

环境信息

  • Python版本:3.9.16
  • kafka-python版本:2.0.2
  • 运行环境:MacBook Pro(macOS 11.6.5)

作为Kafka新手,我尝试将Consumer定位到Topic的指定Offset时,频繁触发ValueError,偶尔会无征兆地正常运行,但不清楚具体原因。

测试代码

from kafka import KafkaConsumer, TopicPartition

consumer = KafkaConsumer(bootstrap_servers=['localhost:9092'])

# import pdb
# pdb.set_trace()

myTP = TopicPartition('my-topic', 0)

consumer.assign([myTP])
print ("this is the consumer assignment: {}".format(consumer.assignment()))
# print ("not sure why this will work but printing position: {} ".format(consumer.position(myTP)))
consumer.seek(myTP, 22) 
# print ("not sure why this will work but printing position: {} ".format(consumer.position(myTP)))

for blah in consumer:
    print ("{}, {}".format(blah.offset, blah.value))

常见错误信息

大多数情况下运行代码会抛出ValueError,错误摘要如下:

this is the consumer assignment: {TopicPartition(topic='my-topic', partition=0)}
Traceback (most recent call last):
  File "/Users/my_secret_username/venvs/kafka/lib/python3.9/site-packages/kafka/protocol/types.py", line 20, in _unpack
    (value,) = f(data)
struct.error: unpack requires a buffer of 4 bytes
...
...
...
ValueError: Error encountered when attempting to convert value: b'' to struct format: '<built-in method unpack of _struct.Struct object at 0x10539a930>', hit error: unpack requires a buffer of 4 bytes

临时解决方案

我发现一个临时解决方法:在seek命令前后打印Consumer的position,此时代码总能正常运行。运行结果如下:

$ python tkCons.py 
this is the consumer assignment: {TopicPartition(topic='my-topic', partition=0)}
not sure why this will work but printing position: 34 
not sure why this will work but printing position: 22 
22, b'{"number": 8}'
23, b'{"number": 9}'
24, b'{"number": 0}'
25, b'{"number": 1}'
26, b'{"number": 2}'
27, b'{"number": 3}'
28, b'{"number": 4}'
29, b'{"number": 5}'
30, b'{"number": 6}'
31, b'{"number": 7}'
32, b'{"number": 8}'
33, b'{"number": 9}'

疑问

但我不明白该方案有效的原因:是否需要添加短延迟?打印position是否重置了Consumer内部的某些状态?

完整错误栈

$ python tkCons.py
this is the consumer assignment: {TopicPartition(topic='my-topic', partition=0)}
Traceback (most recent call last):
  File "/Users/my_secret_username/venvs/kafka/lib/python3.9/site-packages/kafka/protocol/types.py", line 20, in _unpack
    (value,) = f(data)
struct.error: unpack requires a buffer of 4 bytes

During handling of the above exception, another exception occurred:

Traceback (most recent call last):
  File "/Users/my_secret_username/kafka/tkCons.py", line 34, in <module>
    for blah in consumer:
  File "/Users/my_secret_username/venvs/kafka/lib/python3.9/site-packages/kafka/consumer/group.py", line 1193, in __next__
    return self.next_v2()
  File "/Users/my_secret_username/venvs/kafka/lib/python3.9/site-packages/kafka/consumer/group.py", line 1201, in next_v2
    return next(self._iterator)
  File "/Users/my_secret_username/venvs/kafka/lib/python3.9/site-packages/kafka/consumer/group.py", line 1116, in _message_generator_v2
    record_map = self.poll(timeout_ms=timeout_ms, update_offsets=False)
  File "/Users/my_secret_username/venvs/kafka/lib/python3.9/site-packages/kafka/consumer/group.py", line 655, in poll
    records = self._poll_once(remaining, max_records, update_offsets=update_offsets)
  File "/Users/my_secret_username/venvs/kafka/lib/python3.9/site-packages/kafka/consumer/group.py", line 702, in _poll_once
    self._client.poll(timeout_ms=timeout_ms)
  File "/Users/my_secret_username/venvs/kafka/lib/python3.9/site-packages/kafka/client_async.py", line 602, in poll
    self._poll(timeout / 1000)
  File "/Users/my_secret_username/venvs/kafka/lib/python3.9/site-packages/kafka/client_async.py", line 687, in _poll
    self._pending_completion.extend(conn.recv())
  File "/Users/my_secret_username/venvs/kafka/lib/python3.9/site-packages/kafka/conn.py", line 1053, in recv
    responses = self._recv()
  File "/Users/my_secret_username/venvs/kafka/lib/python3.9/site-packages/kafka/conn.py", line 1127, in _recv
    return self._protocol.receive_bytes(recvd_data)
  File "/Users/my_secret_username/venvs/kafka/lib/python3.9/site-packages/kafka/protocol/parser.py", line 132, in receive_bytes
    resp = self._process_response(self._rbuffer)
  File "/Users/my_secret_username/venvs/kafka/lib/python3.9/site-packages/kafka/protocol/parser.py", line 138, in _process_response
    recv_correlation_id = Int32.decode(read_buffer)
  File "/Users/my_secret_username/venvs/kafka/lib/python3.9/site-packages/kafka/protocol/types.py", line 64, in decode
    return _unpack(cls._unpack, data.read(4))
  File "/Users/my_secret_username/venvs/kafka/lib/python3.9/site-packages/kafka/protocol/types.py", line 23, in _unpack
    raise ValueError("Error encountered when attempting to convert value: "
ValueError: Error encountered when attempting to convert value: b'' to struct format: '<built-in method unpack of _struct.Struct object at 0x10539a930>', hit error: unpack requires a buffer of 4 bytes

内容的提问来源于stack exchange,提问作者Classified

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.21 07:02:04