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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.02 12:45:32