阅读提示

这篇讲解 betalens_db_manager 的完整用法——从初始化 schema、导入数据、到 GUI 界面管理。学完这篇,你可以独立完成数据库从零到可用的全部操作,不需要依赖别人给你准备数据。

导言

betalens_db_manager 是 Betalens 生态里专门负责”数据进库”和”数据库维护”的模块。它的职责边界很清晰:

  • :建库、Schema 管理、数据导入、任务记录、GUI 管理界面。
  • 不做:运行时数据查询(那是 betalens.datafeed 的事)。

这套工具对新手最友好的入口是命令行和 GUI;对开发者来说,知道 Python API 可以让你把导入流程自动化、流水线化。

核心模块一览

betalens_db_manager 对外暴露以下主要对象:

1
2
3
4
5
6
7
8
from betalens_db_manager import (
SchemaManager, # Schema 创建、迁移、核验
DatabaseClient, # 连接管理、基础查询
QueryRequest, # 分页查询请求
ImportJobRunner, # 批量导入任务
ImportRecordStore, # 导入记录(防重复导入)
DatabaseManager, # 统一门面(plan / bootstrap / verify)
)

通常不需要直接调用底层对象,用顶层的 DatabaseManager 即可。

一、初始化数据库:Schema 创建

方式一:一键脚本(推荐新手)

1
2
3
4
5
cd C:\Users\Janis\OneDrive\betalens
python -m betalens_db_manager plan # 查看建表计划(dry-run)
.\betalens_db_manager\init_local.bat # 执行初始化
python -m betalens_db_manager verify # 核验
python -m betalens_db_manager verify --deep # 深度核验(检查数据完整性)

init_local.bat 会优先使用仓库的 .venv 中的 Python,找不到才用系统 Python。

方式二:命令行一步到位

1
python -m betalens_db_manager init --yes

--yes 跳过确认提示。init 会自动处理迁移版本(0001 ~ 0009),并保证 schema 版本一致性。

方式三:Python API

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
from betalens_db_manager import DatabaseManager

manager = DatabaseManager()

# 查看建表计划
plan = manager.plan()
print(plan)

# 初始化并验证
report = manager.bootstrap(
create_database_if_missing=True, # 如果数据库不存在,自动创建
create_compat_views=True, # 创建兼容视图(如 daily_market)
verify=True,
)
print(report.status) # "completed" / "failed" 等

二、导入数据:ImportJobRunner

导入前准备

数据导入有两种主流来源:

  1. EDE 格式数据包(常见于券商资管、合规数据供应商)
  2. Wind 数据接口(通过 Wind Python 模块拉取)

此外还支持:

  • 指数成分index_universe 模块导入中证/申万等指数成分
  • 交易日历trade_calendar 导入
  • 自定义 CSV/Parquet:通过 ImportJobRunner 通用入口

Wind 数据导入

1
2
3
4
5
6
7
8
9
10
11
12
13
14
from betalens_db_manager import ImportJobRunner, ImportRecordStore

runner = ImportJobRunner()
store = ImportRecordStore() # 记录导入历史,防重复

# 导入日行情数据(从 Wind 拉取)
result = runner.run(
source_type="wind",
target_table="daily_market",
start_date="2020-01-01",
end_date="2024-12-31",
store=store, # 传入 store 启用去重
)
print(f"导入记录数:{result.rows_imported}")

EDE 数据包导入

1
2
3
4
5
6
result = runner.run(
source_type="ede",
source_path=r"D:\data\market\2024",
target_table="daily_market",
store=store,
)

指数成分导入

1
2
3
4
5
6
from betalens_db_manager import DatabaseManager

manager = DatabaseManager()
manifest = manager.plan_manifest(r"D:\data\indexes.yaml")
# indexes.yaml 中列出要导入的指数代码列表
report = manager.bootstrap(manifest=manifest)

导入记录与防重

ImportRecordStore 用 SQLite 文件(logs/database-manager/jobs.sqlite3)记录每次导入的时间、表名、起止日期、记录数。下次再导入时,自动跳过已有时间段:

1
2
3
# 第二次导入同一时间段,不会重复写入
result = runner.run(...)
print(result.skipped) # True,表示被记录器跳过

