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

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解释器无法正常解析语法。

解决方案

  1. 升级kafka-python库
    执行以下命令升级到兼容Python 3的最新版本:
pip install --upgrade kafka-python
  1. 指定最低兼容版本(可选)
    如果需要锁定版本范围,可直接安装2.0.0及以上版本(该版本修复了关键字冲突问题):
pip install kafka-python>=2.0.0

升级完成后,无需修改现有业务代码,直接重启Flask项目即可恢复正常。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.18 22:39:56