Flask集成KafkaConsumer触发SimpleProducer语法错误的技术求助
Flask集成KafkaConsumer启动报错:async关键字语法冲突问题
问题场景
启动Flask后端项目时,导入KafkaConsumer触发语法错误,错误源于kafka库内部代码使用Python3关键字async作为属性名,导致解析失败。
相关代码文件
run.py
from app import create_app if __name__ == '__main__': app = create_app() app.run(host='0.0.0.0', port=8000, debug=True)
/app/app/init.py
from flask import Flask from flask_cors import CORS from app.config import Config from threading import Thread from app.consumer.kafka_consumer import consume_messages def create_app(): app = Flask(__name__) app.config.from_object(Config) app.config['SQLALCHEMY_DATABASE_URI'] = Config.MYSQL_DATABASE_URI app.config['SQLALCHEMY_TRACK_MODIFICATIONS'] = False CORS(app) consumer_thread = Thread(target=consume_messages) consumer_thread.start() # Register blueprints from app.controller.segmentation_controller import segmentation_bp app.register_blueprint(segmentation_bp) return app
/app/app/consumer/kafka_consumer.py
from kafka import KafkaConsumer from app.config import Config from app.service.segmentation_service import SegmentationService def consume_messages(): segmentation_service = SegmentationService() consumer = KafkaConsumer( Config.KAFKA_TOPIC, bootstrap_servers=['localhost:9092'], auto_offset_reset='latest', # 从最新可用偏移量开始读取 enable_auto_commit=True, group_id='my-group', value_deserializer=lambda x: x.decode('utf-8'), ) for message in consumer: segmentation_service.process_messages(message)
错误信息
2023-06-14 11:44:43 Traceback (most recent call last): 2023-06-14 11:44:43 File "run.py", line 1, in <module> 2023-06-14 11:44:43 from app import create_app 2023-06-14 11:44:43 File "/app/app/__init__.py", line 5, in <module> 2023-06-14 11:44:43 from app.consumer.kafka_consumer import consume_messages 2023-06-14 11:44:43 File "/app/app/consumer/kafka_consumer.py", line 1, in <module> 2023-06-14 11:44:43 from kafka import KafkaConsumer 2023-06-14 11:44:43 File "/usr/local/lib/python3.8/site-packages/kafka/__init__.py", line 23, in <module> 2023-06-14 11:44:43 from kafka.producer import KafkaProducer 2023-06-14 11:44:43 File "/usr/local/lib/python3.8/site-packages/kafka/producer/__init__.py", line 4, in <module> 2023-06-14 11:44:43 from .simple import SimpleProducer 2023-06-14 11:44:43 File "/usr/local/lib/python3.8/site-packages/kafka/producer/simple.py", line 54 2023-06-14 11:44:43 return '<SimpleProducer batch=%s>' % self.async
问题原因
使用的kafka-python库版本过旧,该版本代码中误用了Python 3的保留关键字async作为类属性名,导致Python解释器无法正常解析语法。
解决方案
- 升级kafka-python库
执行以下命令升级到兼容Python 3的最新版本:
pip install --upgrade kafka-python
- 指定最低兼容版本(可选)
如果需要锁定版本范围,可直接安装2.0.0及以上版本(该版本修复了关键字冲突问题):
pip install kafka-python>=2.0.0
升级完成后,无需修改现有业务代码,直接重启Flask项目即可恢复正常。
内容的提问来源于stack exchange,提问作者Nabih Bawazir
相关产品推荐
相关产品推荐

