数据库 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)

property closed: bool

数据库门面是否已经关闭。

返回:

已调用 close() 时为 True,否则为 False

返回类型:

bool

query(sql, params=None, result='dataframe')[源代码]

一次性执行查询并返回全部结果。

大结果集建议改用 stream_query()read_query(),避免一次性占用过多内存。

参数:
  • sql (str) -- 由当前数据库执行的查询 SQL。

  • params (Any) -- 按底层驱动参数风格绑定的 SQL 参数,默认为 None

  • result (str) -- 结果类型,支持 dataframerecordsrows

返回:

dataframe 返回 DataFrame;records 返回记录字典列表;rows 返回原始行列表。

抛出:
返回类型:

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

返回:

适配器报告的影响行数或原生执行结果。

抛出:
返回类型:

Any

executemany(sql, values)[源代码]

使用多组绑定值批量执行同一条 SQL。

参数:
  • sql (str) -- 带驱动占位符的 SQL。

  • values (Any) -- 每行一组参数的可迭代对象。

返回:

适配器报告的累计影响行数或原生执行结果。

抛出:
返回类型:

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) -- 每个分块的结果类型,支持 dataframerecordsrows

返回:

可迭代、可主动停止并可合并已读数据的查询流。

返回类型:

QueryStream

抛出:

参考样例

>>> 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) -- 返回类型,支持 dataframerecordsrows

返回:

完整数据或中断前已读取的部分数据。

抛出:
返回类型:

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) -- 可选数据库或表目标,例如 riskrisk.events;默认扫描适配器可见范围。

  • output (Any | None) -- 可选 .xlsx 输出路径。提供时通过 dataframe2excel 导出。

  • excel_params (Mapping[str, Any] | None) -- 传递给 dataframe2excel 的附加参数。

返回:

中文列名的表和字段信息 DataFrame;数据库原始元数据值保持不变。

返回类型:

pandas.DataFrame

抛出:
create_table(data, table_name, *, dialect_options=None)[源代码]

根据 DataFrame 字段和后端方言参数创建表。

参数:
  • data (DataFrame) -- 用于推断字段结构的非空列 DataFrame。

  • table_name (str) -- 表名或 数据库名.表名 等限定名。

  • dialect_options (Mapping[str, Any] | None) -- 后端专有建表参数,例如键、引擎、分区、字段类型和注释。

返回:

适配器返回的已执行 DDL 或原生结果。

抛出:
返回类型:

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) -- ar 使用的显式键字段;未提供时由支持的适配器读取表元数据。

  • columns (Sequence[str] | None) -- 位置行字段名,或用于校验 DataFrame/记录字段顺序。

  • dialect_options (Mapping[str, Any] | None) -- 传递给目标适配器的建表和写入选项。

返回:

完成状态、接收/插入/更新/跳过行数及已提交批次数。

返回类型:

WriteResult

抛出:
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) -- dataframerecordsrows

  • 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

close()[源代码]

关闭查询资源但不把状态标记为主动中断。

尚在运行时状态会变为 closed;重复调用不会重复关闭游标或连接。

返回类型:

None

to_dataframe()[源代码]

把已保留分块合并为 DataFrame,并附加读取状态属性。

返回:

连续索引的合并 DataFrame。attrs 包含 completedrows_readtotal_rowsstateinterrupted_atinterrupt_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.querystream_queryread_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

适配器是否已关闭。

ensure_open()[源代码]

确保适配器仍可使用。

返回类型:

None

capabilities_for_table(table_name, table_metadata=None)[源代码]

返回目标表可保证的能力。

参数:
  • table_name (str)

  • table_metadata (Mapping[str, Any] | None)

返回类型:

DatabaseCapabilities

require_write_mode(table_name, mode, table_metadata=None)[源代码]

校验目标表是否支持指定写入模式。

参数:
  • table_name (str)

  • mode (str)

  • table_metadata (Mapping[str, Any] | None)

返回类型:

DatabaseCapabilities

build_count_sql(sql)[源代码]

生成通用子查询计数 SQL。

参数:

sql (str)

返回类型:

str

quote_identifier(identifier)[源代码]

按当前数据库方言引用单个标识符。

参数:

identifier (str)

返回类型:

str

quote_qualified_name(name)[源代码]

逐段引用数据库对象限定名。

参数:

name (str)

返回类型:

str

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。

抛出:
返回类型:

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

query(sql, params=None, result='dataframe')[源代码]

执行查询。

参数:
  • sql (str)

  • params (Any)

  • result (str)

返回类型:

Any

close()[源代码]

关闭适配器持有的资源。

返回类型:

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]

hscredit.database.registry.available_adapters()[源代码]

返回当前已注册的规范数据库类型。

返回:

按名称排序且不包含别名的数据库类型元组。

返回类型:

tuple[str, ...]

公共类型

class hscredit.database.types.PoolOptions(mincached=0, maxcached=0, maxshared=0, maxconnections=0, blocking=False, maxusage=None, setsession=None, ping=1)[源代码]

基类:object

DBUtils 兼容的连接池配置。

参数

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
maxshared: 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 -- 包含未知参数或参数范围无效。

返回类型:

PoolOptions

to_dbutils_kwargs()[源代码]

转换为 DBUtils PooledDB 关键字参数。

返回:

可直接传给 PooledDB 的新字典。

返回类型:

dict

抛出:

ValidationError -- 当前配置组合无效。

class hscredit.database.types.RedisPoolOptions(max_connections=None, blocking=False, timeout=None)[源代码]

基类:object

Redis 原生连接池配置。

参数:
  • max_connections (int | None)

  • blocking (bool)

  • timeout (float | None)

max_connections: int | None = None
blocking: bool = False
timeout: float | None = None
classmethod from_mapping(value=None)[源代码]

从映射创建 Redis 连接池配置。

参数:

value (Mapping[str, Any] | None)

返回类型:

RedisPoolOptions

to_redis_kwargs()[源代码]

转换为 redis-py 连接池参数。

返回类型:

Dict[str, Any]

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)[源代码]

基类:object

PyMongo 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
classmethod from_mapping(value=None)[源代码]

从映射创建 MongoDB 连接池配置。

参数:

value (Mapping[str, Any] | None)

返回类型:

MongoPoolOptions

to_mongo_kwargs()[源代码]

转换为 MongoClient 驼峰连接池参数。

返回类型:

Dict[str, Any]

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>)[源代码]

基类:object

Redis 与 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 -- 目标为空或包含空片段。

返回类型:

QualifiedTarget

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.DatabaseError[源代码]

基类:HSCreditError

数据库模块基础异常。

exception hscredit.database.exceptions.DatabaseConnectionError[源代码]

基类:DatabaseError

数据库连接或连接池操作失败。

exception hscredit.database.exceptions.DatabaseQueryError[源代码]

基类:DatabaseError

数据库查询或 SQL 执行失败。

exception hscredit.database.exceptions.DatabaseWriteError(message, *, result=None)[源代码]

基类:DatabaseError

数据库写入失败,并可携带部分写入结果。

参数:
exception hscredit.database.exceptions.DatabaseMetadataError[源代码]

基类:DatabaseError

数据库元数据读取或导出失败。

exception hscredit.database.exceptions.DatabaseCapabilityError[源代码]

基类:DatabaseError

数据库或目标表不支持所请求的能力。