Python KafkaConsumer初始化报错:Invalid file object: None 求助
Python Kafka消费者报错:ValueError: Invalid file object: None
问题代码
from kafka import KafkaConsumer consumer = KafkaConsumer( 'mykafkatopic', bootstrap_servers='mykafkaserver:9092', group_id='airflow_consumer_test', auto_offset_reset='earliest' ) for msg in consumer: print(msg.value)
报错信息
Traceback (most recent call last): File "C:/Users/ivlasov/pyCharm/pythonProject/main.py", line 9, in <module> consumer = KafkaConsumer( File "C:\Users\ivlasov\AppData\Local\Programs\Python\Python38\lib\site-packages\kafka\consumer\group.py", line 356, in __init__ self._client = KafkaClient(metrics=self._metrics, **self.config) File "C:\Users\ivlasov\AppData\Local\Programs\Python\Python38\lib\site-packages\kafka\client_async.py", line 244, in __init__ self.config['api_version'] = self.check_version(timeout=check_timeout) File "C:\Users\ivlasov\AppData\Local\Programs\Python\Python38\lib\site-packages\kafka\client_async.py", line 909, in check_version version = conn.check_version(timeout=remaining, strict=strict, topics=list(self.config['bootstrap_topics_filter'])) File "C:\Users\ivlasov\AppData\Local\Programs\Python\Python38\lib\site-packages\kafka\conn.py", line 1254, in check_version selector.register(self._sock, selectors.EVENT_READ) File "C:\Users\ivlasov\AppData\Local\Programs\Python\Python38\lib\selectors.py", line 299, in register key = super().register(fileobj, events, data) File "C:\Users\ivlasov\AppData\Local\Programs\Python\Python38\lib\selectors.py", line 238, in register key = SelectorKey(fileobj, self._fileobj_lookup(fileobj), events, data) File "C:\Users\ivlasov\AppData\Local\Programs\Python\Python38\lib\selectors.py", line 225, in _fileobj_lookup return _fileobj_to_fd(fileobj) File "C:\Users\ivlasov\AppData\Local\Programs\Python\Python38\lib\selectors.py", line 39, in _fileobj_to_fd raise ValueError("Invalid file object: " ValueError: Invalid file object: None Process finished with exit code 1
解决方案
这个错误大多是客户端无法成功连接Kafka集群导致的,可按以下方向排查解决:
验证Kafka服务器连通性
先确认mykafkaserver:9092地址能被你的Python程序访问。可以用telnet mykafkaserver 9092或nc -zv mykafkaserver 9092命令测试,若连不通,优先解决网络问题,比如防火墙限制、端口未开放、服务器地址是否输入错误。手动指定Kafka API版本
客户端自动检测API版本失败时,手动指定版本可绕过该问题。修改代码添加api_version参数,示例如下:consumer = KafkaConsumer( 'mykafkatopic', bootstrap_servers='mykafkaserver:9092', group_id='airflow_consumer_test', auto_offset_reset='earliest', api_version=(2, 8, 0) # 替换为你的Kafka实际版本 )可通过Kafka集群的
kafka-topics.sh --version命令查看版本号。升级kafka-python库版本
旧版本的kafka-python可能存在这类连接bug,尝试升级到最新稳定版:pip install --upgrade kafka-python补充SSL配置(若集群启用SSL)
如果你的Kafka集群需要SSL认证,代码未配置相关参数会导致连接失败。需添加SSL相关配置,示例如下:consumer = KafkaConsumer( 'mykafkatopic', bootstrap_servers='mykafkaserver:9092', group_id='airflow_consumer_test', auto_offset_reset='earliest', security_protocol='SSL', ssl_cafile='/path/to/ca.pem', ssl_certfile='/path/to/client-cert.pem', ssl_keyfile='/path/to/client-key.pem' )
内容的提问来源于stack exchange,提问作者Yurii Vlasov
相关产品推荐
相关产品推荐

