使用LangChain搭配AlloyDB进行向量搜索时出现未关闭Client Session及连接器问题
问题描述
我在使用LangChain和Google AlloyDB做向量搜索时碰到了个棘手的问题:我专门写了异步上下文管理器来管理AlloyDB引擎和aiohttp.ClientSession这类资源的初始化与清理,明明已经调用了这些资源的关闭方法,但日志里还是跳出了未关闭的错误:
2025-01-22 10:22:19 - [BOT-ESPECIALIST] - INFO: Closed connections with AlloyDB. 2025-01-22 10:22:20 - [BOT-ESPECIALIST] - ERROR: Unclosed client session client_session: <aiohttp.client.ClientSession object at 0x7cccbee73490> 2025-01-22 10:22:20 - [BOT-ESPECIALIST] - ERROR: Unclosed connector connections: ['deque([(<aiohttp.client_proto.ResponseHandler object at 0x7cccbee67ac0>, 112391.933409363), (<aiohttp.client_proto.ResponseHandler object at 0x7cccbee1e040>, 112392.240980781)])'] connector: <aiohttp.connector.TCPConnector object at 0x7cccbee734c0>
我的技术栈
- LangChain:用来和AlloyDB交互,处理向量存储与搜索逻辑
- Google AlloyDB:存储和查询向量数据的托管数据库
- aiohttp.ClientSession:AlloyDB的LangChain集成内部用来发起异步HTTP请求
相关代码
带异步上下文管理器的AlloyDB类
import aiohttp import asyncio import logging from typing import List, Union from langchain_google_alloydb_pg import AlloyDBEngine, AlloyDBVectorStore from langchain_core.documents import Document from langchain_core.embeddings.embeddings import Embeddings class AlloyDB: def __init__(self, connection, embedding_model): self.engine: Union[AlloyDBEngine, None] = None self.connection = connection self.embedding_model = embedding_model self.vector_store: Union[AlloyDBVectorStore, None] = None self.session = aiohttp.ClientSession(timeout=aiohttp.ClientTimeout(total=4)) async def __aenter__(self): """Initialize resources.""" self.engine = await AlloyDBEngine.afrom_instance( project_id=self.connection.project_id, region=self.connection.region, cluster=self.connection.cluster, instance=self.connection.instance, database=self.connection.database, user=self.connection.db_user, password=self.connection.db_password, ) return self async def __aexit__(self, exc_type, exc_value, traceback): """Clean up resources.""" if self.engine: self.engine.close() if not self.session.closed: await self.session.close() logging.info("Closed connections with AlloyDB.")
主函数使用示例
async def main(): async with AlloyDB(connection, embedding_model) as db: await db.init_vector_storage_table(table="products", table_config=table_config) query = "I'd like a fruit." docs = await db.search_documents(query) print(docs) loop = asyncio.get_event_loop() try: loop.run_until_complete(main()) finally: loop.close()
我的观察
- 资源清理操作:我在
__aexit__方法里明确调用了self.engine.close()和await self.session.close(),但错误依然没有消失。 - 错误细节:日志显示
ClientSession和关联的连接器未关闭,尽管我已经执行了关闭操作。
我的疑问
- 明明显式调用了
close(),为什么aiohttp.ClientSession还是没彻底关闭? - LangChain的AlloyDB集成内部使用aiohttp的方式会不会有问题,导致这个错误出现?
- 在这种异步上下文管理器里管理
aiohttp.ClientSession,有什么推荐的模式或最佳实践吗?
我的分析与解决方案
我来帮你捋捋可能的原因和解决办法:
1. ClientSession未真正关闭的核心原因
你在__init__里自己创建了self.session,但很可能LangChain的AlloyDB集成内部自己单独创建了ClientSession实例,根本没用到你传入的这个——也就是说你关闭的只是自己创建的session,而集成内部的session完全没被处理。另外,虽然你调用了await self.session.close(),但如果你的业务方法(比如init_vector_storage_table或search_documents)里还有未完成的异步请求,这些请求会持有session的连接,导致session无法彻底关闭。
2. LangChain AlloyDB集成的潜在问题
AlloyDBEngine在异步初始化时,大概率会自行创建aiohttp.ClientSession,而且没有把这个实例暴露给外部让你关闭。你可以去看看AlloyDBEngine的源码,是不是它内部维护了自己的session,而你调用的self.engine.close()是同步方法,根本没法处理异步的session关闭逻辑?毕竟关闭session是异步操作,同步方法里没法等待它完成。
3. 推荐的资源管理优化方案
- 让LangChain复用你创建的session:检查
AlloyDBEngine.afrom_instance的参数列表,看看有没有可以传入自定义aiohttp.ClientSession的选项。如果有的话,把你创建的self.session传进去,这样你关闭自己的session时,就能同时关闭LangChain内部使用的那个实例。 - 改用异步方式关闭引擎:如果
AlloyDBEngine提供了异步关闭方法(比如aclose()),别再用同步的close(),换成await self.engine.aclose(),这样它内部的异步资源(比如自己创建的session)才能被正确清理。 - 增加短延迟确保关闭完成:在
__aexit__里,调用await asyncio.sleep(0.2)之类的短延迟,确保session的关闭操作完全完成后再退出上下文管理器。这算是临时的 workaround,但能帮你验证是不是关闭操作没来得及执行完。 - 用
async with管理ClientSession:把self.session的初始化放到__aenter__里,用async with aiohttp.ClientSession(...) as self.session:的方式创建,这样session会自动在上下文退出时关闭,减少手动管理的风险。
举个修改后的__aexit__例子,如果引擎有异步关闭方法:
async def __aexit__(self, exc_type, exc_value, traceback): """Clean up resources.""" if self.engine: # 优先用异步关闭方法 if hasattr(self.engine, 'aclose'): await self.engine.aclose() else: self.engine.close() if not self.session.closed: await self.session.close() # 等待片刻确保连接完全释放 await asyncio.sleep(0.2) logging.info("Closed connections with AlloyDB.")
另外,你可以在关闭session后加个日志,确认session的状态:
await self.session.close() logging.info(f"Session closed status: {self.session.closed}")
如果日志显示True但还是报错,那基本可以确定是LangChain内部还有其他未被处理的session。
备注:内容来源于stack exchange,提问作者Lucas Vital

