多彩编程 多彩编程MZPH · CODE BLOG
ARTICLE DETAIL

文章详情

深耕前端与后端开发技术的一线实战笔记与踩坑复盘。

openJiuwen agent-core GaussDB 异步存储适配:GaussDbStore 用法与 SQLAlchemy Dialect 实现深度解析

openJiuwen agent-core GaussDB 异步存储适配:GaussDbStore 用法与 SQLAlchemy Dialect 实现深度解析 人工智能AI AgentAgent 框架大模型工具调用RAG提示工程强化学习【免费下载链接】agent-coreopenJiuwen agent-core可提供AI Agent开发、运行、调优与演进相关的全套SDK能力项目地址https://gitcode.com/openJiuwen/agent-core点击查看免费下载导读GaussDbStore是 openJiuwen agent-core 扩展模块openjiuwen.extensions.store.db中面向 GaussDB 数据库的异步存储实现。它基于 SQLAlchemy 的方言Dialect扩展机制通过import_dbapi()局部加载async_gaussdb驱动向上层提供与DefaultDbStore完全一致的AsyncEngine访问接口。阅读本文后你将掌握如何配置gaussdbasync_gaussdb连接串、如何用GaussDbStore完成 CRUD 与事务操作、如何将其作为LongTermMemory的db_store替换默认存储以及该实现避开全局sys.modules污染、实现最小核心、可选扩展的设计原理与源码级细节。一、定位与模块结构扩展而非核心GaussDbStore位于扩展模块openjiuwen/extensions/store/db/与核心模块openjiuwen.core.foundation.store.db中的DefaultDbStore形成对照。这一目录划分本身就是设计的一部分gauss_db_store.pyGaussDbStore类定义整个文件仅 31 行是一个极薄封装gauss_dialect.pySQLAlchemy 方言dialect注册与驱动适配逻辑这是本模块的技术核心init.py对外导出GaussDbStore一个符号。从 pyproject.toml 可以看出gaussdb是独立可选依赖gaussdb [async-gaussdb~0.30.0, asyncpg0.29.0]而all-storagepyproject.toml则将sqlite、postgres、mysql、gaussdb、redis、elasticsearch、obs全部打包方便一次装齐all-storage [openjiuwen[sqlite,postgres,mysql,gaussdb,redis,elasticsearch,obs]]这正体现了项目的最小核心、可选扩展设计哲学核心模块不因某个特定数据库而膨胀特定数据库支持以可选扩展方式按需安装。二、设计哲学四条核心原则GaussDbStore的构建遵循四条原则这四条原则也解释了本模块几乎所有实现细节最小依赖Minimal Dependencies只依赖标准的 SQLAlchemy 接口避免与特定数据库强耦合隔离性Isolation不污染全局命名空间保证与其他数据库驱动兼容共存可扩展性Extensibility基于 SQLAlchemy 方言机制实现易于维护与升级标准合规Standard Compliance遵循 Python 与 SQLAlchemy 的最佳实践。后续的import_dbapi()替代 monkey patch、继承BaseDbStore而非DefaultDbStore等决策都可以在这四条原则上找到依据。三、方言实现为什么用import_dbapi()而不是 Monkey Patching3.1 旧方案Monkey Patching的缺陷文档中给出了旧实现的示意代码通过把sys.modules[asyncpg]指向async_gaussdb来偷梁换柱# Old implementation: Global pollution of sys.modules def _apply_monkey_patch(): import async_gaussdb sys.modules[asyncpg] async_gaussdb # Global pollution sys.modules[asyncpg.exceptions] async_gaussdb.exceptions这种做法的问题非常明显❌全局污染改写sys.modules会影响进程内后续所有代码❌不可预期其他使用真实asyncpg的代码会被连带影响❌不合规违背 SQLAlchemy 官方推荐的方言实现方式❌难调试隐式修改全局状态出现问题难以追踪。3.2 新方案import_dbapi()局部加载SQLAlchemy 官方为方言提供了import_dbapi()这个 classmethod 扩展点方言在真正需要底层驱动时才调用它。新实现正是利用这一机制# New implementation: Local driver loading class GaussDialect_asyncpg(PGDialect_asyncpg): classmethod def import_dbapi(cls): import async_gaussdb from sqlalchemy.dialects.postgresql.asyncpg import AsyncAdapt_asyncpg_dbapi return AsyncAdapt_asyncpg_dbapi(async_gaussdb)好处✅无全局污染不修改sys.modules完全隔离✅标准合规使用 SQLAlchemy 官方提供的扩展点✅易维护代码结构清晰职责明确✅类型安全通过AsyncAdapt_asyncpg_dbapi包装器保证类型兼容。3.3 源码中的实际实现仓库中 gauss_dialect.py 的真实实现比文档示意代码更完整还包含了依赖缺失时的友好报错classmethod def import_dbapi(cls): try: import async_gaussdb patched_driver _patch_gaussdb_driver(async_gaussdb) logger.info([GaussDialect] Loaded async_gaussdb via import_dbapi) return AsyncAdapt_asyncpg_dbapi(patched_driver) except ImportError as import_error: raise ImportError( Please install async-gaussdb to use the gaussdb dialect. Run: pip install openjiuwen[gaussdb] or pip install openjiuwen[all-storage] ) from import_error如果直接import async_gaussdb失败异常信息会明确提示安装方式而不是抛出一个难以理解的底层错误。3.4 方言类的完整配置GaussDialectAsyncpg继承自PGDialect_asyncpg源码中定义了大量方言级属性gauss_dialect.pyclass GaussDialectAsyncpg(PGDialect_asyncpg): statement_compiler GaussCompiler name gaussdb driver async_gaussdb supports_statement_cache True supports_native_enum False supports_native_uuid False use_insertmanyvalues False colspecs { **PGDialect_asyncpg.colspecs, String: GaussString, }这些配置各有讲究属性值含义namegaussdb方言名称对应连接串 URL schemedriverasync_gaussdb驱动名称与name一起构成gaussdbasync_gaussdb://supports_statement_cacheTrue启用 SQLAlchemy 语句缓存提升重复语句执行性能supports_native_enum/supports_native_uuidFalse不依赖 GaussDB 原生枚举与 UUID 类型退回通用实现use_insertmanyvaluesFalse不使用批量 insertmanyvalues 特性规避 GaussDB 兼容性问题colspecs扩展String→GaussString自定义类型映射见下文此外还重写了三个反射/元数据相关方法_get_server_version_info固定返回(9, 2)绕过版本探测可能引发的兼容问题_domain_query与_enum_query返回SELECT 1 WHERE FALSE告知 SQLAlchemy GaussDB 没有 domain/enum 类型可查询避免建表或反射时报错。3.5 方言注册与模块缓存在 gauss_dialect.py 末尾模块导入时自动执行注册registry.register(gaussdb.async_gaussdb, __name__, GaussDialectAsyncpg) registry.register(gaussdb, __name__, GaussDialectAsyncpg)同时注册两个名字保证gaussdbasync_gaussdb://与gaussdb://两种连接串都能被解析。Python 的模块缓存机制保证该注册逻辑在整个进程生命周期中只执行一次——只要任意代码import openjiuwen.extensions.store.db.gauss_dialect或导入GaussDbStore间接触发后续所有create_async_engine(gaussdb...)都会命中已注册的方言无需手动调用任何初始化函数。值得注意的容错处理整个导入过程包在try/except ImportError中gauss_dialect.py如果环境缺少 SQLAlchemy 相关依赖只记录 warning 而不会让导入链崩溃。四、驱动补丁机制让 async_gaussdb 符合 DB-API 2.0async_gaussdb驱动本身缺少部分 DB-API 2.0 标准属性SQLAlchemy 方言加载时无法直接识别因此 gauss_dialect.py 提供了_patch_gaussdb_driverdef _patch_gaussdb_driver(driver_module): if not hasattr(driver_module, paramstyle): driver_module.paramstyle format if not hasattr(driver_module, Error): driver_module.Error getattr(driver_module, GaussDBError, Exception) if not hasattr(driver_module, apilevel): driver_module.apilevel 2.0 if not hasattr(driver_module, threadsafety): driver_module.threadsafety 0 return driver_module各补丁项的语义属性补丁值说明paramstyleformat参数占位符风格%s形式SQLAlchemy 据此格式化 SQL 参数ErrorGaussDBError缺省回退ExceptionDB-API 规定的异常基类供上层统一捕获apilevel2.0声明驱动遵循 DB-API 2.0 规范threadsafety0标记驱动非线程安全SQLAlchemy 据此调整连接使用策略所有补丁都先判断hasattr再赋值幂等且不覆盖驱动自身已有的属性。补丁完成后驱动再被AsyncAdapt_asyncpg_dbapi包装from sqlalchemy.dialects.postgresql.asyncpg import AsyncAdapt_asyncpg_dbapi return AsyncAdapt_asyncpg_dbapi(patched_driver)AsyncAdapt_asyncpg_dbapi是 SQLAlchemy 提供的异步连接适配层负责提供await_()等方法确保create_async_engine创建的引擎在 async 上下文中正常工作。五、源码中的三项额外兼容性补丁除文档明示的驱动补丁外gauss_dialect.py 还包含三项针对 GaussDB 行为差异的定制值得实际使用者了解1. 反射 SQL 修复before_cursor_execute 事件GaussDB 的pg_type.typcollation子查询可能返回 NULL 导致元数据反射失败因此在每次游标执行前拦截并替换gauss_dialect.pyevent.listens_for(Engine, before_cursor_execute, retvalTrue) def _patch_gaussdb_reflection_sql(*args): conn, _cursor, statement, parameters, _context, _executemany args if getattr(conn.dialect, name, ) gaussdb: if pg_type.typcollation in statement: statement re.sub( r\(\s*SELECT\s[^)]*?pg_type\.typcollation[^)]*?\), NULL, statement, flagsre.IGNORECASE | re.DOTALL ) return statement, parameters事件处理器只在dialect.name gaussdb时生效对 PostgreSQL 等其他方言零影响符合隔离性原则。2.GaussString类型入参强制字符串化确保所有传入String列的数据在进入驱动前被转为字符串其中对datetime特别处理为%Y-%m-%d %H:%M:%S.%f格式gauss_dialect.py。这在 GaussDB 对某些类型绑定比较严格时能避免隐式转换错误。3.GaussCompilerFOR UPDATE 子句兼容PGDialect_asyncpg默认的行锁语法可能带 GaussDB 不接受的子句因此用子类简化输出gauss_dialect.pyclass GaussCompiler(PGCompiler): def for_update_clause(self, select, **kw): return FOR UPDATE六、为什么继承BaseDbStore而不是DefaultDbStore文档明确指出虽然GaussDbStore与DefaultDbStore当前实现几乎相同对比 gauss_db_store.py 与 default_db_store.py 可见两者都只是持有async_conn并返回引擎但继承BaseDbStore是深思熟虑的设计语义清晰BaseDbStore是抽象基类定义统一存储接口base_db_store.py 中get_async_engine()是唯一抽象方法GaussDbStore由此明确标识自己是 GaussDB 专用实现便于在代码中区分不同数据库存储。模块隔离DefaultDbStore位于核心模块openjiuwen.core.foundation.store.dbGaussDbStore位于扩展模块openjiuwen.extensions.store.db核心模块保持干净具体数据库实现放在扩展侧。未来可扩展GaussDbStore未来可独立增加 GaussDB 专属能力如连接池配置、监控指标、性能优化不影响核心模块稳定性。依赖管理DefaultDbStore随核心模块分发GaussDbStore是可选扩展通过pip install openjiuwen[gaussdb]或pip install openjiuwen[all-storage]安装与最小核心、可选扩展哲学一致。从协议层面看LongTermMemory.register_store要求db_store是BaseDbStore实例见 long_term_memory.py 的 isinstance 校验因此任何继承BaseDbStore的存储都可以无缝替换这也是插件化设计的直接体现。七、类 API 详解7.1GaussDbStoreclass openjiuwen.extensions.store.db.gauss_db_store.GaussDbStore(async_conn: AsyncEngine)GaussDB 数据库存储实现继承自BaseDbStore。它包装一个 SQLAlchemyAsyncEngine上层业务通过get_async_engine()获取引擎后即可使用标准 SQLAlchemy async 接口执行 CRUD。导入该模块会自动触发gauss_dialect注册通过import_dbapi()局部加载驱动无需手动调用。参数async_connAsyncEngineSQLAlchemy 异步引擎实例通常由create_async_engine创建。示例 from sqlalchemy.ext.asyncio import create_async_engine from openjiuwen.extensions.store.db import GaussDbStore engine create_async_engine(gaussdbasync_gaussdb://user:passwordhost:port/database) store GaussDbStore(async_connengine) # Use store.get_async_engine() to get the engine for database operations7.2get_async_engineget_async_engine() - AsyncEngine返回内部持有的AsyncEngine实例。多次调用返回同一个实例单例语义可通过is断言验证 engine store.get_async_engine() assert engine is store.get_async_engine() # Same instance该方法是BaseDbStore定义的抽象接口实现base_db_store.py保证所有数据库存储对上层暴露一致的引擎获取入口。7.3async_connasync_conn: AsyncEngine内部持有的AsyncEngine实例属性由构造时传入并保存gauss_db_store.py。八、调用链与初始化时序文档给出了清晰的调用链理解它有助于排查初始化问题from openjiuwen.extensions.store.db import GaussDbStore │ ▼ gauss_db_store.py → import gauss_dialect (module load, executed once only) │ ▼ gauss_dialect.py module load automatically completes: ├── registry.register(gaussdb, ...) # Register GaussDialect_asyncpg with SQLAlchemy └── GaussDialect_asyncpg.import_dbapi() # Load async_gaussdb driver locally │ ▼ GaussDbStore(async_connengine) # Create store instance │ ▼ store.get_async_engine() # Get AsyncEngine → Execute database operations via SQLAlchemy两个关键语义值得强调Python 模块缓存机制保证gauss_dialect.py中的初始化逻辑含方言注册只执行一次无需手动调用新实现通过import_dbapi()局部加载驱动不污染全局sys.modules。九、完整 CRUD 实战示例以下示例完整演示了从建引擎、建表到增删改查的全过程可直接运行替换连接串为真实 GaussDB 地址import asyncio from datetime import datetime from sqlalchemy import String, Integer, Float, DateTime, Text, select, update, delete from sqlalchemy.ext.asyncio import create_async_engine, AsyncSession from sqlalchemy.orm import DeclarativeBase, sessionmaker, Mapped, mapped_column from openjiuwen.extensions.store.db import GaussDbStore class Base(DeclarativeBase): pass class Product(Base): __tablename__ products id: Mapped[int] mapped_column(primary_keyTrue, autoincrementTrue) name: Mapped[str] mapped_column(String(100), nullableFalse) price: Mapped[float] mapped_column(Float, nullableFalse) category: Mapped[str] mapped_column(String(50)) stock: Mapped[int] mapped_column(Integer, default0) created_at: Mapped[datetime] mapped_column(DateTime, defaultdatetime.now) async def main(): GAUSSDB_URL gaussdbasync_gaussdb://user:passwordhost:port/database engine create_async_engine(GAUSSDB_URL, echoFalse) store GaussDbStore(async_connengine) engine store.get_async_engine() async_session sessionmaker(engine, class_AsyncSession, expire_on_commitFalse) # Create tables async with engine.begin() as conn: await conn.run_sync(Base.metadata.create_all) # Insert async with async_session() as session: product Product(nameSmartphone, price6999.00, categoryElectronics, stock50) session.add(product) await session.commit() await session.refresh(product) print(fInserted: ID{product.id}) # Query async with async_session() as session: result await session.execute(select(Product).where(Product.name Smartphone)) p result.scalar_one() print(fQuery result: {p.name}, Price{p.price}) # Update async with async_session() as session: stmt update(Product).where(Product.name Smartphone).values(price5999.00) await session.execute(stmt) await session.commit() # Delete async with async_session() as session: stmt delete(Product).where(Product.name Smartphone) await session.execute(stmt) await session.commit() await engine.dispose() if __name__ __main__: asyncio.run(main())示例要点store.get_async_engine()返回的引擎可直接交给sessionmaker使用说明GaussDbStore与标准 SQLAlchemy async ORM 工作流完全兼容echoFalse关闭 SQL 日志调试时可改为True结尾await engine.dispose()显式释放连接池资源。十、与 LongTermMemory 内存引擎的集成GaussDbStore可作为LongTermMemory的db_store参数替换默认的DefaultDbStore。这正是它在 agent-core 中最典型的应用场景——把记忆持久化到 GaussDB from openjiuwen.core.memory import LongTermMemory from openjiuwen.extensions.store.db import GaussDbStore from sqlalchemy.ext.asyncio import create_async_engine engine create_async_engine(gaussdbasync_gaussdb://user:passwordhost:port/agentmgr) db_store GaussDbStore(async_connengine) memory LongTermMemory() await memory.register_store( ... kv_storekv_store, ... vector_storevector_store, ... db_storedb_store, ... )源码层面的集成链路long_term_memory.py如下register_store首先做类型校验db_store必须是BaseDbStore实例否则抛出MEMORY_REGISTER_STORE_EXECUTION_ERRORlong_term_memory.py若提供了db_store自动调用create_tables初始化持久化表结构long_term_memory.py若未显式提供message_store则会用SqlDbStore(db_store)自动构造内部SqlMessageStorelong_term_memory.py。SqlDbStore通过db_store.get_async_engine()获取引擎绑定AsyncSession见 sql_db_store.py因此上层记忆读写最终都汇聚到GaussDbStore持有的引擎上。这意味着你只需替换这一个db_store整个记忆系统的 SQL 持久化层就整体迁移到了 GaussDB。十一、连接字符串格式gaussdbasync_gaussdb://username:passwordhost:port/database格式解析方言名驱动名://用户名:密码主机:端口/数据库名。其中gaussdb方言名对应GaussDialectAsyncpg.nameasync_gaussdb驱动名对应GaussDialectAsyncpg.driver即async-gaussdb包该连接串能够被正确解析的前提是gauss_dialect.py已完成注册导入GaussDbStore时自动完成。十二、依赖说明与安装方式12.1 依赖清单包版本要求说明sqlalchemy 2.0.41项目主依赖自带async-gaussdb~ 0.30.0GaussDB 异步驱动兼容 0.30.x可选依赖asyncpg 0.29.0PostgreSQL 异步驱动SQLAlchemyPGDialect_asyncpg的错误翻译_asyncpg_error_translate与异步连接适配在运行时必需12.2 为什么asyncpg是必需依赖虽然GaussDbStore主驱动是async_gaussdb但asyncpg仍是运行时硬依赖原因有三SQLAlchemy asyncpg 方言基础设施GaussDialectAsyncpg继承PGDialect_asyncpg并复用其适配层AsyncAdapt_asyncpg_dbapi。SQLAlchemy asyncpg 方言的错误翻译机制_asyncpg_error_translate在数据库报错时会直接import asyncpg并引用asyncpg.exceptions.*运行时错误处理查询抛错时AsyncAdapt_asyncpg_connection._handle_exception会访问_asyncpg_error_translate触发import asyncpg。若未安装asyncpg错误处理过程本身会抛出ImportError导致真正的数据库错误信息被掩盖连接适配AsyncAdapt_asyncpg_dbapi提供的异步连接适配层内部依赖 asyncpg 兼容接口。这也解释了 pyproject.toml 中gaussdbextra 为何同时包含async-gaussdb与asyncpg两个包。12.3 安装方式# Option 1: Install GaussDB optional dependency (includes asyncpg) pip install openjiuwen[gaussdb] # Option 2: Install all storage dependencies (recommended) pip install openjiuwen[all-storage] # Option 3: Install drivers separately pip install async-gaussdb asyncpg仅使用 GaussDB 存储时Option 1 最精简需要在多种存储后端间切换时Option 2 一步到位手动管理依赖的定制环境可用 Option 3。十三、使用限制与兼容性边界13.1 支持的能力✅标准 SQLAlchemy 操作基础 CRUD 操作增删改查事务管理连接池管理异步查询。✅PostgreSQL 兼容语法GaussDialect_asyncpg继承自PGDialect_asyncpg支持大部分 PostgreSQL 数据类型与函数适用于标准 SQL 查询与事务操作。13.2 不支持的能力❌GaussDB 专属高级语法不支持 GaussDB 专属存储过程与函数调用不支持 GaussDB 专属数据类型扩展不支持 GaussDB 专有性能优化特性如特定索引策略。❌asyncpg 专属特性由于使用async_gaussdb驱动不支持 asyncpg 自定义类型编码/解码器不支持 asyncpg 的通知/监听机制。❌高级连接特性不支持服务端游标的某些高级用法不支持 PostgreSQL 专属复制特性不支持 PostgreSQL 的 LISTEN/NOTIFY 机制。13.3 兼容性说明完全兼容场景标准 ORM 操作SQLAlchemy Core 与 ORM基础数据类型INTEGER、VARCHAR、TEXT、TIMESTAMP 等标准事务与连接管理简单查询与更新操作。需要适配的场景使用 PostgreSQL 专属函数的复杂查询需要数据库专属优化的场景使用高级并发控制如高级行级锁。十四、最佳实践与优雅降级14.1 使用标准 SQLAlchemy 接口# Recommended: Use standard ORM operations async with async_session() as session: user User(nameAlice) session.add(user) await session.commit()14.2 避免数据库专属语法# Avoid: Using PostgreSQL-specific functions # stmt select(User).where(User.created_at func.now()) # Recommended: Use standard SQL functions from sqlalchemy import func stmt select(User).where(User.created_at func.current_timestamp())14.3 处理可选依赖降级策略由于GaussDbStore依赖可选安装的驱动生产代码应做好 ImportError 降级缺依赖时回退到核心模块的DefaultDbStoretry: from openjiuwen.extensions.store.db import GaussDbStore except ImportError: # Fallback to default storage from openjiuwen.core.foundation.store.db import DefaultDbStore该模式与项目的可选扩展设计呼应核心功能不因扩展缺失而不可用。14.4 三条实践小结保持标准接口尽量只使用跨数据库通用的 ORM 能力可移植性最好关注反射行为建表/反射依赖元数据查询若遇异常可确认是否触发pg_type.typcollation相关的自动修复逻辑按需安装仅 GaussDB 场景用openjiuwen[gaussdb]多存储场景用openjiuwen[all-storage]并始终确保asyncpg同时就位否则运行时报错时错误翻译链路会失效。总结GaussDbStore是 openJiuwen agent-core 最小核心、可选扩展架构的一个典型样本核心只暴露BaseDbStore抽象接口GaussDB 支持完全收敛在openjiuwen.extensions.store.db扩展模块中。通过import_dbapi()局部加载驱动、注册gaussdb方言、并以AsyncAdapt_asyncpg_dbapi做异步适配它既避免了旧式 monkey patch 的全局污染又让上层业务可以用完全标准的 SQLAlchemy async 接口读写 GaussDB。配合LongTermMemory.register_store一行db_store参数即可把整个记忆持久化层切换到 GaussDB——这正是该扩展在 agent-core 中的核心价值所在。如需进一步研究建议继续阅读gauss_dialect.py、gauss_db_store.py、base_db_store.py、long_term_memory.py以及 pyproject.toml 中gaussdb与all-storage的依赖声明。赞分享人工智能AI AgentAgent 框架大模型工具调用RAG提示工程强化学习【免费下载链接】agent-coreopenJiuwen agent-core可提供AI Agent开发、运行、调优与演进相关的全套SDK能力项目地址https://gitcode.com/openJiuwen/agent-core点击查看免费下载相关推荐openJiuwen agent-core 对象存储实践AioBotoClient 异步 S3/OBS 客户端全解析openJiuwen agent core 对象存储实践AioBotoClient 异步 S3/OBS 客户端全解析 本指南围绕 openJiuwen age人工智能AI AgentAgent 框架大模型工具调用RAG提示工程强化学习openJiuwen agent-core 存储层完全指南BaseKVStore / BaseDbStore / BaseVectorStore 抽象与内置实现深度解析openJiuwen agent core 存储层完全指南BaseKVStore / BaseDbStore / BaseVectorStore 抽象与内置实人工智能AI AgentAgent 框架大模型工具调用RAG提示工程强化学习openJiuwen agent-core 对象存储客户端BaseObjectStorageClient 抽象接口与 AioBotoClient 异步 S3 实现详解openJiuwen agent core 对象存储客户端BaseObjectStorageClient 抽象接口与 AioBotoClient 异步 S3人工智能AI AgentAgent 框架大模型工具调用RAG提示工程强化学习上一篇Pwndbg与SystemTap调试示例常用调试场景下一篇gh_mirrors/ma/materials深度探索从基础语法到高级框架的完整学习路径创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表