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

使用LangChain搭配AlloyDB进行向量搜索时出现未关闭Client Session及连接器问题

使用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.14 14:40:28