数据库 hscredit.database
统一 SQL/NoSQL 连接池、参数化查询、流式读取、自动建表、分批写入、Redis/MongoDB CRUD、元数据导出与适配器扩展 API。 使用流程和后端能力矩阵见 数据库与 NoSQL 连接池、读写及表结构导出。
数据库门面
- class hscredit.database.client.Database(database_type, *, pool_options=None, adapter_options=None, **connect_kwargs)[源代码]
基类:
object统一数据库连接与操作门面。
参数
- database_typestr
已注册数据库类型或别名。
- pool_optionsmapping or PoolOptions, optional
连接池配置。
- adapter_optionsmapping, optional
后端专有配置。
**connect_kwargs直接传递给后端驱动的连接参数。
属性
- adapter
当前数据库适配器实例。
参考样例
>>> with Database("mysql", host="127.0.0.1", database="risk") as db: ... rows = db.query("SELECT 1", result="rows")
- 参数:
database_type (str)
pool_options (Mapping[str, Any] | None)
adapter_options (Mapping[str, Any] | None)
connect_kwargs (Any)
- query(sql, params=None, result='dataframe')[源代码]
一次性执行查询并返回全部结果。
大结果集建议改用
stream_query()或read_query(),避免一次性占用过多内存。- 参数:
sql (str) -- 由当前数据库执行的查询 SQL。
params (Any) -- 按底层驱动参数风格绑定的 SQL 参数,默认为
None。result (str) -- 结果类型,支持
dataframe、records和rows。
- 返回:
dataframe返回 DataFrame;records返回记录字典列表;rows返回原始行列表。- 抛出:
StateError -- 数据库门面已经关闭。
ValidationError --
result不是支持的结果类型。DatabaseQueryError -- 数据库执行查询失败。
- 返回类型:
Any
参考样例
>>> frame = db.query("SELECT id FROM events WHERE id > %s", params=(10,)) >>> records = db.query("SELECT id FROM events", result="records")
- execute(sql, params=None)[源代码]
执行单条 DDL 或 DML SQL。
- 参数:
sql (str) -- 要执行的 SQL。
params (Any) -- 按底层驱动参数风格绑定的 SQL 参数,默认为
None。
- 返回:
适配器报告的影响行数或原生执行结果。
- 抛出:
StateError -- 数据库门面已经关闭。
DatabaseQueryError -- SQL 执行失败。
- 返回类型:
Any
- executemany(sql, values)[源代码]
使用多组绑定值批量执行同一条 SQL。
- 参数:
sql (str) -- 带驱动占位符的 SQL。
values (Any) -- 每行一组参数的可迭代对象。
- 返回:
适配器报告的累计影响行数或原生执行结果。
- 抛出:
StateError -- 数据库门面已经关闭。
DatabaseQueryError -- 批量执行失败。
- 返回类型:
Any
- property native_client: Any
返回 Redis 或 MongoDB 适配器持有的原生客户端。
- read_one(resource, selector=None, **options)[源代码]
读取单个 Redis key 或 MongoDB 文档。
- 参数:
resource (Any)
selector (Any)
options (Any)
- 返回类型:
Any
- read_many(resource, selector=None, **options)[源代码]
批量读取 Redis keys 或 MongoDB 文档。
- 参数:
resource (Any)
selector (Any)
options (Any)
- 返回类型:
Any
- read(resource, selector=None, **options)[源代码]
根据输入形态自适应执行单条或批量读取。
- 参数:
resource (Any)
selector (Any)
options (Any)
- 返回类型:
Any
- write_one(resource, data, **options)[源代码]
写入单个 Redis key 或 MongoDB 文档。
- 参数:
resource (Any)
data (Any)
options (Any)
- 返回类型:
Any
- write_many(resource, data=None, **options)[源代码]
批量写入 Redis key-value 或 MongoDB 文档。
- 参数:
resource (Any)
data (Any)
options (Any)
- 返回类型:
Any
- write(resource, data=None, **options)[源代码]
根据输入形态自适应执行单条或批量写入。
- 参数:
resource (Any)
data (Any)
options (Any)
- 返回类型:
Any
- delete_one(resource, selector=None, **options)[源代码]
删除单个 Redis key 或首个匹配 MongoDB 文档。
- 参数:
resource (Any)
selector (Any)
options (Any)
- 返回类型:
Any
- delete_many(resource, selector=None, **options)[源代码]
批量删除 Redis keys 或 MongoDB 文档。
- 参数:
resource (Any)
selector (Any)
options (Any)
- 返回类型:
Any
- delete(resource, selector=None, **options)[源代码]
根据输入形态与
many选项自适应执行删除。- 参数:
resource (Any)
selector (Any)
options (Any)
- 返回类型:
Any
- exists(resource, selector=None, **options)[源代码]
判断 Redis key 或 MongoDB 匹配文档是否存在。
- 参数:
resource (Any)
selector (Any)
options (Any)
- 返回类型:
bool
- stream_query(sql, params=None, *, chunksize=50000, progress=False, retain=True, count_total=False, count_sql=None, total_rows=None, columns=None, json_fields=None, result='dataframe')[源代码]
打开可中断的分块查询流。
当
json_fields不为空时,适配器把 JSON 路径提取下推到数据库,只传输columns和指定的 JSON 子字段,不返回原始大 JSON。JSON 字段定义格式为{源字段: {输出字段: 路径或(路径, 默认值)}}。- 参数:
sql (str) -- 原始查询 SQL。JSON 投影会把该查询包装为子查询。
params (Any) -- 原始 SQL 的绑定参数,默认为
None。chunksize (int) -- 每次向 DB-API 流式游标请求的最大行数,默认 50000。
progress (bool) -- 是否显示读取进度,默认
False。未知总数时显示累计行数、速度和耗时。retain (bool) -- 是否保留已经产生的分块,默认
True。设为False后不能合并历史数据。count_total (bool) -- 是否为进度条自动执行
COUNT(1),默认False。count_sql (str | None) -- 为进度条显式指定的统计 SQL;提供后会执行该 SQL。
total_rows (int | None) -- 已知总行数;提供后不执行统计 SQL。
columns (Sequence[str] | None) -- JSON 投影时原样保留的普通输出字段;不能包含 JSON 源字段。
json_fields (Mapping[str, Mapping[str, Any]] | None) -- JSON 源字段、输出字段、JSONPath 和可选默认值的嵌套映射。
result (str) -- 每个分块的结果类型,支持
dataframe、records和rows。
- 返回:
可迭代、可主动停止并可合并已读数据的查询流。
- 返回类型:
- 抛出:
ValidationError -- 分块、进度、结果类型或 JSON 投影参数无效。
DatabaseCapabilityError -- 当前适配器不支持 JSON 字段投影。
DatabaseQueryError -- 统计查询或打开流式查询失败。
参考样例
>>> stream = db.stream_query( ... "SELECT id, huge_json FROM user_profile", ... columns=["id"], ... json_fields={ ... "huge_json": { ... "city": ("$.address.city", "未知"), ... "customer_id": "$.customer.id", ... } ... }, ... result="records", ... ) >>> for records in stream: ... consume(records)
- read_query(sql, params=None, *, chunksize=50000, progress=False, count_total=False, count_sql=None, total_rows=None, columns=None, json_fields=None, result='dataframe')[源代码]
消费完整查询流,并在中断后直接返回已经读取的数据。
参数语义与
stream_query()一致,但本方法自动消费所有分块。发生KeyboardInterrupt时会关闭底层资源,并按照result返回当前已合并数据。- 参数:
sql (str) -- 原始查询 SQL。
params (Any) -- 原始 SQL 的绑定参数,默认为
None。chunksize (int) -- 每次请求的最大行数,默认 50000。
progress (bool) -- 是否显示进度,默认
False。count_total (bool) -- 是否自动执行
COUNT(1)获取进度条总数,默认False。count_sql (str | None) -- 为进度条显式指定的统计 SQL。
total_rows (int | None) -- 已知总行数。
columns (Sequence[str] | None) -- JSON 投影时原样保留的普通输出字段;不能包含 JSON 源字段。
json_fields (Mapping[str, Mapping[str, Any]] | None) -- JSON 子字段投影映射。
result (str) -- 返回类型,支持
dataframe、records和rows。
- 返回:
完整数据或中断前已读取的部分数据。
- 抛出:
ValidationError -- 查询或投影参数无效。
DatabaseQueryError -- 流式查询失败。
- 返回类型:
Any
参考样例
>>> frame = db.read_query("SELECT * FROM events", progress=True) >>> rows = db.read_query("SELECT id FROM events", result="rows")
- export_schema(targets=None, *, output=None, excel_params=None)[源代码]
读取数据库表结构并生成中文字段元数据宽表。
- 参数:
targets (Sequence[str] | None) -- 可选数据库或表目标,例如
risk、risk.events;默认扫描适配器可见范围。output (Any | None) -- 可选
.xlsx输出路径。提供时通过dataframe2excel导出。excel_params (Mapping[str, Any] | None) -- 传递给
dataframe2excel的附加参数。
- 返回:
中文列名的表和字段信息 DataFrame;数据库原始元数据值保持不变。
- 返回类型:
pandas.DataFrame
- 抛出:
ValidationError -- 目标、输出扩展名或 Excel 参数无效。
DatabaseMetadataError -- 元数据读取、目标匹配或 Excel 导出失败。
- create_table(data, table_name, *, dialect_options=None)[源代码]
根据 DataFrame 字段和后端方言参数创建表。
- 参数:
data (DataFrame) -- 用于推断字段结构的非空列 DataFrame。
table_name (str) -- 表名或
数据库名.表名等限定名。dialect_options (Mapping[str, Any] | None) -- 后端专有建表参数,例如键、引擎、分区、字段类型和注释。
- 返回:
适配器返回的已执行 DDL 或原生结果。
- 抛出:
InputValidationError --
data不是带字段的 DataFrame。ValidationError -- 表名或方言参数无效。
DatabaseQueryError -- 建表 SQL 执行失败。
- 返回类型:
Any
- stream_write(data, table_name, *, mode='a', batch_size=10000, key_columns=None, columns=None, dialect_options=None)[源代码]
把 DataFrame、DataFrame 分块或行记录迭代器流式写入目标表。
- 参数:
data (Any) -- DataFrame、DataFrame 分块、记录字典或位置行的可迭代对象。
table_name (str) -- 目标表限定名。
mode (str) -- 写入模式。
a追加且主键重复不覆盖;r追加且主键重复覆盖;o保留表结构并清空重写;d删除并按首批数据重建表后写入。batch_size (int) -- 每个写入批次的最大行数,默认 10000。
key_columns (Sequence[str] | None) --
a或r使用的显式键字段;未提供时由支持的适配器读取表元数据。columns (Sequence[str] | None) -- 位置行字段名,或用于校验 DataFrame/记录字段顺序。
dialect_options (Mapping[str, Any] | None) -- 传递给目标适配器的建表和写入选项。
- 返回:
完成状态、接收/插入/更新/跳过行数及已提交批次数。
- 返回类型:
- 抛出:
InputValidationError -- 数据为空、批次字段不一致或缺少键字段。
ValidationError -- 模式、批次大小、表名或方言参数无效。
DatabaseCapabilityError -- 目标数据库或表不支持指定模式。
DatabaseWriteError -- 准备、批次写入或收尾失败;异常的
result保留部分统计。
- close()[源代码]
关闭连接池或原生客户端;重复调用不产生副作用。
关闭后所有查询和写入方法都会抛出
StateError。- 返回类型:
None
类外快捷操作
数据库与 NoSQL 的类外快捷操作。
快捷函数接受以下任一 source,并把实际操作委托给 Database 的同名方法:
含
db_type的连接配置映射;已创建的
Database实例;PyMySQL、python-oracledb、Impyla、PyODPS 等原生 DB-API 连接。
配置创建的连接池在操作结束后自动关闭;传入的 Database 和原生连接均视为借用,不由 快捷函数关闭。配置创建的流式查询在查询流完成、停止或关闭时释放连接池。
参考样例
>>> from hscredit.database import read_query
>>> config = {"db_type": "mysql", "host": "127.0.0.1", "user": "risk", "password": "***"}
>>> frame = read_query(config, "SELECT id FROM events")
- hscredit.database.shortcuts.query(source, sql, params=None, result='dataframe', *, db_type=None)[源代码]
快捷执行一次性 SQL 查询。
- 参数:
source (Any) -- 数据库配置、Database 实例或原生 DB-API 连接。
sql (str) -- 查询 SQL。
params (Any) -- 驱动绑定参数。
result (str) --
dataframe、records或rows。db_type (str | None) -- 原生连接无法自动识别时使用的数据库类型。
- 返回:
对应格式的完整查询结果。
- 返回类型:
Any
- hscredit.database.shortcuts.execute(source, sql, params=None, *, db_type=None)[源代码]
快捷执行单条 DDL 或 DML SQL。
- 参数:
source (Any) -- 数据库配置、Database 实例或原生 DB-API 连接。
sql (str) -- 要执行的 SQL。
params (Any) -- 驱动绑定参数。
db_type (str | None) -- 可选数据库类型。
- 返回:
适配器报告的影响行数或原生结果。
- 返回类型:
Any
- hscredit.database.shortcuts.executemany(source, sql, values, *, db_type=None)[源代码]
快捷批量执行同一条 SQL。
- 参数:
source (Any) -- 数据库配置、Database 实例或原生 DB-API 连接。
sql (str) -- 带驱动占位符的 SQL。
values (Any) -- 多组绑定值。
db_type (str | None) -- 可选数据库类型。
- 返回:
累计影响行数或原生结果。
- 返回类型:
Any
- hscredit.database.shortcuts.stream_query(source, sql, params=None, *, db_type=None, **options)[源代码]
快捷打开流式 SQL 查询;拥有的数据源随查询流关闭。
- 参数:
source (Any) -- 数据库配置、Database 实例或原生 DB-API 连接。
sql (str) -- 查询 SQL。
params (Any) -- 驱动绑定参数。
db_type (str | None) -- 可选数据库类型。
options (Any) -- 传递给
Database.stream_query()的分块、进度和 JSON 投影参数。
- 返回:
可中断并可合并已读数据的 QueryStream。
- 返回类型:
Any
- hscredit.database.shortcuts.read_query(source, sql, params=None, *, db_type=None, **options)[源代码]
快捷消费流式 SQL 查询并返回合并结果。
- 参数:
source (Any) -- 数据库配置、Database 实例或原生 DB-API 连接。
sql (str) -- 查询 SQL。
params (Any) -- 驱动绑定参数。
db_type (str | None) -- 可选数据库类型。
options (Any) -- 传递给
Database.read_query()的选项。
- 返回:
DataFrame、记录字典列表或原始行列表。
- 返回类型:
Any
- hscredit.database.shortcuts.export_schema(source, targets=None, *, db_type=None, **options)[源代码]
快捷导出数据库表和字段元数据。
- 参数:
source (Any) -- 数据库配置、Database 实例或原生 DB-API 连接。
targets (Any) -- 数据库或
数据库.表目标。db_type (str | None) -- 可选数据库类型。
options (Any) -- 输出路径和
excel_params。
- 返回:
中文字段元数据 DataFrame。
- 返回类型:
Any
- hscredit.database.shortcuts.create_table(source, data, table_name, *, db_type=None, **options)[源代码]
快捷根据 DataFrame 创建表。
- 参数:
source (Any) -- 数据库配置、Database 实例或原生 DB-API 连接。
data (Any) -- 用于推断表结构的 DataFrame。
table_name (str) -- 目标表限定名。
db_type (str | None) -- 可选数据库类型。
options (Any) -- 传递给
Database.create_table()的方言参数。
- 返回:
已执行 DDL 或适配器原生结果。
- 返回类型:
Any
- hscredit.database.shortcuts.stream_write(source, data, table_name, *, db_type=None, **options)[源代码]
快捷流式写入 DataFrame 或记录迭代器。
- 参数:
source (Any) -- 数据库配置、Database 实例或原生 DB-API 连接。
data (Any) -- DataFrame、分块或记录迭代器。
table_name (str) -- 目标表限定名。
db_type (str | None) -- 可选数据库类型。
options (Any) -- 写入模式、批次、键字段和方言参数。
- 返回:
WriteResult 写入统计。
- 返回类型:
Any
- hscredit.database.shortcuts.read_one(source, resource, selector=None, *, db_type=None, **options)[源代码]
快捷读取单个 Redis key 或 MongoDB 文档。
- 参数:
source (Any) -- 数据库配置或 Database 实例。
resource (Any) -- Redis key 或 MongoDB collection。
selector (Any) -- MongoDB 查询条件。
db_type (str | None) -- 可选数据库类型。
options (Any) -- 后端读取选项。
- 返回:
单值或单个文档。
- 返回类型:
Any
- hscredit.database.shortcuts.read_many(source, resource, selector=None, *, db_type=None, **options)[源代码]
快捷批量读取 Redis keys 或 MongoDB 文档。
参数与
read_one()一致,返回后端批量读取结果。- 参数:
source (Any)
resource (Any)
selector (Any)
db_type (str | None)
options (Any)
- 返回类型:
Any
- hscredit.database.shortcuts.read(source, resource, selector=None, *, db_type=None, **options)[源代码]
快捷自适应执行单条或批量 NoSQL 读取。
参数与
read_one()一致,具体单条/批量语义由适配器根据输入和选项确定。- 参数:
source (Any)
resource (Any)
selector (Any)
db_type (str | None)
options (Any)
- 返回类型:
Any
- hscredit.database.shortcuts.write_one(source, resource, data, *, db_type=None, **options)[源代码]
快捷写入单个 Redis key 或 MongoDB 文档。
- 参数:
source (Any) -- 数据库配置或 Database 实例。
resource (Any) -- Redis key 或 MongoDB collection。
data (Any) -- 要写入的值或文档。
db_type (str | None) -- 可选数据库类型。
options (Any) -- 后端写入选项。
- 返回:
后端写入结果。
- 返回类型:
Any
- hscredit.database.shortcuts.write_many(source, resource, data=None, *, db_type=None, **options)[源代码]
快捷批量写入 Redis key-value 或 MongoDB 文档。
参数与
write_one()一致,data为批量输入。- 参数:
source (Any)
resource (Any)
data (Any)
db_type (str | None)
options (Any)
- 返回类型:
Any
- hscredit.database.shortcuts.write(source, resource, data=None, *, db_type=None, **options)[源代码]
快捷自适应执行单条或批量 NoSQL 写入。
参数与
write_one()一致,具体单条/批量语义由适配器确定。- 参数:
source (Any)
resource (Any)
data (Any)
db_type (str | None)
options (Any)
- 返回类型:
Any
- hscredit.database.shortcuts.delete_one(source, resource, selector=None, *, db_type=None, **options)[源代码]
快捷删除单个 Redis key 或 MongoDB 文档。
- 参数:
source (Any) -- 数据库配置或 Database 实例。
resource (Any) -- Redis key 或 MongoDB collection。
selector (Any) -- MongoDB 查询条件。
db_type (str | None) -- 可选数据库类型。
options (Any) -- 后端删除选项。
- 返回:
后端删除结果。
- 返回类型:
Any
- hscredit.database.shortcuts.delete_many(source, resource, selector=None, *, db_type=None, **options)[源代码]
快捷批量删除 Redis keys 或 MongoDB 文档。
参数与
delete_one()一致,返回后端批量删除结果。- 参数:
source (Any)
resource (Any)
selector (Any)
db_type (str | None)
options (Any)
- 返回类型:
Any
- hscredit.database.shortcuts.delete(source, resource, selector=None, *, db_type=None, **options)[源代码]
快捷自适应执行单条或批量 NoSQL 删除。
参数与
delete_one()一致,具体单条/批量语义由适配器确定。- 参数:
source (Any)
resource (Any)
selector (Any)
db_type (str | None)
options (Any)
- 返回类型:
Any
- hscredit.database.shortcuts.exists(source, resource, selector=None, *, db_type=None, **options)[源代码]
快捷判断 Redis key 或 MongoDB 文档是否存在。
- 参数:
source (Any) -- 数据库配置或 Database 实例。
resource (Any) -- Redis key 或 MongoDB collection。
selector (Any) -- MongoDB 查询条件。
db_type (str | None) -- 可选数据库类型。
options (Any) -- 后端查询选项。
- 返回:
是否存在。
- 返回类型:
bool
流式查询
- class hscredit.database.stream.QueryStream(resource, *, chunksize, retain=True, total_rows=None, progress=False, result='dataframe', defaults=None)[源代码]
基类:
object可中断并可合并已读数据的分块查询迭代器。
QueryStream 由
Database.stream_query创建。迭代期间只持有当前数据库资源;retain=True时额外保留已经返回的标准化 DataFrame 分块,用于停止或中断后合并结果。参数
- resource
实现
fetchmany(size)、close()和columns的适配器查询资源。- chunksizeint
每次请求的最大行数。
- retainbool, default=True
是否保留已读取分块。
- total_rowsint, optional
进度条总行数。
- progressbool, default=False
是否显示 tqdm 进度条。
- result{"dataframe", "records", "rows"}, default="dataframe"
每次迭代和最终合并使用的结果类型。
- defaultsmapping, optional
JSON 投影字段在数据库返回
NULL时使用的默认值。
属性
- stateStreamState
当前运行、完成、中断、失败或关闭状态。
- rows_readint
已经读取并返回的累计行数。
- interrupt_reasonstr, optional
主动停止、中断或失败原因。
参考样例
>>> stream = db.stream_query("SELECT id FROM events", chunksize=1000) >>> first = next(stream) >>> stream.stop("仅抽取首批") >>> partial = stream.to_dataframe()
- 参数:
resource (Any)
chunksize (int)
retain (bool)
total_rows (int | None)
progress (bool)
result (str)
defaults (Mapping[str, Any] | None)
- stop(reason='用户主动停止')[源代码]
安全停止读取并关闭底层查询资源。
- 参数:
reason (str) -- 写入流状态和最终 DataFrame 属性的停止原因。
- 返回:
None。已经读取的数据可继续通过to_result()获取。- 返回类型:
None
- to_dataframe()[源代码]
把已保留分块合并为 DataFrame,并附加读取状态属性。
- 返回:
连续索引的合并 DataFrame。
attrs包含completed、rows_read、total_rows、state、interrupted_at和interrupt_reason。- 返回类型:
pandas.DataFrame
- 抛出:
StateError -- 创建流时设置了
retain=False。
- to_records()[源代码]
把已保留分块合并为记录字典列表。
- 返回:
每行一个字段名到原始值映射的列表。
- 抛出:
StateError -- 创建流时设置了
retain=False。- 返回类型:
List[Mapping[str, Any]]
- to_rows()[源代码]
把已保留分块合并为原始行元组列表。
- 返回:
字段顺序与查询结果一致的行元组列表。
- 抛出:
StateError -- 创建流时设置了
retain=False。- 返回类型:
List[Any]
- to_result()[源代码]
按创建流时的
result配置返回合并结果。- 返回:
DataFrame、
list[dict]或行元组列表。- 抛出:
StateError -- 创建流时设置了
retain=False。- 返回类型:
Any
- hscredit.database.types.RESULT_TYPES = frozenset({'dataframe', 'records', 'rows'})
Database.query、stream_query和read_query共用的结果类型。
JSON 投影扩展
- class hscredit.database.adapters.base.BaseDatabaseAdapter(*, connect_kwargs, pool_options, adapter_options=None)[源代码]
基类:
object数据库适配器公共基类。
参数
- connect_kwargsdict
传递给底层数据库驱动的连接参数。
- pool_optionsPoolOptions
已校验的连接池配置。
- adapter_optionsdict
仅由适配器解释的方言或原生通道参数。
属性
- database_typestr
注册表使用的数据库类型名称。
- capabilitiesDatabaseCapabilities
适配器默认保证的事务、读取、写入与元数据能力。
参考样例
自定义后端至少实现查询、流资源和所需写入方法;支持 JSON 字段投影时覆盖
json_extract_expression(),公共 SQL 包装由build_json_projection_sql()完成。- 参数:
connect_kwargs (Mapping[str, Any])
pool_options (PoolOptions)
adapter_options (Mapping[str, Any] | None)
- database_type = 'base'
- identifier_quote = '"'
- capabilities = DatabaseCapabilities(transactions=True, streaming_read=True, native_bulk_write=False, metadata_export=True, write_modes=frozenset({'d', 'o'}))
- property closed: bool
适配器是否已关闭。
- capabilities_for_table(table_name, table_metadata=None)[源代码]
返回目标表可保证的能力。
- 参数:
table_name (str)
table_metadata (Mapping[str, Any] | None)
- 返回类型:
- require_write_mode(table_name, mode, table_metadata=None)[源代码]
校验目标表是否支持指定写入模式。
- 参数:
table_name (str)
mode (str)
table_metadata (Mapping[str, Any] | None)
- 返回类型:
- json_extract_expression(column_sql, path)[源代码]
生成从单个 JSON 字段提取路径值的 SQL 表达式。
适配器只负责返回表达式,不添加输出别名。
path已经由公共层完成安全校验。- 参数:
column_sql (str) -- 已按当前方言引用的 JSON 源字段表达式。
path (str) -- 以
$开头的 JSONPath。
- 返回:
返回标量或嵌套 JSON 原始值的 SQL 表达式。
- 抛出:
DatabaseCapabilityError -- 适配器未实现 JSON 字段投影。
- 返回类型:
str
- build_json_projection_sql(sql, *, columns=None, json_fields=None)[源代码]
把原查询包装为仅返回普通字段和指定 JSON 子字段的 SQL。
- 参数:
sql (str) -- 输出中必须包含
columns和所有 JSON 源字段的原始查询。columns (Sequence[str] | None) -- 原样保留的普通字段名序列;不能包含
json_fields的源字段。json_fields (Mapping[str, Mapping[str, Any]] | None) --
{源字段: {输出字段: 路径或(路径, 默认值)}}映射。 默认值由 QueryStream 在分块返回前处理,不写入 SQL。
- 返回:
仅选择目标字段的后端方言 SQL;没有投影时返回原 SQL。
- 抛出:
ValidationError -- 字段名、JSONPath、默认值简写或输出名重复无效。
DatabaseCapabilityError -- 当前适配器没有 JSON 路径提取实现。
- 返回类型:
str
- create_table(data, table_name, *, dialect_options=None)[源代码]
创建目标表。
- 参数:
data (Any)
table_name (str)
dialect_options (Mapping[str, Any] | None)
- 返回类型:
Any
- resolve_key_columns(table_name, key_columns, first_batch, *, dialect_options=None)[源代码]
解析显式或数据库元数据中的主键字段。
- 参数:
table_name (str)
key_columns (Sequence[str] | None)
first_batch (Any)
dialect_options (Mapping[str, Any] | None)
- 返回类型:
Sequence[str] | None
- prepare_write(table_name, mode, first_batch, *, key_columns=None, dialect_options=None)[源代码]
在首个批次写入前校验并准备目标表。
- 参数:
table_name (str)
mode (str)
first_batch (Any)
key_columns (Sequence[str] | None)
dialect_options (Mapping[str, Any] | None)
- 返回类型:
None
- write_batch(table_name, batch, mode, batch_index, *, key_columns=None, dialect_options=None)[源代码]
写入一个已经校验的 DataFrame 批次。
- 参数:
table_name (str)
batch (Any)
mode (str)
batch_index (int)
key_columns (Sequence[str] | None)
dialect_options (Mapping[str, Any] | None)
- 返回类型:
Any
- finish_write(table_name, mode, result, *, dialect_options=None)[源代码]
完成适配器专有的写入收尾。
- 参数:
table_name (str)
mode (str)
result (Any)
dialect_options (Mapping[str, Any] | None)
- 返回类型:
None
适配器注册
- hscredit.database.registry.register_adapter(name, adapter_class, *, aliases=(), replace=False)[源代码]
注册自定义数据库适配器及其别名。
- 参数:
name (str) -- 新适配器的规范名称。
adapter_class (Type[BaseDatabaseAdapter]) --
BaseDatabaseAdapter的实现类。aliases (Iterable[str]) -- 可选别名集合。
replace (bool) -- 是否允许替换同名适配器和冲突别名,默认
False。
- 返回:
None。- 抛出:
ValidationError -- 名称、适配器类型或别名冲突无效。
- 返回类型:
None
参考样例
>>> register_adapter("custom_db", CustomDatabaseAdapter, aliases=["custom"]) >>> Database("custom", host="127.0.0.1")
- hscredit.database.registry.get_adapter_class(name)[源代码]
按数据库类型或别名获取适配器类。
内置适配器在第一次获取时才导入,从而保持所有数据库驱动均为可选依赖。
- 参数:
name (str) -- 数据库类型或别名。
- 返回:
已加载的适配器类。
- 抛出:
ValidationError -- 数据库类型未注册或导入入口无效。
- 返回类型:
Type[BaseDatabaseAdapter]
公共类型
- class hscredit.database.types.PoolOptions(mincached=0, maxcached=0, maxshared=0, maxconnections=0, blocking=False, maxusage=None, setsession=None, ping=1)[源代码]
基类:
objectDBUtils 兼容的连接池配置。
参数
mincached、maxcached、maxshared、maxconnections 与 DBUtils
PooledDB含义一致,值为 0 时沿用 DBUtils 的“不限制”语义。- 参数:
mincached (int)
maxcached (int)
maxshared (int)
maxconnections (int)
blocking (bool)
maxusage (int | None)
setsession (Tuple[str, ...] | None)
ping (int)
- mincached: int = 0
- maxcached: int = 0
- maxconnections: int = 0
- blocking: bool = False
- maxusage: int | None = None
- setsession: Tuple[str, ...] | None = None
- ping: int = 1
- classmethod from_mapping(value=None)[源代码]
从映射创建并校验连接池配置。
- 参数:
value (Mapping[str, Any] | None) --
PoolOptions、参数映射或None。- 返回:
已校验的不可变连接池配置。
- 抛出:
ValidationError -- 包含未知参数或参数范围无效。
- 返回类型:
- to_dbutils_kwargs()[源代码]
转换为 DBUtils
PooledDB关键字参数。- 返回:
可直接传给
PooledDB的新字典。- 返回类型:
dict
- 抛出:
ValidationError -- 当前配置组合无效。
- class hscredit.database.types.RedisPoolOptions(max_connections=None, blocking=False, timeout=None)[源代码]
基类:
objectRedis 原生连接池配置。
- 参数:
max_connections (int | None)
blocking (bool)
timeout (float | None)
- max_connections: int | None = None
- blocking: bool = False
- timeout: float | None = None
- class hscredit.database.types.MongoPoolOptions(min_pool_size=0, max_pool_size=100, max_connecting=2, wait_queue_timeout_ms=None, max_idle_time_ms=None)[源代码]
基类:
objectPyMongo
MongoClient原生连接池配置。- 参数:
min_pool_size (int)
max_pool_size (int)
max_connecting (int)
wait_queue_timeout_ms (int | None)
max_idle_time_ms (int | None)
- min_pool_size: int = 0
- max_pool_size: int = 100
- max_connecting: int = 2
- wait_queue_timeout_ms: int | None = None
- max_idle_time_ms: int | None = None
- class hscredit.database.types.DatabaseCapabilities(transactions=True, streaming_read=True, native_bulk_write=False, metadata_export=True, write_modes=<factory>)[源代码]
基类:
object数据库或目标表可保证的能力。
参数
- transactionsbool
是否支持事务提交和回滚。
- streaming_readbool
是否支持流式读取。
- native_bulk_writebool
是否具有后端原生批量写入通道。
- metadata_exportbool
是否支持表结构扫描。
- write_modesfrozenset[str]
可保证的
a/r/o/d写入模式集合。
- 参数:
transactions (bool)
streaming_read (bool)
native_bulk_write (bool)
metadata_export (bool)
write_modes (FrozenSet[str])
- transactions: bool = True
- streaming_read: bool = True
- native_bulk_write: bool = False
- metadata_export: bool = True
- write_modes: FrozenSet[str]
- class hscredit.database.types.WriteResult(mode, completed, rows_received=0, rows_inserted=None, rows_updated=None, rows_skipped=None, batches_committed=0, failed_batch=None, consistency=None, details=<factory>)[源代码]
基类:
object流式写入结果。
参数
- mode{"a", "r", "o", "d"}
本次写入模式。
- completedbool
是否完成全部批次和适配器收尾。
- rows_received、rows_inserted、rows_updated、rows_skippedint, optional
输入及后端报告的行数统计。
- batches_committedint
已成功提交的批次数。
- failed_batchint, optional
失败批次编号;准备阶段为 0,收尾阶段为已提交批次数加 1。
- consistencystr, optional
后端最终一致性说明。
- detailsdict
适配器附加的原始统计信息。
- 参数:
mode (str)
completed (bool)
rows_received (int)
rows_inserted (int | None)
rows_updated (int | None)
rows_skipped (int | None)
batches_committed (int)
failed_batch (int | None)
consistency (str | None)
details (Dict[str, Any])
- mode: str
- completed: bool
- rows_received: int = 0
- rows_inserted: int | None = None
- rows_updated: int | None = None
- rows_skipped: int | None = None
- batches_committed: int = 0
- failed_batch: int | None = None
- consistency: str | None = None
- details: Dict[str, Any]
- class hscredit.database.types.NoSQLWriteResult(operation, acknowledged, affected_count=None, matched_count=None, modified_count=None, identifiers=<factory>, details=<factory>)[源代码]
基类:
objectRedis 与 MongoDB 共用的写入、更新和删除结果。
- 参数:
operation (str)
acknowledged (bool)
affected_count (int | None)
matched_count (int | None)
modified_count (int | None)
identifiers (Tuple[Any, ...])
details (Mapping[str, Any])
- operation: str
- acknowledged: bool
- affected_count: int | None = None
- matched_count: int | None = None
- modified_count: int | None = None
- identifiers: Tuple[Any, ...]
- details: Mapping[str, Any]
- class hscredit.database.types.StreamState(value)[源代码]
基类:
str,Enum流式查询生命周期状态。
running表示仍可读取;completed表示自然耗尽;interrupted表示主动停止 或键盘中断;failed表示读取失败;closed表示在耗尽前关闭。- RUNNING = 'running'
- COMPLETED = 'completed'
- INTERRUPTED = 'interrupted'
- FAILED = 'failed'
- CLOSED = 'closed'
- class hscredit.database.metadata.QualifiedTarget(raw, parts)[源代码]
基类:
object未改写大小写的数据库对象限定名。
参数
- rawstr
用户提供并去除首尾空白后的原始目标。
- partstuple[str, ...]
以点拆分的数据库、模式和表名片段。
- 参数:
raw (str)
parts (Tuple[str, ...])
- raw: str
- parts: Tuple[str, ...]
- classmethod parse(value)[源代码]
解析以点分隔的数据库、模式和表目标。
- 参数:
value (str) --
数据库、数据库.表或多级限定名。- 返回:
保持原始大小写的限定目标。
- 抛出:
ValidationError -- 目标为空或包含空片段。
- 返回类型:
- class hscredit.database.metadata.MetadataInspection(rows=<factory>, errors=<factory>)[源代码]
基类:
object适配器元数据扫描结果。
参数
- rowsiterable of mapping
使用内部英文字段键保存的表/字段记录。
- errorslist
部分数据库或表扫描失败时保留的原始错误信息。
- 参数:
rows (Iterable[Mapping[str, Any]])
errors (List[Any])
- rows: Iterable[Mapping[str, Any]]
- errors: List[Any]
异常
数据库模块异常。
所有异常均兼容 hscredit 统一异常体系,并为调用方保留底层驱动异常链。
- exception hscredit.database.exceptions.DatabaseWriteError(message, *, result=None)[源代码]
-
数据库写入失败,并可携带部分写入结果。
- 参数:
message (str)
result (WriteResult | None)