FastAPI微服务通过Kafka通信的实现疑问与用户ID赋值需求
基于FastAPI+Kafka的微服务用户ID赋值解决方案
一、原逻辑的可行性判断
原逻辑不可行,核心问题有两个:
- Auth服务的消费者是后台异步任务,不存在FastAPI的请求上下文,
current_user()依赖请求中的认证信息(如JWT Token)才能工作,后台任务没有请求对象,直接调用会抛出上下文缺失的错误。 - 从业务流程看,创建Expense的请求由用户发起,Tracker服务本身就能直接获取用户的认证信息,完全不需要绕到Auth服务去拿user_id,这会增加不必要的系统复杂度和请求延迟。
二、推荐实现方案:直接在Tracker获取用户ID
1. 同步Auth的认证配置到Tracker
Tracker需要和Auth服务共享相同的认证逻辑(比如使用fastapi-users的JWT认证),这样就能在接口中通过依赖获取当前用户:
# Tracker服务的认证依赖配置 from fastapi import Depends from fastapi_users import FastAPIUsers, models from fastapi_users.authentication import JWTAuthentication # 与Auth服务保持完全一致的JWT配置 SECRET = "your-shared-jwt-secret" ALGORITHM = "HS256" jwt_authentication = JWTAuthentication( secret=SECRET, algorithm=ALGORITHM, lifetime_seconds=3600 ) # 复用Auth服务的用户模型(或保持结构一致) class User(models.BaseUser): pass class UserCreate(models.BaseUserCreate): pass class UserUpdate(models.BaseUserUpdate): pass class UserDB(models.BaseUserDB, User): pass # 初始化fastapi_users(需对接Tracker自己的用户查询逻辑,或直接调用Auth的用户接口) fastapi_users = FastAPIUsers[User, models.ID]( user_db=your_user_db_implementation, auth_backends=[jwt_authentication], ) # 定义获取当前活跃用户的依赖 current_active_user = fastapi_users.current_user(active=True)
2. 修改创建Expense的接口
在创建Expense的接口中,通过依赖获取当前用户,直接为Expense的user_id字段赋值:
from fastapi import APIRouter, Depends from pydantic import BaseModel from your_db_module import save_expense_to_db from your_models import Expense, ExpenseCreate router = APIRouter(prefix="/expenses") # 定义前端传入的费用创建请求模型 class ExpenseCreateRequest(BaseModel): source: str amount: float category_id: int @router.post("/", response_model=Expense) async def create_expense( expense_data: ExpenseCreateRequest, current_user: User = Depends(current_active_user) ): # 直接从当前用户对象提取user_id expense_create_obj = ExpenseCreate( **expense_data.dict(), user_id=current_user.id ) # 保存到数据库 db_expense = await save_expense_to_db(expense_create_obj) # 若需通知Auth服务新增费用记录,可保留Kafka生产者调用,但无需返回user_id # await send_expense_info(db_expense.dict()) return db_expense
三、特殊场景下的Kafka跨服务方案(不推荐)
如果因业务限制必须通过Auth服务返回user_id,可按以下流程调整:
1. 修改Tracker生产者:携带认证Token发送消息
在创建Expense时,先获取请求中的用户Token,将Token与费用信息一起发送到Kafka:
import json from aiokafka import AIOKafkaProducer async def send_expense_with_token(expense_info: dict, user_token: str): producer = AIOKafkaProducer( bootstrap_servers="localhost:9092", value_serializer=lambda v: json.dumps(v).encode("utf-8") ) await producer.start() try: # 封装费用信息和用户Token message_payload = { "expense": expense_info, "auth_token": user_token } await producer.send_and_wait("tracker_to_auth", value=message_payload) finally: await producer.stop()
2. 修改Auth服务:验证Token并返回user_id
Auth服务新增一个Kafka生产者,将验证后的user_id和费用ID回传给Tracker:
import json from aiokafka import AIOKafkaConsumer, AIOKafkaProducer from fastapi_users.jwt import decode_jwt from your_auth_config import SECRET, ALGORITHM # 新增回复Tracker的生产者 async def send_user_id_response(expense_id: int, user_id: str): producer = AIOKafkaProducer( bootstrap_servers="localhost:9092", value_serializer=lambda v: json.dumps(v).encode("utf-8") ) await producer.start() try: await producer.send_and_wait( "auth_to_tracker", value={"expense_id": expense_id, "user_id": user_id} ) finally: await producer.stop() # 修改消费者逻辑:验证Token并返回user_id async def process_expense_messages(): consumer = AIOKafkaConsumer( "tracker_to_auth", bootstrap_servers="localhost:9092", value_deserializer=lambda v: json.loads(v.decode("utf-8")) ) await consumer.start() try: async for msg in consumer: payload = msg.value expense_info = payload["expense"] auth_token = payload["auth_token"] # 验证JWT Token获取user_id try: decoded_token = decode_jwt(auth_token, SECRET, [ALGORITHM]) user_id = decoded_token["sub"] # 发送回复给Tracker await send_user_id_response(expense_info["id"], user_id) except Exception as e: print(f"Token validation failed: {str(e)}") finally: await consumer.stop()
3. Tracker新增消费者:接收user_id并更新Expense
Tracker添加消费者监听Auth的回复消息,更新数据库中Expense的user_id字段:
import json from aiokafka import AIOKafkaConsumer from your_db_module import update_expense_user_id async def process_auth_response(): consumer = AIOKafkaConsumer( "auth_to_tracker", bootstrap_servers="localhost:9092", value_deserializer=lambda v: json.loads(v.decode("utf-8")) ) await consumer.start() try: async for msg in consumer: response_data = msg.value await update_expense_user_id( response_data["expense_id"], response_data["user_id"] ) finally: await consumer.stop()
4. Tracker启动时启动后台消费者
在FastAPI应用启动事件中,将消费者任务加入异步后台任务:
from fastapi import FastAPI import asyncio app = FastAPI() @app.on_event("startup") async def startup(): # 启动Auth回复的消费者 asyncio.create_task(process_auth_response())
关键注意事项
- 优先选择直接在Tracker获取user_id的方案,避免跨服务消息传递带来的延迟、消息丢失、重复消费等问题。
- 若使用Kafka方案,需实现幂等更新(比如通过费用ID唯一标识),防止重复消息导致数据异常。
- 两个服务的JWT配置(密钥、算法、过期时间)必须完全一致,否则Token验证会失败。
内容的提问来源于stack exchange,提问作者Valerka288
相关产品推荐
相关产品推荐

