ARTICLE · INTELLIGENCE

战地情报 · 详情页

来自尧图项目组的一线实战观察与深度解析

如何在 FastAPI 等异步服务中用 Mem0 AsyncMemory 并发读写记忆

如何在 FastAPI 等异步服务中用 Mem0 AsyncMemory 并发读写记忆 如何在 FastAPI 等异步服务中用 Mem0 AsyncMemory 并发读写记忆【免费下载链接】embedchainThe Memory Layer for AI Agents - Drop-in memory infrastructure for AI agents and apps. Context that persists. Built for production.项目地址: https://gitcode.com/GitHub_Trending/em/embedchain如果你在 FastAPI、后台 worker 或其他基于asyncio的 Python 服务里需要读写 Mem0 记忆同步客户端会阻塞事件循环。Mem0 的 AsyncMemory 功能文档 提供了非阻塞的AsyncMemory接口它直连同一套存储后端、所有方法都与同步版一一对应并且可以直接在 FastAPI 路由里await配合asyncio.gather并发执行多路记忆操作。本文的目标是把AsyncMemory接入一个 FastAPI 服务完成一次并发写、一次并发读并给出文档中给出的验证方式。前提条件依据 Python SDK QuickstartPython 3.10 或更高版本一个 OpenAI API key并设置环境变量export OPENAI_API_KEYyour-openai-api-keyyour-openai-api-key替换为你自己的 key通过 pip 安装 SDKpip install mem0ai。默认的MemoryConfig会装配 OpenAIgpt-5-mini做事实抽取、text-embedding-3-small做嵌入1536 维、Qdrant 向量库数据落在/tmp/qdrant、SQLite 历史记录~/.mem0/history.db不带 reranker。需要换 LLM、embedder 或向量库时走 Configure components 一节的配置方式。初始化客户端默认配置与 MemoryConfigAsyncMemory从mem0顶层导出支持默认配置和自定义MemoryConfig两种方式import asyncio from mem0 import AsyncMemory # Default configuration memory AsyncMemory() # Custom configuration from mem0.configs.base import MemoryConfig custom_config MemoryConfig( # Your custom configuration here ) memory AsyncMemory(configcustom_config)文档给出一条初始化后的快速自检初始化后立刻跑一次await memory.search(...)如果能无错误返回记忆说明配置可用。在 FastAPI 这类长进程服务里文档建议每个进程只创建一次AsyncMemory实例在启动钩子里完成配置并复用它避免重复建后端连接。需要干净关闭时可以用异步上下文管理器包住客户端import asyncio from contextlib import asynccontextmanager from mem0 import AsyncMemory asynccontextmanager async def get_memory(): memory AsyncMemory() try: yield memory finally: # Clean up resources if needed pass async def safe_memory_usage(): async with get_memory() as memory: return await memory.search(test query, filters{user_id: alice})一条硬约束AsyncMemory依赖运行中的事件循环必须在async def函数里调用或通过asyncio.run()之类的辅助入口调用否则会出运行时错误。在 FastAPI 路由中读写记忆下面是文档中的 FastAPI 接线示例模块级共享一个实例写与查两个端点分别awaitfrom fastapi import FastAPI, HTTPException from mem0 import AsyncMemory app FastAPI() memory AsyncMemory() app.post(/memories/) async def add_memory(messages: list, user_id: str): try: result await memory.add(messagesmessages, user_iduser_id) return {status: success, data: result} except Exception as exc: raise HTTPException(status_code500, detailstr(exc)) app.get(/memories/search) async def search_memories(query: str, user_id: str, limit: int 10): try: result await memory.search(queryquery, filters{user_id: user_id}, top_klimit) return {status: success, data: result} except Exception as exc: raise HTTPException(status_code500, detailstr(exc))两个 API 使用上要注意的点来自 AsyncMemory 源码与文档add()接受顶层user_id/agent_id/run_id参数而search()和get_all()只接受filters{user_id: ..., agent_id: ..., run_id: ...}形式且 filters 必须至少包含user_id、agent_id、run_id之一把顶层参数传给search()/get_all()会抛ValueError。delete_all必须至少带user_id、agent_id、run_id中一个三者全给可以把删除范围收窄到单个会话。get(memory_id...)在 ID 无效时抛ValueErrorsearch返回 dict记忆在results键下与同步客户端形状一致。用 asyncio.gather 并发写入并发是本场景的核心。文档给出的批处理模式是把多个memory.add(...)协程塞进任务列表用asyncio.gather(*tasks, return_exceptionsTrue)一次性执行再逐个检查结果是异常还是成功import asyncio from mem0 import AsyncMemory async def batch_operations(): memory AsyncMemory() tasks [ memory.add( messages[{role: user, content: fMessage {i}}], user_idfuser_{i} ) for i in range(5) ] results await asyncio.gather(*tasks, return_exceptionsTrue) for i, result in enumerate(results): if isinstance(result, Exception): print(fTask {i} failed: {result}) else: print(fTask {i} completed successfully)文档对成功形态的说明是并发工作正常时成功的任务返回 memory ID失败的任务以异常形式出现在results列表中。return_exceptionsTrue的作用是让单个任务失败不会中断整批而是以Exception实例的形式留在结果里方便逐条上报。文档同时建议用asyncio.gather提吞吐时按后端容量控制并发上限避免过大批次触发内存问题。检索、更新与审计与同步客户端相同的一整套操作在异步版中都有对应方法# Create memories result await memory.add( messages[ {role: user, content: Im travelling to SF}, {role: assistant, content: Thats great to hear!} ], user_idalice ) # Search memories results await memory.search( queryWhere am I travelling?, filters{user_id: alice} ) # List memories all_memories await memory.get_all(filters{user_id: alice}) # Get a specific memorymemory-id-here 替换为 add/get_all 返回的真实 ID specific_memory await memory.get(memory_idmemory-id-here) # Update a memory updated_memory await memory.update( memory_idmemory-id-here, textIm travelling to Seattle ) # Delete a memory await memory.delete(memory_idmemory-id-here) # Delete scoped memories await memory.delete_all(user_idalice)按user_id/agent_id/run_id三级范围组织记忆并用history拉审计日志也是异步版支持的用法await memory.add( messages[{role: user, content: I prefer vegetarian food}], user_idalice, agent_iddiet-assistant, run_idconsultation-001 ) agent_memories await memory.get_all(filters{user_id: alice, agent_id: diet-assistant}) session_memories await memory.get_all(filters{user_id: alice, run_id: consultation-001}) history await memory.history(memory_idmemory-id-here) # 同上替换为真实 ID给记忆操作加超时与重试如果服务依赖外部 LLM/网络后端文档给了一层带超时和指数退避的重试包装import asyncio from mem0 import AsyncMemory async def with_timeout_and_retry(operation, max_retries3, timeout10.0): for attempt in range(max_retries): try: return await asyncio.wait_for(operation(), timeouttimeout) except asyncio.TimeoutError: print(fTimeout on attempt {attempt 1}) except Exception as exc: print(fError on attempt {attempt 1}: {exc}) if attempt max_retries - 1: await asyncio.sleep(2 ** attempt) raise Exception(fOperation failed after {max_retries} attempts) async def robust_memory_search(): memory AsyncMemory() async def search_operation(): return await memory.search(test query, filters{user_id: alice}) return await with_timeout_and_retry(search_operation)文档对此的提醒是重试次数必须设上限失控的重试循环会长时间占住事件循环、阻塞其他任务。验证接入是否成功文档给出的验证清单初始化后立即跑一次await memory.search(...)能无错返回记忆说明配置正确跑一次快速的 add/search 循环确认返回的记忆内容与写入的输入匹配检查应用日志确认异步任务正常完成、没有阻塞事件循环在 FastAPI 中通过端点如健康检查确认共享客户端能处理并发请求观察重试计数器异常上涨通常指向配置或连通性问题。响应形状也是判断依据每次调用应返回与同步客户端相同的字段ID、results或确认对象如果返回缺键通常是协程没有被await。已知问题与边界文档的故障对照表按现象定位原因IssuePossible causesFixInitialization failsMissing dependencies, invalid configValidateMemoryConfigsettings and environment variables.Slow operationsLarge datasets, network latencyCache heavy queries and tune vector store parameters.Memory not foundInvalid ID or deleted recordCheck ID source and handle soft-deleted states.Connection timeoutsNetwork issues, overloaded backendApply retries/backoff and inspect infrastructure health.Out-of-memory errorsOversized batchesReduce concurrency or chunk operations into smaller sets.其他边界初始化时的配置错误ValueError和连接错误ConnectionError建议在构造AsyncMemory(config...)时就捕获并打印文档示例见 async-memory.mdx 的 Handle errors gracefully 一节ValueError类错误无效 memory ID、空查询串要显式捕获并记录否则异步堆栈在后台任务里容易丢失忘写await是最常见的漏写记忆的成因文档建议用 lint 或辅助包装强制检查如果向量库不支持关键词检索初始化时会打 warning 并禁用混合BM25打分search退化为纯语义相似度——这是 mem0/memory/main.py 中AsyncMemory构造逻辑的既有行为。更完整的端点、参数与错误说明见 Async Memory 功能文档如果要换掉默认的 OpenAI/Qdrant 组件从 Configure components 入手。【免费下载链接】embedchainThe Memory Layer for AI Agents - Drop-in memory infrastructure for AI agents and apps. Context that persists. Built for production.项目地址: https://gitcode.com/GitHub_Trending/em/embedchain创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
RELATED READING

延伸阅读

更多一线实战笔记与深度复盘,助您持续精进