简体中文 English 繁體中文 日本語 Русский
本文为静态镜像,内容以交互版为准 在交互式文档中心打开 →

ErisPulse.Core.Bases.sql_base 模块


模块概述

ErisPulse SQL 存储共享基类

提供 SQLite / MySQL / PostgreSQL 三种内置后端共享的通用 SQL 逻辑: 嵌套键 KV 存取、事务编排、DDL、查询构建与 ALTER TABLE。 方言差异(占位符、标识符引号、UPSERT、类型翻译等)收敛到 :class:SQLDialect 及其子类。

提示

  1. 新增 SQL 后端只需:定义 Dialect 子类 + 继承 SQLStorageBase 实现连接管理与执行漏斗
  2. KV 值统一以 JSON 文本存储,三种后端行为完全一致

函数列表

_validate_identifier(name: str, context: str = 'identifier')

内部方法 验证 SQL 标识符(表名/列名)是否安全


_validate_select_column(name: str, context: str = 'column')

内部方法 验证 SELECT/ORDER BY 列表达式是否安全

采用黑名单模式:仅拦截 SQL 注入危险字符(; ' " -- /* */ 换行), 允许任意合法 SQL 列表达式,包括简单列名、聚合函数、别名、表达式等。


_validate_column_type(col_type: str)

内部方法 验证列类型定义是否安全(防止通过类型定义注入 SQL)


类列表

class _SingletonMixin

内部方法 进程级单例混入(__new__ 级去重,可 __new__ 绕过用于测试)

class SQLDialect

SQL 方言基类

收敛三种内置后端的方言差异:占位符翻译、标识符引号、UPSERT 语法、 列类型翻译(含 AUTOINCREMENT 自增主键)、表存在性查询等。

提示

  1. 构建期统一使用 ? 占位符记号,由 :meth:translate_sql 翻译为方言样式
  2. 标识符一律经 :meth:quote 引用(列名/表名均先经合法性校验)

方法列表

translate_sql(sql: str)

将内部 ? 占位符记号翻译为方言占位符


quote(name: str)

引用 SQL 标识符


kv_table_ddl(table: str, key_col: str, value_col: str)

生成 KV 表建表 DDL


upsert_kv_sql(table: str, key_col: str, value_col: str)

生成 KV UPSERT 语句(内部 ? 占位符)


has_table_sql(table_name: str)

生成表存在性查询


create_index_sql(table: str, index: str, column: str)

生成普通索引创建 DDL

基类语义为"不存在则创建"(IF NOT EXISTS);不支持该语法的方言 (MySQL)覆写本方法为裸 CREATE INDEX,幂等由 :meth:has_index_sql 存在性预检承接。


has_index_sql(table: str, index: str)

生成索引存在性查询(内部 ? 占位符)


autoincrement_column(base_type: str)

翻译自增主键定义("INTEGER PRIMARY KEY AUTOINCREMENT" 中的类型部分)


last_insert_id_sql()

生成"最近一次插入的自增主键"查询 SQL

供 ORM(Core/Bases/model.py)在 INSERT 后同一连接上回填自增主键。 须与 INSERT 在同一连接/事务内执行(MySQL LAST_INSERT_ID() 与 PostgreSQL lastval() 均为连接级语义)。

返回值 (查询): SQL(内部 ? 占位符记号,无参数)


_map_first_word(word: str)

内部方法 按首词映射列类型,保留括号参数与后缀


translate_column_type(col_type: str)

翻译列类型定义为方言等价写法

处理自增主键重排与首词类型映射;无法识别的类型原样保留。


is_missing_table_error(exc: BaseException)

判断异常是否为"KV 表不存在"(用于自动重建后重试)


class SQLQueryBuilder(BaseQueryBuilder)

SQL 查询构建器(方言无关的共享实现)

链式构建标准 SQL(内部 ? 占位符),终止方法通过所属后端的 :meth:SQLStorageBase._execute_query 漏斗执行,由方言完成占位符翻译。

提示 使用方式:

  1. await storage.Table("users").Insert({"name": "Alice"}).aExecute()
  2. storage.Table("users").Select("name").Where("age > ?", 18).Execute()

方法列表

_rows_to_dict(rows: list[tuple], columns: list[str] | None)

内部方法 将 tuple 行列表转为字典列表(columns 不可用时原样返回)


async aExecute()

执行构建的查询

>>> rows = await storage.Table("users").Select("name", "age").aExecute()
>>> affected = await storage.Table("users").Delete().Where("age < ?", 18).aExecute()

async _aexecute_insert_multi()

内部方法 批量插入执行


async aExecuteOne()

执行查询并返回单条结果

示例:

>>> row = await storage.Table("users").Select("*").Where("id = ?", 1).aExecuteOne()

async aCount()

执行 COUNT 查询

>>> total = await storage.Table("users").Where("age > ?", 18).aCount()

async aExists()

检查是否存在匹配的记录

>>> if await storage.Table("users").Where("name = ?", "Alice").aExists():

class AlterTableBuilder

ALTER TABLE 构建器

链式收集表结构修改操作,aExecute() 原生异步执行, Execute() 同步兼容桥接。

提示 使用方式:

  1. await storage.AlterTable("users").AddColumn("email", "TEXT").aExecute()
  2. storage.AlterTable("users").RenameTo("members").Execute()

方法列表

AddColumn(column_name: str, column_type: str)

添加列

>>> storage.AlterTable("users").AddColumn("email", "TEXT").Execute()

RenameTo(new_name: str)

重命名表

>>> storage.AlterTable("users").RenameTo("members").Execute()

async aExecute()

执行所有已收集的 ALTER TABLE 操作

返回值: 操作是否成功


Execute()

执行所有已收集的 ALTER TABLE 操作(同步兼容,桥接到 :meth:aExecute)

返回值: 操作是否成功


class SQLStorageBase(BaseStorage)

SQL 存储后端共享基类

实现三种内置 SQL 后端共享的 KV 存取(含嵌套键)、批量操作、DDL、 事务编排与 ALTER TABLE;子类只需提供连接管理与方言执行漏斗。

提示 子类必须实现:

  1. _create_loop_resource / _destroy_loop_resource:每事件循环资源(连接池)
  2. _acquire_resource_conn / _release_resource_conn:非事务连接获取/归还
  3. _open_txn_conn / _close_txn_conn:事务专用连接获取/释放
  4. _exec_query_on:方言执行漏斗(占位符翻译 + 游标语义)

方法列表

_ensure_resource_state()

内部方法 确保实例级资源表存在(new 绕过构造时兜底)


_is_ready()

内部方法 检查存储后端是否已初始化完成


_finish_init()

内部方法 子类完成后调用:标记就绪并订阅配置热更新告警


_watch_storage_config(check: Callable[[], str | None])

内部方法 订阅存储配置热更新;check() 返回告警文案(None 表示无需告警)


_default_project_path(filename: str)

内部方法 项目目录下的默认数据库文件路径


async _create_loop_resource()

内部方法 为当前事件循环创建连接资源(连接池/共享连接)


async _destroy_loop_resource(resource: Any)

内部方法 销毁事件循环连接资源


async _acquire_resource_conn(resource: Any)

内部方法 从资源获取非事务连接


async _release_resource_conn(resource: Any, conn: Any)

内部方法 归还非事务连接


async _open_txn_conn()

内部方法 获取事务专用连接


async _close_txn_conn(conn: Any)

内部方法 释放事务专用连接


async _get_loop_resource()

内部方法 获取当前事件循环的连接资源(惰性创建,幂等,瞬时失败自动重试)


async _emit_storage_event(event: str)

内部方法 发出存储生命周期事件(storage.ready / storage.unreachable / storage.recovered),供外部感知连接状态变化(如 Dashboard 告警)。

后台发射(fire):存储操作绝不等待观测者;lifecycle 未就绪或 处理异常时静默跳过,不影响存储操作本身。


async _acquire()

内部方法 获取连接的统一入口:优先异步事务上下文绑定连接,否则走事件循环资源


async _run_with_conn(impl: Callable[..., Any])

内部方法 执行 impl(conn, *args):显式连接优先,否则自动获取并在退出时归还


async _exec_query_on(kind: str, sql: str, params: Any, conn: Any)

内部方法 方言执行漏斗:翻译占位符后按 kind 执行


async _execute_query(kind: str, sql: str, params: Any)

内部方法 查询执行统一入口(连接路由 → 执行漏斗)


async _acquire_txn_conn()

内部方法 获取事务专用连接


async _begin_txn(conn: Any)

内部方法 开启事务(SQL BEGIN,方言可覆写)


async _commit_txn(conn: Any, handle: Any = None)

内部方法 提交事务


async _rollback_txn(conn: Any, handle: Any = None)

内部方法 回滚事务


async _release_txn_conn(conn: Any)

内部方法 释放事务专用连接


async _run_alter(table_name: str, operations: 'list[tuple[str, tuple[Any, ...]]]')

内部方法 顺序执行 ALTER TABLE 操作


_parse_nested_key(key: str)

内部方法 解析嵌套键:点号(.)总是表示嵌套访问,即使根键不存在也会创建嵌套结构


_get_nested_value(obj: Any, key_path: list[str])

内部方法 从嵌套对象中获取值


_set_nested_value(obj: Any, key_path: list[str], value: Any)

内部方法 在嵌套对象中设置值(中间层一律预创建为字典)


_delete_nested_value(obj: Any, key_path: list[str])

内部方法 从嵌套对象中删除值,返回 (更新后的对象, 是否删除成功)


_loads(raw: Any)

内部方法 JSON 反序列化(失败时返回原始文本)


_kv_table_sql()

内部方法 返回 (已引用表名, 已引用键列, 已引用值列)


async _init_kv_table()

内部方法 创建默认 KV 表


async _recover_missing_table()

内部方法 缺表自动恢复


_shadow_overlay()

内部方法 影子存储覆盖层(方向十一):当前 owner 属影子时返回其覆盖层, 否则 None。KV 三件套(aget / aset / adelete)据此实现"写隔离、 读透传、删为墓碑";ORM 读写不在覆盖层语义内(文档声明)。


async aget(key: str, default: Any = None)

异步获取存储项的值

支持嵌套键访问,如 "user.settings.theme" 会从存储的嵌套对象中获取值

>>> timeout = await storage.aget("network.timeout", 30)

async aset(key: str, value: Any)

异步设置存储项的值

支持嵌套键设置,如 "user.settings.theme" 会更新存储的嵌套对象中的对应字段

- **key** (`存储项键名,支持嵌套路径`): - **value**: 存储项的值
- **conn** (`内部参数:事务连接路由(勿手动传入)`): **返回值** (`操作是否成功`): 

示例:

>>> await storage.aset("user.settings.theme", "dark")

async adelete(key: str)

异步删除存储项

支持嵌套键删除,如 "user.settings.theme" 会删除嵌套对象中的对应字段

- **key** (`存储项键名,支持嵌套路径`): - **conn**: 内部参数:事务连接路由(勿手动传入)
**返回值** (`操作是否成功`): 

示例:

>>> await storage.adelete("user.settings.theme")

async aget_all_keys()

异步获取所有存储项的键名

>>> all_keys = await storage.aget_all_keys()

async aclear()

异步清空所有存储项

>>> await storage.aclear()

async aget_multi(keys: list[str])

异步批量获取多个存储项的值(单次 IN 查询,仅返回存在的键)

>>> settings = await storage.aget_multi(["app.name", "app.version"])

get_multi(keys: list[str])

批量获取多个存储项的值(同步兼容,桥接到 :meth:aget_multi)

>>> settings = storage.get_multi(["app.name", "app.version"])

async aset_multi(items: dict[str, Any])

异步批量设置多个存储项

与旧版语义一致:键按字面量直写,不做点号嵌套解析 (需要嵌套行为请逐键调用 :meth:aset)。

>>> await storage.aset_multi({"app.name": "MyApp", "app.debug": True})

set_multi(items: dict[str, Any])

批量设置多个存储项(同步兼容,桥接到 :meth:aset_multi)

>>> storage.set_multi({"app.name": "MyApp", "app.debug": True})

async adelete_multi(keys: list[str])

异步批量删除多个存储项(单次批量 DELETE)

>>> await storage.adelete_multi(["temp.key1", "temp.key2"])

delete_multi(keys: list[str])

批量删除多个存储项(同步兼容,桥接到 :meth:adelete_multi)

>>> storage.delete_multi(["temp.key1", "temp.key2"])

Table(table_name: str)

获取指定表的查询构建器

示例:

>>> rows = await storage.Table("users").Select("name", "age").Where("age > ?", 18).aExecute()

async aCreateTable(table_name: str, columns: dict[str, str])

异步创建表

列类型支持 SQLite 风格定义(如 "INTEGER PRIMARY KEY AUTOINCREMENT"), 由方言自动翻译为目标后端等价写法。

>>> await storage.aCreateTable("users", {
...     "id": "INTEGER PRIMARY KEY AUTOINCREMENT",
...     "name": "TEXT NOT NULL"
... })

async aDropTable(table_name: str)

异步删除表

>>> await storage.aDropTable("users")

async aHasTable(table_name: str)

异步检查表是否存在

>>> if await storage.aHasTable("users"):

async aGetTableColumns(table_name: str)

异步列举表的现有列名(ORM 自动迁移用)


AlterTable(table_name: str)

获取 ALTER TABLE 构建器

示例:

>>> storage.AlterTable("users").AddColumn("email", "TEXT").Execute()

async aclose()

异步关闭当前事件循环上绑定的连接资源(连接池/共享连接)

事务专用连接不受影响(由事务自行管理)。

示例:

>>> await storage.aclose()