三、冲突检测与数据回滚

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
from betalens_db_manager import DatabaseClient

client = DatabaseClient()

# 检测某时间段是否有数据冲突
conflicts = client.check_conflicts(
table="daily_market",
date_range=("2024-01-01", "2024-01-31"),
)
if conflicts:
print("发现冲突数据:", conflicts)
# 决定是删除旧数据还是跳过新数据
client.delete_range(
table="daily_market",
date_range=("2024-01-01", "2024-01-31"),
)

四、GUI 管理界面

如果不想用命令行,betalens_db_manager 提供了 PySide6 图形界面:

1
python -m betalens_db_manager gui

界面包含:

  • 连接管理:填写 host/port/dbname/user/password,连接测试。
  • Schema 视图:查看当前数据库有哪些表、字段、数据量。
  • 导入任务:拖拽文件或填入 Wind 参数,发起导入任务。
  • 导入记录:查看历史导入、重复导入检测结果。
  • 冲突管理:删除/覆盖冲突数据。

GUI 的连接信息只用于当前进程,不会写入本地文件,适合临时查询。

五、Migration 与版本控制

Betalens 的 Schema 通过版本化的 migration 文件管理(0001 ~ 0009 等)。每个 migration 有规范化 checksum(LF 格式):

  • init_local.bat 执行时按顺序应用未执行的 migration。
  • verify --deep 会检测 checksum,如果数据库中的 migration 历史被篡改,会拒绝启动。
  • 从旧版本数据库升级时,0009 的早期阻断式 checksum 被允许兼容,平滑升级。
  • 禁止 schema 降级——如果当前 schema 版本高于 migration 文件,框架会报错。

六、任务记录

所有导入任务都记录在 logs/database-manager/jobs.sqlite3 中:

1
2
3
4
5
6
from betalens_db_manager import ImportRecordStore

store = ImportRecordStore()
records = store.get_history(table="daily_market")
for r in records:
print(f"{r.date}: {r.rows_imported} rows, status={r.status}")

团队协作时,可以把这个 SQLite 文件也提交到仓库,确保大家有相同的导入历史记录。

常见错误

1. betalens_db_manager verify 报错 schema checksum 不匹配

原因:数据库被手动修改,或从备份恢复了旧版本。解决:

1
2
3
# 删库重建(谨慎操作,会丢失所有数据!)
python -m betalens_db_manager init --yes --force
python -m betalens_db_manager verify --deep

2. 导入 Wind 数据报权限错误

检查 Wind Python 模块是否正确安装,以及数据库用户是否有写权限:

1
2
GRANT ALL PRIVILEGES ON DATABASE datafeed TO your_user;
GRANT ALL ON SCHEMA betalens TO your_user;

3. jobs.sqlite3 锁死

如果上次导入异常中断,SQLite 文件可能被锁。删掉 logs/database-manager/jobs.sqlite3,重新导入(会丢失历史记录,但导入任务本身不受影响)。

4. GUI 连接失败

GUI 不会保存密码到文件,需要每次手动输入。如果忘了密码,用 pgAdmin 重置后再填入 GUI。

开发者侧:导入流水线的自动化

如果你的数据来源比较复杂(比如每天定时从多个数据源拉取并入库),可以写一个 Python 脚本做全自动化:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
from betalens_db_manager import DatabaseManager, ImportJobRunner, ImportRecordStore
from datetime import date, timedelta

def daily_import(end_date: date):
manager = DatabaseManager()
runner = ImportJobRunner()
store = ImportRecordStore()
start_date = end_date - timedelta(days=30)

# 导入日行情
result = runner.run(
source_type="wind",
target_table="daily_market",
start_date=start_date.isoformat(),
end_date=end_date.isoformat(),
store=store,
)
print(f"daily_market: {result.rows_imported} rows")

# 导入指数成分(以中证800为例)
index_result = runner.run(
source_type="wind",
target_table="index_constituent",
codes=["000906.SH"], # 中证800
trade_date=end_date.isoformat(),
store=store,
)
print(f"index_constituent: {index_result.rows_imported} rows")

if __name__ == "__main__":
from datetime import date
daily_import(date.today())

配合 Windows 任务计划程序或 Linux cron,可以每天定时跑一次。

延伸阅读