基于大模型的个人消费分析和理财助手:开发日志 9

AI 账单分类(下):SSE 流式上传与实时进度反馈

背景与问题

上一篇文章实现了 AI 账单分类的核心逻辑,但接入到实际的上传 API 时,遇到了一个明显的体验问题:

用户上传一个包含 300 条账单的 CSV 文件,后端需要:

  1. 解析 CSV → 300 条 Bill 对象
  2. AI 分类 → 每批 50 条,需要 6 次 LLM API 调用
  3. 存入数据库 → 写入 300 条记录

整个过程同步执行需要 10-30 秒,客户端会一直处于"等待中",用户不知道系统是否还在工作。

解决方案:SSE 流式响应

采用 Server-Sent Events(SSE)将原本同步的请求-响应模型改造为流式模型。

后端改造(FastAPI):

from starlette.responses import StreamingResponse

def _sse_event(data: dict) -> str:
    """生成 SSE 格式的事件文本"""
    return f"data: {json.dumps(data, ensure_ascii=False)}\n\n"

async def event_stream():
    total = len(parsed_data)
    processed = 0
    
    for i in range(0, total, 50):
        batch = parsed_data[i:i+50]
        try:
            results = await classify_bill_batch(batch)
            for j, r in enumerate(results):
                batch[j].transaction_type = r.transaction_type
        except Exception:
            pass  # AI 分类失败不会导致整个上传失败
        
        # 立即写入数据库
        from sqlmodel import Session as DBSession
        with DBSession(engine) as db_session:
            db_session.add_all(batch)
            db_session.commit()
        
        processed += len(batch)
        yield _sse_event({"processed": processed, "total": total})
    
    yield _sse_event({"status": "completed", "total": total})

return StreamingResponse(event_stream(), media_type="text/event-stream")

核心设计意图:

  1. 不再注入 SessionDep 依赖——原始代码使用 FastAPI 的 SessionDep(依赖注入的数据库会话),但在异步生成器函数中,依赖注入的生命周期管理变得复杂。这里改为直接使用 DBSession(engine) 在生成器内部手动管理会话,每批完成后立即提交——这样即使生成器中间出错,已处理的批次不会丢失。

  2. AI 分类失败不阻塞整个流程——如果某批次的 AI 分类失败(LLM API 超时),try/except 静默跳过,保留原始的 transaction_type。账单仍然能成功上传,只是分类可能不够准确。这在设计上遵循了**"功能降级"原则**——核心功能(上传)不能因为辅助功能(分类)的失败而失败。

  3. 进度粒度 = 批次粒度——每处理完 50 条就 yield 一次,不是每 1 条。如果每 1 条都 yield,SSE 事件过于频繁,前端频繁 setState 反而卡顿;如果最后才 yield,又失去了流式的意义。50 条一批是一个合理的平衡点。

前端改造(Dart + Dio):

// API 层 —— 返回 Stream 而非 Future
static Stream<Map<String, dynamic>> uploadWechatBillData(File file) async* {
    final formData = FormData.fromMap({
        "excelData": await MultipartFile.fromFile(file.path, filename: "wechat_bill.xlsx"),
    });

    final response = await HttpUtils().dio.post(
        "/bills/uploadWechatBillData",
        data: formData,
        options: Options(
            contentType: "multipart/form-data",
            responseType: ResponseType.stream,  // 关键:流式响应
            receiveTimeout: null,  // 禁止 Dio 超时中断长连接
        ),
    );

    final stream = response.data.stream as Stream<List<int>>;
    String buffer = '';
    await for (final chunk in stream) {
        buffer += utf8.decode(chunk);
        final lines = buffer.split('\n');
        buffer = lines.removeLast();  // 保留不完整的最后一行
        for (final line in lines) {
            if (line.startsWith('data: ')) {
                final jsonStr = line.substring(6);
                if (jsonStr.isNotEmpty) {
                    yield jsonDecode(jsonStr) as Map<String, dynamic>;
                }
            }
        }
    }
}

Dart Stream 的适配细节:

  1. async* + yield——Dart 的生成器模式天然适配 SSE 的事件流语义。后端 yield 一次,前端 yield 一次,两端对称。

  2. ResponseType.stream——Dio 默认会将完整响应体读入内存再返回,但 SSE 需要逐块处理数据,所以必须改为 stream 模式。

  3. 缓冲池(buffer)处理 TCP 分包——TCP 传输层不保证数据包边界与应用层消息对齐,一个 SSE 事件可能被拆到两个 TCP chunk 中。代码通过维护一个 buffer 变量拼接半行数据,用 split('\n') 切分出完整行,最后一个不完整的行回退到 buffer 等待下一个 chunk。

  4. receiveTimeout: null——如果不禁止 Dio 的超时,长时间的 AI 分类过程会被中断。

UI 层进度显示:

// bill_import_page.dart
await for (final event in widget.onUpload(file)) {
    if (event['status'] == 'completed') {
        showToast("上传成功,共导入 $total 条账单");
    } else {
        setState(() {
            _processedBills = event['processed'] as int? ?? 0;
            _totalBills = event['total'] as int? ?? 0;
        });
    }
}

// UI 渲染
Text(
    _totalBills > 0
        ? "正在AI识别分类... 已处理 $_processedBills/$_totalBills 条"
        : "正在解析账单...",
)

从简单的 CircularProgressIndicator 升级为"已处理 23/150 条"的量化进度提示——这个小小的变化让用户从"不知道要等多久"变成"知道大概还有多久",体验提升显著。

总结

改造前 改造后
后端 API 签名 → Response(同步) → StreamingResponse(SSE)
数据库管理 FastAPI SessionDep 注入 生成器内手动管理
前端 API 返回值 Future<Result<void>> Stream<Map<String, dynamic>>
HTTP 响应模式 ResponseType.json ResponseType.stream
用户界面 转圈菊花 "已处理 N/M 条"量化进度

这个改造的设计意图很明确:让不可预测的长时间操作变得可感知。AI API 的调用延迟是不稳定的——有时 2 秒、有时 10 秒。如果用户面对着静止的转圈,焦虑感会随时间指数增长。流式进度反馈把"我还在等你"变成了"你已经完成了 N 条,还有 M 条"——用户对当前状态有清晰的 mental model,等待体验完全不同。

Logo

AtomGit 是由开放原子开源基金会联合 CSDN 等生态伙伴共同推出的新一代开源与人工智能协作平台。平台坚持“开放、中立、公益”的理念,把代码托管、模型共享、数据集托管、智能体开发体验和算力服务整合在一起,为开发者提供从开发、训练到部署的一站式体验。

更多推荐