基于大模型的个人消费分析和理财助手:开发日志 9
基于大模型的个人消费分析和理财助手:开发日志 9
AI 账单分类(下):SSE 流式上传与实时进度反馈
背景与问题
上一篇文章实现了 AI 账单分类的核心逻辑,但接入到实际的上传 API 时,遇到了一个明显的体验问题:
用户上传一个包含 300 条账单的 CSV 文件,后端需要:
- 解析 CSV → 300 条 Bill 对象
- AI 分类 → 每批 50 条,需要 6 次 LLM API 调用
- 存入数据库 → 写入 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")
核心设计意图:
-
不再注入
SessionDep依赖——原始代码使用 FastAPI 的SessionDep(依赖注入的数据库会话),但在异步生成器函数中,依赖注入的生命周期管理变得复杂。这里改为直接使用DBSession(engine)在生成器内部手动管理会话,每批完成后立即提交——这样即使生成器中间出错,已处理的批次不会丢失。 -
AI 分类失败不阻塞整个流程——如果某批次的 AI 分类失败(LLM API 超时),
try/except静默跳过,保留原始的transaction_type。账单仍然能成功上传,只是分类可能不够准确。这在设计上遵循了**"功能降级"原则**——核心功能(上传)不能因为辅助功能(分类)的失败而失败。 -
进度粒度 = 批次粒度——每处理完 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 的适配细节:
-
async*+yield——Dart 的生成器模式天然适配 SSE 的事件流语义。后端 yield 一次,前端 yield 一次,两端对称。 -
ResponseType.stream——Dio 默认会将完整响应体读入内存再返回,但 SSE 需要逐块处理数据,所以必须改为 stream 模式。 -
缓冲池(buffer)处理 TCP 分包——TCP 传输层不保证数据包边界与应用层消息对齐,一个 SSE 事件可能被拆到两个 TCP chunk 中。代码通过维护一个
buffer变量拼接半行数据,用split('\n')切分出完整行,最后一个不完整的行回退到 buffer 等待下一个 chunk。 -
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,等待体验完全不同。
AtomGit 是由开放原子开源基金会联合 CSDN 等生态伙伴共同推出的新一代开源与人工智能协作平台。平台坚持“开放、中立、公益”的理念,把代码托管、模型共享、数据集托管、智能体开发体验和算力服务整合在一起,为开发者提供从开发、训练到部署的一站式体验。
更多推荐



所有评论(0)