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
相关产品推荐
相关产品推荐

