
Stock Research Group
- 53 installs
- 61 repo stars
- Updated March 16, 2026
- kirkluokun/awesome-a-stock-openclawskills
Helps with ai & agent building tasks.
About
stock-research-group is a Claude Code skill for ai & agent building. It helps solo builders move faster with AI-assisted development.
- stock-research-group
- AI & Agent Building
- AI-coding skill
Stock Research Group by the numbers
- 53 all-time installs (skills.sh)
- Ranked #6,979 of 16,546 AI & Agent Building skills by installs in the Skillselion catalog
- Data as of Aug 4, 2026 (Skillselion catalog sync)
npx skills add https://github.com/kirkluokun/awesome-a-stock-openclawskills --skill stock-research-groupAdd your badge
Show developers this skill is listed on Skillselion. Paste this into your README.
| Installs | 53 |
|---|---|
| repo stars | ★ 61 |
| Last updated | March 16, 2026 |
| Repository | kirkluokun/awesome-a-stock-openclawskills ↗ |
What it does
Helps with ai & agent building tasks.
Files
# analyzers/balancesheet/_shared/db.py
"""资产负债表数据库工具"""
import json
import sqlite3
from datetime import datetime
from pathlib import Path
from config import DB_PATH
SCHEMA_PATH = Path(__file__).parent.parent / "schema.sql"
def connect(db_path: Path = DB_PATH) -> sqlite3.Connection:
"""连接数据库"""
conn = sqlite3.connect(str(db_path))
conn.row_factory = sqlite3.Row
return conn
def init_table(conn: sqlite3.Connection = None):
"""初始化表结构"""
if conn is None:
conn = connect()
should_close = True
else:
should_close = False
try:
schema = SCHEMA_PATH.read_text(encoding="utf-8")
conn.executescript(schema)
conn.commit()
finally:
if should_close:
conn.close()
def upsert_balancesheet(conn: sqlite3.Connection, rows: list) -> int:
"""批量插入/更新资产负债表数据
Args:
conn: 数据库连接
rows: 数据行列表,每行包含所有字段
Returns:
插入/更新的行数
"""
if not rows:
return 0
# 基础列(常用字段)
base_cols = [
'ts_code', 'ann_date', 'f_ann_date', 'end_date', 'report_type',
'comp_type', 'end_type', 'total_share', 'money_cap',
'accounts_receiv', 'inventories', 'total_cur_assets', 'fix_assets',
'total_assets', 'st_borr', 'acct_payable', 'total_cur_liab',
'total_liab', 'undistr_porfit', 'total_hldr_eqy_exc_min_int'
]
# 唯一键:如果ann_date为空,使用ts_code+end_date+report_type作为唯一键
# 这样可以确保同一报告期的数据能正确更新
unique_cols = ['ts_code', 'end_date', 'report_type']
now = datetime.now().strftime("%Y-%m-%d %H:%M:%S")
# 准备数据
db_rows = []
for row in rows:
db_row = {}
# 提取基础列
for col in base_cols:
db_row[col] = row.get(col)
# 存储完整JSON
db_row['payload_json'] = json.dumps(row, ensure_ascii=False, default=str)
db_row['source'] = 'tushare'
db_row['created_at'] = now
db_row['updated_at'] = now
# 处理唯一键中的NULL值
for col in unique_cols:
if db_row.get(col) is None:
db_row[col] = ''
# 如果ann_date为空,设置为空字符串(不影响唯一键判断)
if db_row.get('ann_date') is None:
db_row['ann_date'] = ''
db_rows.append(db_row)
# 构建SQL
all_cols = list(db_rows[0].keys())
placeholders = ", ".join(["?"] * len(all_cols))
columns_sql = ", ".join(all_cols)
conflict_sql = ", ".join(unique_cols)
update_cols = [c for c in all_cols if c not in unique_cols + ['created_at']]
update_sql = ", ".join([f"{c} = excluded.{c}" for c in update_cols])
sql = (
f"INSERT INTO balancesheet ({columns_sql}) "
f"VALUES ({placeholders}) "
f"ON CONFLICT({conflict_sql}) DO UPDATE SET {update_sql}"
)
values = [[r.get(col) for col in all_cols] for r in db_rows]
try:
conn.executemany(sql, values)
conn.commit()
return len(db_rows)
except Exception as e:
conn.rollback()
raise
if __name__ == "__main__":
# CLI入口:初始化表
import argparse
parser = argparse.ArgumentParser()
parser.add_argument("--init", action="store_true", help="初始化表结构")
args = parser.parse_args()
if args.init:
init_table()
print("表结构初始化完成")
# analyzers/balancesheet/pipeline.py
"""资产负债表分析流程入口"""
import sys
from pathlib import Path
PROJECT_ROOT = Path(__file__).resolve().parents[2]
sys.path.insert(0, str(PROJECT_ROOT))
from analyzers.balancesheet.search import (
get_balancesheet,
get_balancesheet_history,
get_field_value,
search_by_field,
ensure_data
)
from fetchers.balancesheet import fetch_and_save
def run_fetch(**kwargs):
"""拉取数据"""
count = fetch_and_save(**kwargs)
print(f"保存 {count} 条记录")
return count
def run_query(ts_code: str, field: str = None, end_date: str = None):
"""查询数据"""
if field:
value = get_field_value(ts_code, field, end_date)
if value is not None:
print(f"{ts_code} {field} = {value:,.0f}")
else:
print(f"未找到数据")
else:
record = get_balancesheet(ts_code, end_date)
if record:
print(f"报告期: {record.get('end_date')}")
total_assets = record.get('total_assets') or 0
inventories = record.get('inventories') or 0
accounts_receiv = record.get('accounts_receiv') or 0
print(f"总资产: {total_assets:,.0f}")
print(f"存货: {inventories:,.0f}")
print(f"应收账款: {accounts_receiv:,.0f}")
else:
print("未找到数据")
def run_history(ts_code: str, limit: int = 4):
"""查询历史数据"""
records = get_balancesheet_history(ts_code, limit=limit)
print(f"\n{ts_code} 最近{len(records)}期数据:")
print("-" * 60)
for r in records:
total_assets = r.get('total_assets') or 0
inventories = r.get('inventories') or 0
print(f"{r['end_date']}: 总资产={total_assets:,.0f}, "
f"存货={inventories:,.0f}")
if __name__ == "__main__":
import argparse
parser = argparse.ArgumentParser(description="资产负债表分析")
subparsers = parser.add_subparsers(dest="command")
# fetch命令
p_fetch = subparsers.add_parser("fetch", help="拉取数据")
p_fetch.add_argument("--ts-code")
p_fetch.add_argument("--start-date")
p_fetch.add_argument("--end-date")
p_fetch.add_argument("--period", help="报告期 YYYYMMDD(VIP接口必需)")
p_fetch.add_argument("--report-type", default="1")
p_fetch.add_argument("--comp-type", help="公司类型")
p_fetch.add_argument("--vip", action="store_true", help="使用VIP接口(全量拉取,需提供period)")
# query命令
p_query = subparsers.add_parser("query", help="查询数据")
p_query.add_argument("--ts-code", required=True)
p_query.add_argument("--field")
p_query.add_argument("--end-date")
# history命令
p_history = subparsers.add_parser("history", help="查询历史")
p_history.add_argument("--ts-code", required=True)
p_history.add_argument("--limit", type=int, default=4)
args = parser.parse_args()
if args.command == "fetch":
if args.vip and not args.period:
parser.error("使用VIP接口时必须提供 --period 参数")
run_fetch(
ts_code=args.ts_code,
start_date=args.start_date,
end_date=args.end_date,
period=args.period,
report_type=args.report_type,
comp_type=args.comp_type,
use_vip=args.vip
)
elif args.command == "query":
run_query(args.ts_code, args.field, args.end_date)
elif args.command == "history":
run_history(args.ts_code, args.limit)
PRAGMA journal_mode=WAL;
CREATE TABLE IF NOT EXISTS balancesheet (
ts_code TEXT NOT NULL,
ann_date TEXT,
f_ann_date TEXT,
end_date TEXT NOT NULL,
report_type TEXT,
comp_type TEXT,
end_type TEXT,
total_share REAL,
money_cap REAL,
accounts_receiv REAL,
inventories REAL,
total_cur_assets REAL,
fix_assets REAL,
total_assets REAL,
st_borr REAL,
acct_payable REAL,
total_cur_liab REAL,
total_liab REAL,
undistr_porfit REAL,
total_hldr_eqy_exc_min_int REAL,
source TEXT,
payload_json TEXT NOT NULL,
created_at TEXT NOT NULL,
updated_at TEXT NOT NULL,
UNIQUE(ts_code, end_date, report_type)
);
CREATE INDEX IF NOT EXISTS idx_balancesheet_ts_code ON balancesheet(ts_code);
CREATE INDEX IF NOT EXISTS idx_balancesheet_end_date ON balancesheet(end_date);
CREATE INDEX IF NOT EXISTS idx_balancesheet_ts_end ON balancesheet(ts_code, end_date);
# analyzers/balancesheet/search.py
"""资产负债表查询功能"""
import json
import sys
from pathlib import Path
from typing import Optional, List, Dict
PROJECT_ROOT = Path(__file__).resolve().parents[2]
sys.path.insert(0, str(PROJECT_ROOT))
from analyzers.balancesheet._shared.db import connect
def get_balancesheet(
ts_code: str,
end_date: str = None,
report_type: str = "1"
) -> Optional[Dict]:
"""
查询单条资产负债表记录
Args:
ts_code: 股票代码
end_date: 报告期 YYYYMMDD,默认最新
report_type: 报表类型,默认1(合并报表)
Returns:
记录字典,包含基础列和完整字段
"""
conn = connect()
try:
if end_date:
sql = """
SELECT * FROM balancesheet
WHERE ts_code = ? AND end_date = ? AND report_type = ?
ORDER BY ann_date DESC
LIMIT 1
"""
cur = conn.execute(sql, (ts_code, end_date, report_type))
else:
sql = """
SELECT * FROM balancesheet
WHERE ts_code = ? AND report_type = ?
ORDER BY end_date DESC, ann_date DESC
LIMIT 1
"""
cur = conn.execute(sql, (ts_code, report_type))
row = cur.fetchone()
if not row:
return None
result = dict(row)
# 解析完整字段
if result.get('payload_json'):
full_data = json.loads(result['payload_json'])
result.update(full_data)
return result
finally:
conn.close()
def get_balancesheet_history(
ts_code: str,
start_date: str = None,
end_date: str = None,
report_type: str = "1",
limit: int = None
) -> List[Dict]:
"""
查询历史资产负债表记录
Args:
ts_code: 股票代码
start_date: 开始日期 YYYYMMDD
end_date: 结束日期 YYYYMMDD
report_type: 报表类型
limit: 限制返回数量(如4表示最近4个季度)
Returns:
记录列表,按end_date降序
"""
conn = connect()
try:
conditions = ["ts_code = ?", "report_type = ?"]
params = [ts_code, report_type]
if start_date:
conditions.append("end_date >= ?")
params.append(start_date)
if end_date:
conditions.append("end_date <= ?")
params.append(end_date)
sql = f"""
SELECT * FROM balancesheet
WHERE {' AND '.join(conditions)}
ORDER BY end_date DESC, ann_date DESC
"""
if limit:
sql += f" LIMIT {limit}"
cur = conn.execute(sql, params)
rows = cur.fetchall()
results = []
for row in rows:
result = dict(row)
# 解析完整字段
if result.get('payload_json'):
full_data = json.loads(result['payload_json'])
result.update(full_data)
results.append(result)
return results
finally:
conn.close()
def ensure_data(
ts_code: str,
end_date: str = None,
years: int = 1,
report_type: str = "1"
) -> bool:
"""
确保有足够的历史数据,不足时自动拉取
Args:
ts_code: 股票代码
end_date: 截止日期,默认当前日期
years: 需要的数据年数
report_type: 报表类型
Returns:
是否成功获取数据
"""
from datetime import datetime, timedelta
from fetchers.balancesheet import fetch_and_save
if end_date is None:
end_date = datetime.now().strftime("%Y%m%d")
# 检查现有数据
existing = get_balancesheet_history(
ts_code=ts_code,
end_date=end_date,
report_type=report_type,
limit=years * 4 + 4 # 多查一些确保
)
# 计算需要的日期范围(N+1年,确保有足够数据计算同比)
need_years = years + 1
start_date_dt = datetime.strptime(end_date, "%Y%m%d") - timedelta(days=need_years * 365)
start_date = start_date_dt.strftime("%Y%m%d")
# 如果数据不足,拉取
if len(existing) < years * 4:
print(f"数据不足({len(existing)}条),拉取 {start_date} 至 {end_date} 的数据...")
try:
count = fetch_and_save(
ts_code=ts_code,
start_date=start_date,
end_date=end_date,
report_type=report_type
)
print(f"拉取完成,新增 {count} 条记录")
return True
except Exception as e:
print(f"拉取失败: {e}")
return False
return True
def get_field_value(
ts_code: str,
field_name: str,
end_date: str = None,
report_type: str = "1",
auto_fetch: bool = True
) -> Optional[float]:
"""
查询特定字段值,支持自动拉取
Args:
ts_code: 股票代码
field_name: 字段名(如'inventories')
end_date: 报告期,默认最新
report_type: 报表类型
auto_fetch: 是否自动拉取数据(如果数据库无数据)
Returns:
字段值(float),不存在返回None
"""
if auto_fetch:
ensure_data(ts_code, end_date, years=1, report_type=report_type)
record = get_balancesheet(ts_code, end_date, report_type)
if not record:
return None
# 先查基础列
if field_name in record:
value = record[field_name]
if value is not None:
return float(value)
# 再从payload_json中查找
if 'payload_json' in record:
full_data = json.loads(record['payload_json'])
value = full_data.get(field_name)
if value is not None:
return float(value)
return None
def search_by_field(
field_name: str,
end_date: str,
report_type: str = "1",
comp_type: str = None,
limit: int = 100,
order: str = "DESC"
) -> List[Dict]:
"""
跨公司查询特定字段
Args:
field_name: 字段名
end_date: 报告期
report_type: 报表类型
comp_type: 公司类型过滤
limit: 返回数量限制
order: 排序方向(DESC/ASC)
Returns:
列表,每项包含ts_code和字段值
"""
conn = connect()
try:
# 先检查是否是基础列
base_cols = [
'total_share', 'money_cap', 'accounts_receiv', 'inventories',
'total_cur_assets', 'fix_assets', 'total_assets', 'st_borr',
'acct_payable', 'total_cur_liab', 'total_liab', 'undistr_porfit',
'total_hldr_eqy_exc_min_int'
]
conditions = ["end_date = ?", "report_type = ?"]
params = [end_date, report_type]
if comp_type:
conditions.append("comp_type = ?")
params.append(comp_type)
if field_name in base_cols:
# 直接从基础列查询
sql = f"""
SELECT ts_code, {field_name} as value
FROM balancesheet
WHERE {' AND '.join(conditions)} AND {field_name} IS NOT NULL
ORDER BY {field_name} {order}
LIMIT ?
"""
params.append(limit)
cur = conn.execute(sql, params)
else:
# 需要解析JSON(性能较差,但支持所有字段)
sql = f"""
SELECT ts_code, payload_json
FROM balancesheet
WHERE {' AND '.join(conditions)}
LIMIT ?
"""
params.append(limit * 2) # 多查一些,因为可能有些记录没有该字段
cur = conn.execute(sql, params)
# 解析并过滤
results = []
for row in cur.fetchall():
try:
data = json.loads(row['payload_json'])
value = data.get(field_name)
if value is not None:
results.append({
'ts_code': row['ts_code'],
'value': float(value)
})
except:
continue
# 排序
results.sort(key=lambda x: x['value'], reverse=(order == "DESC"))
return results[:limit]
return [dict(row) for row in cur.fetchall()]
finally:
conn.close()
"""券商金股数据库工具"""
import sqlite3
from pathlib import Path
from config import DB_PATH
SCHEMA_PATH = Path(__file__).parent.parent / "schema.sql"
def connect() -> sqlite3.Connection:
conn = sqlite3.connect(str(DB_PATH))
conn.row_factory = sqlite3.Row
return conn
def init_table():
"""初始化券商金股表"""
conn = connect()
try:
schema = SCHEMA_PATH.read_text(encoding="utf-8")
conn.executescript(schema)
conn.commit()
finally:
conn.close()
def upsert_recommendations(rows: list[dict]) -> int:
"""批量插入金股数据,返回插入条数"""
if not rows:
return 0
conn = connect()
try:
sql = """
INSERT INTO broker_recommend (month, broker, ts_code, name)
VALUES (?, ?, ?, ?)
ON CONFLICT(month, broker, ts_code) DO UPDATE SET
name = excluded.name
"""
values = [
(r.get("month"), r.get("broker"), r.get("ts_code"), r.get("name"))
for r in rows
]
conn.executemany(sql, values)
conn.commit()
return len(rows)
finally:
conn.close()
if __name__ == "__main__":
import argparse
parser = argparse.ArgumentParser(description="券商金股数据库工具")
parser.add_argument("--init", action="store_true", help="初始化表")
args = parser.parse_args()
if args.init:
init_table()
print(f"表已初始化: {DB_PATH}")
else:
parser.print_help()
# analyzers/broker_recommend/pipeline.py
"""券商金股分析流程入口"""
import json
import sys
from datetime import datetime
from pathlib import Path
PROJECT_ROOT = Path(__file__).resolve().parents[2]
sys.path.insert(0, str(PROJECT_ROOT))
from fetchers.broker_recommend import fetch_broker_recommend
from analyzers.broker_recommend.stats import get_monthly_stats
def run_fetch(month: str = None) -> int:
"""
拉取数据
Args:
month: 月度 YYYYMM,None 时使用当前月份
Returns:
新增条数
"""
if not month:
month = datetime.now().strftime("%Y%m")
return fetch_broker_recommend(month)
def run_stats(month: str = None) -> dict:
"""
运行统计
Args:
month: 月度 YYYYMM,None 时使用当前月份
Returns:
统计结果
"""
if not month:
month = datetime.now().strftime("%Y%m")
return get_monthly_stats(month)
if __name__ == "__main__":
import argparse
parser = argparse.ArgumentParser(description="券商金股分析流程")
subparsers = parser.add_subparsers(dest="command")
# 拉取
fetch_p = subparsers.add_parser("fetch", help="拉取金股数据")
fetch_p.add_argument("--month", help="月度 YYYYMM,默认当前月份")
# 统计
stats_p = subparsers.add_parser("stats", help="查看统计")
stats_p.add_argument("--month", help="月度 YYYYMM,默认当前月份")
args = parser.parse_args()
if args.command == "fetch":
count = run_fetch(args.month)
print(f"已入库 {count} 条")
elif args.command == "stats":
result = run_stats(args.month)
print(json.dumps(result, ensure_ascii=False, indent=2))
else:
parser.print_help()
-- analyzers/broker_recommend/schema.sql
CREATE TABLE IF NOT EXISTS broker_recommend (
id INTEGER PRIMARY KEY AUTOINCREMENT,
month TEXT NOT NULL, -- 月度 YYYYMM
broker TEXT NOT NULL, -- 券商名称
ts_code TEXT NOT NULL, -- 股票代码
name TEXT NOT NULL, -- 股票简称
created_at TEXT DEFAULT CURRENT_TIMESTAMP,
UNIQUE(month, broker, ts_code)
);
CREATE INDEX IF NOT EXISTS idx_br_month ON broker_recommend(month);
CREATE INDEX IF NOT EXISTS idx_br_ts_code ON broker_recommend(ts_code);
CREATE INDEX IF NOT EXISTS idx_br_broker ON broker_recommend(broker);
# analyzers/broker_recommend/stats.py
"""券商金股统计功能"""
from collections import defaultdict
from ._shared.db import connect
def get_monthly_stats(month: str) -> dict:
"""
获取指定月份的统计
Returns:
{
"month": "202602",
"by_company": [{"ts_code": "...", "name": "...", "count": 5, "brokers": [...]}],
"by_broker": [{"broker": "...", "count": 10}]
}
"""
conn = connect()
try:
cur = conn.execute(
"SELECT ts_code, name, broker FROM broker_recommend WHERE month = ?",
(month,)
)
rows = [dict(r) for r in cur.fetchall()]
finally:
conn.close()
# 按公司统计
company_stats = defaultdict(lambda: {"ts_code": "", "name": "", "count": 0, "brokers": []})
for row in rows:
key = row["ts_code"]
company_stats[key]["ts_code"] = row["ts_code"]
company_stats[key]["name"] = row["name"]
company_stats[key]["count"] += 1
company_stats[key]["brokers"].append(row["broker"])
by_company = sorted(
[
{
"ts_code": v["ts_code"],
"name": v["name"],
"count": v["count"],
"brokers": list(set(v["brokers"]))
}
for v in company_stats.values()
],
key=lambda x: x["count"],
reverse=True
)
# 按券商统计
broker_stats = defaultdict(int)
for row in rows:
broker_stats[row["broker"]] += 1
by_broker = sorted(
[{"broker": k, "count": v} for k, v in broker_stats.items()],
key=lambda x: x["count"],
reverse=True
)
return {
"month": month,
"by_company": by_company,
"by_broker": by_broker
}
def get_continuous_recommendations(ts_code: str, min_months: int = 2) -> list[str]:
"""
获取连续被推荐的月份列表
Args:
ts_code: 股票代码
min_months: 最少连续月份数
Returns:
连续月份列表,如 ["202501", "202502", "202503"]
"""
conn = connect()
try:
cur = conn.execute(
"SELECT DISTINCT month FROM broker_recommend WHERE ts_code = ? ORDER BY month",
(ts_code,)
)
months = [r["month"] for r in cur.fetchall()]
finally:
conn.close()
if len(months) < min_months:
return []
# 找出连续月份(简化版:按字符串排序后检查是否连续)
continuous = []
current_seq = [months[0]]
for i in range(1, len(months)):
prev = int(months[i-1])
curr = int(months[i])
# 计算月份差
prev_year = prev // 100
prev_month = prev % 100
curr_year = curr // 100
curr_month = curr % 100
# 判断是否连续
if curr_year == prev_year and curr_month == prev_month + 1:
# 同一年,月份+1
current_seq.append(months[i])
elif curr_year == prev_year + 1 and curr_month == 1 and prev_month == 12:
# 跨年:12月 -> 1月
current_seq.append(months[i])
else:
# 不连续
if len(current_seq) >= min_months:
continuous.append(current_seq)
current_seq = [months[i]]
if len(current_seq) >= min_months:
continuous.append(current_seq)
# 返回最长的连续序列
return max(continuous, key=len) if continuous else []
def get_broker_recommendations(broker: str, month: str = None) -> list[dict]:
"""
获取某券商的金股列表
Args:
broker: 券商名称
month: 可选,指定月份
Returns:
金股列表
"""
conn = connect()
try:
if month:
cur = conn.execute(
"SELECT month, ts_code, name FROM broker_recommend WHERE broker = ? AND month = ? ORDER BY ts_code",
(broker, month)
)
else:
cur = conn.execute(
"SELECT month, ts_code, name FROM broker_recommend WHERE broker = ? ORDER BY month DESC, ts_code",
(broker,)
)
return [dict(r) for r in cur.fetchall()]
finally:
conn.close()
"""现金流量表数据库工具"""
import json
import sqlite3
from datetime import datetime
from pathlib import Path
from config import DB_PATH
SCHEMA_PATH = Path(__file__).parent.parent / "schema.sql"
def connect(db_path: Path = DB_PATH) -> sqlite3.Connection:
conn = sqlite3.connect(str(db_path))
conn.row_factory = sqlite3.Row
return conn
def init_table(conn: sqlite3.Connection = None):
if conn is None:
conn = connect()
should_close = True
else:
should_close = False
try:
schema = SCHEMA_PATH.read_text(encoding="utf-8")
conn.executescript(schema)
conn.commit()
finally:
if should_close:
conn.close()
def upsert_cashflow(conn: sqlite3.Connection, rows: list) -> int:
if not rows:
return 0
base_cols = [
"ts_code", "ann_date", "end_date",
"n_cashflow_act", "n_cashflow_inv", "n_cashflow_fin",
"c_cash_equ_beg_period", "c_cash_equ_end_period",
"c_inf_fr_operate_a", "c_outf_operate_a",
]
unique_cols = ["ts_code", "end_date"]
now = datetime.now().strftime("%Y-%m-%d %H:%M:%S")
db_rows = []
for row in rows:
db_row = {col: row.get(col) for col in base_cols}
if not db_row.get("ts_code") or not db_row.get("end_date"):
continue
db_row["payload_json"] = json.dumps(row, ensure_ascii=False, default=str)
db_row["source"] = "tushare"
db_row["created_at"] = now
db_row["updated_at"] = now
db_rows.append(db_row)
if not db_rows:
return 0
all_cols = list(db_rows[0].keys())
placeholders = ", ".join(["?"] * len(all_cols))
columns_sql = ", ".join(all_cols)
conflict_sql = ", ".join(unique_cols)
update_cols = [c for c in all_cols if c not in unique_cols + ["created_at"]]
update_sql = ", ".join([f"{c} = excluded.{c}" for c in update_cols])
sql = (
f"INSERT INTO cashflow ({columns_sql}) "
f"VALUES ({placeholders}) "
f"ON CONFLICT({conflict_sql}) DO UPDATE SET {update_sql}"
)
values = [[r.get(col) for col in all_cols] for r in db_rows]
try:
conn.executemany(sql, values)
conn.commit()
return len(db_rows)
except Exception:
conn.rollback()
raise
if __name__ == "__main__":
import argparse
parser = argparse.ArgumentParser()
parser.add_argument("--init", action="store_true", help="初始化表结构")
args = parser.parse_args()
if args.init:
init_table()
print("表结构初始化完成")
"""现金流量表流程入口"""
import sys
from pathlib import Path
PROJECT_ROOT = Path(__file__).resolve().parents[2]
sys.path.insert(0, str(PROJECT_ROOT))
from analyzers.cashflow.search import get_cashflow, get_cashflow_history, get_field_value, ensure_data
from fetchers.cashflow import fetch_and_save
def run_fetch(**kwargs):
count = fetch_and_save(**kwargs)
print(f"保存 {count} 条记录")
return count
def run_query(ts_code: str, field: str = None, end_date: str = None):
ensure_data(ts_code, end_date, years=4)
if field:
value = get_field_value(ts_code, field, end_date)
print("未找到数据" if value is None else f"{ts_code} {field} = {value}")
else:
record = get_cashflow(ts_code, end_date)
print("未找到数据" if not record else record)
def run_history(ts_code: str, limit: int = 4):
records = get_cashflow_history(ts_code, limit=limit)
for r in records:
print(r.get("end_date"), r.get("n_cashflow_act"), r.get("n_cashflow_inv"))
if __name__ == "__main__":
import argparse
parser = argparse.ArgumentParser(description="现金流量表查询")
subparsers = parser.add_subparsers(dest="command", required=True)
p_fetch = subparsers.add_parser("fetch", help="拉取数据")
p_fetch.add_argument("--ts-code")
p_fetch.add_argument("--ann-date")
p_fetch.add_argument("--f-ann-date")
p_fetch.add_argument("--start-date")
p_fetch.add_argument("--end-date")
p_fetch.add_argument("--period")
p_fetch.add_argument("--report-type")
p_fetch.add_argument("--comp-type")
p_fetch.add_argument("--is-calc", type=int)
p_fetch.add_argument("--vip", action="store_true")
p_query = subparsers.add_parser("query", help="查询数据")
p_query.add_argument("--ts-code", required=True)
p_query.add_argument("--field")
p_query.add_argument("--end-date")
p_history = subparsers.add_parser("history", help="查询历史")
p_history.add_argument("--ts-code", required=True)
p_history.add_argument("--limit", type=int, default=4)
args = parser.parse_args()
if args.command == "fetch":
if not any([args.ts_code, args.period, args.start_date, args.end_date]):
parser.error("请至少提供一个查询条件:--ts-code / --period / --start-date / --end-date")
if args.vip and not args.period:
parser.error("使用VIP接口时必须提供 --period 参数")
run_fetch(
ts_code=args.ts_code,
ann_date=args.ann_date,
f_ann_date=args.f_ann_date,
start_date=args.start_date,
end_date=args.end_date,
period=args.period,
report_type=args.report_type,
comp_type=args.comp_type,
is_calc=args.is_calc,
use_vip=args.vip,
)
elif args.command == "query":
run_query(args.ts_code, args.field, args.end_date)
elif args.command == "history":
run_history(args.ts_code, args.limit)
PRAGMA journal_mode=WAL;
CREATE TABLE IF NOT EXISTS cashflow (
ts_code TEXT NOT NULL,
ann_date TEXT,
end_date TEXT NOT NULL,
n_cashflow_act REAL,
n_cashflow_inv REAL,
n_cashflow_fin REAL,
c_cash_equ_beg_period REAL,
c_cash_equ_end_period REAL,
c_inf_fr_operate_a REAL,
c_outf_operate_a REAL,
source TEXT,
payload_json TEXT NOT NULL,
created_at TEXT NOT NULL,
updated_at TEXT NOT NULL,
UNIQUE(ts_code, end_date)
);
CREATE INDEX IF NOT EXISTS idx_cashflow_ts_code ON cashflow(ts_code);
CREATE INDEX IF NOT EXISTS idx_cashflow_end_date ON cashflow(end_date);
CREATE INDEX IF NOT EXISTS idx_cashflow_ts_end ON cashflow(ts_code, end_date);
"""现金流量表查询功能"""
import json
import sys
from pathlib import Path
from typing import Optional, List, Dict
PROJECT_ROOT = Path(__file__).resolve().parents[2]
sys.path.insert(0, str(PROJECT_ROOT))
from analyzers.cashflow._shared.db import connect
def get_cashflow(ts_code: str, end_date: str = None) -> Optional[Dict]:
conn = connect()
try:
if end_date:
sql = "SELECT * FROM cashflow WHERE ts_code = ? AND end_date = ? ORDER BY ann_date DESC LIMIT 1"
cur = conn.execute(sql, (ts_code, end_date))
else:
sql = "SELECT * FROM cashflow WHERE ts_code = ? ORDER BY end_date DESC, ann_date DESC LIMIT 1"
cur = conn.execute(sql, (ts_code,))
row = cur.fetchone()
if not row:
return None
result = dict(row)
if result.get("payload_json"):
try:
result.update(json.loads(result["payload_json"]))
except (json.JSONDecodeError, TypeError):
pass
return result
finally:
conn.close()
def get_cashflow_history(
ts_code: str,
start_date: str = None,
end_date: str = None,
limit: int = None
) -> List[Dict]:
conn = connect()
try:
conditions = ["ts_code = ?"]
params = [ts_code]
if start_date:
conditions.append("end_date >= ?")
params.append(start_date)
if end_date:
conditions.append("end_date <= ?")
params.append(end_date)
sql = f"SELECT * FROM cashflow WHERE {' AND '.join(conditions)} ORDER BY end_date DESC, ann_date DESC"
if limit is not None:
if not isinstance(limit, int) or limit <= 0:
raise ValueError("limit must be a positive int")
sql += " LIMIT ?"
params.append(limit)
cur = conn.execute(sql, params)
rows = cur.fetchall()
results = []
for row in rows:
result = dict(row)
if result.get("payload_json"):
try:
result.update(json.loads(result["payload_json"]))
except (json.JSONDecodeError, TypeError):
pass
results.append(result)
return results
finally:
conn.close()
def get_field_value(ts_code: str, field_name: str, end_date: str = None) -> Optional[float]:
record = get_cashflow(ts_code, end_date)
if not record:
return None
value = record.get(field_name)
if value is not None:
return float(value)
return None
def ensure_data(ts_code: str, end_date: str = None, years: int = 4) -> bool:
from datetime import datetime, timedelta
from fetchers.cashflow import fetch_and_save
if end_date is None:
end_date = datetime.now().strftime("%Y%m%d")
existing = get_cashflow_history(ts_code, end_date=end_date, limit=years * 4 + 4)
need_years = years + 1
start_date_dt = datetime.strptime(end_date, "%Y%m%d") - timedelta(days=need_years * 365)
start_date = start_date_dt.strftime("%Y%m%d")
if len(existing) < years * 4:
print(f"数据不足({len(existing)}条),拉取 {start_date} 至 {end_date} 的数据...")
try:
fetch_and_save(ts_code=ts_code, start_date=start_date, end_date=end_date)
return True
except Exception as e:
print(f"拉取失败: {e}")
return False
return True
import os
import subprocess
import sys
from pathlib import Path
from config import DB_PATH
def parse_dates(value: str):
if not value:
return []
parts = [p.strip() for p in value.split(",")]
return [p for p in parts if p]
def run_cmd(cmd):
subprocess.run(cmd, check=True)
def reset_db(db_path: Path):
if db_path.exists():
db_path.unlink()
def build_ingest_cmd(ingest_path: Path, date_str: str, args):
cmd = [sys.executable, str(ingest_path), "--date", date_str]
if args.report_rc_file:
cmd += ["--report-rc-file", args.report_rc_file]
if args.forecast_file:
cmd += ["--forecast-file", args.forecast_file]
if args.express_file:
cmd += ["--express-file", args.express_file]
if args.income_file:
cmd += ["--income-file", args.income_file]
if args.disclosure_file:
cmd += ["--disclosure-file", args.disclosure_file]
if args.only_report_rc:
cmd.append("--only-report-rc")
if args.no_vip:
cmd.append("--no-vip")
if args.env_path:
cmd += ["--env-path", args.env_path]
return cmd
def build_simple_cmd(script_path: Path):
return [sys.executable, str(script_path)]
def resolve_db_path(root_dir: Path = None):
"""返回数据库路径(使用集中配置)"""
return DB_PATH
def normalize_path(value: str):
if not value:
return value
return os.path.abspath(value)
#!/usr/bin/env python3
"""
对比实际业绩与市场预期,生成告警
"""
import argparse
import json
from datetime import datetime
import pandas as pd
from db import connect, run_migrations, upsert_rows
def now_ts():
return datetime.now().strftime("%Y-%m-%d %H:%M:%S")
def safe_ratio(actual, expected):
if expected is None:
return None
if expected == 0:
return None
return (actual - expected) / abs(expected)
def load_table(conn, sql):
return pd.read_sql_query(sql, conn)
def load_table_with_columns(conn, table, columns):
info = pd.read_sql_query(f"PRAGMA table_info({table})", conn)
available = set(info["name"].tolist())
selected = [col for col in columns if col in available]
if not selected:
return pd.DataFrame(columns=columns)
sql = f"SELECT {', '.join(selected)} FROM {table}"
df = pd.read_sql_query(sql, conn)
for col in columns:
if col not in df.columns:
df[col] = None
return df
def build_expectations(conn, actuals):
df = load_table(
conn,
"""
SELECT ts_code, period, report_date, np
FROM report_rc
WHERE np IS NOT NULL
"""
)
if df.empty or actuals.empty:
return pd.DataFrame()
df["period"] = df["period"].fillna(df["report_date"])
df = df.dropna(subset=["period", "report_date"])
df["period"] = df["period"].astype(str).apply(normalize_quarter_label)
df["report_date"] = df["report_date"].astype(str)
actuals = actuals.copy()
actuals["period"] = actuals["period"].astype(str).apply(normalize_quarter_label)
actuals["ann_date"] = actuals["ann_date"].astype(str)
merged = actuals.merge(
df,
on=["ts_code", "period"],
how="left"
)
merged["report_date_dt"] = pd.to_datetime(merged["report_date"], errors="coerce")
merged["ann_date_dt"] = pd.to_datetime(merged["ann_date"], errors="coerce")
merged = merged.dropna(subset=["report_date_dt", "ann_date_dt"])
merged = merged[merged["report_date_dt"] < merged["ann_date_dt"]]
merged = merged[merged["report_date_dt"] >= (merged["ann_date_dt"] - pd.Timedelta(days=365))]
if merged.empty:
return pd.DataFrame()
agg = merged.groupby(["ts_code", "period", "ann_date"]).agg(
expected_min=("np", "min"),
expected_mean=("np", "mean"),
expected_median=("np", "median"),
expected_max=("np", "max"),
expected_report_date_min=("report_date", "min"),
expected_report_date_max=("report_date", "max"),
).reset_index()
return agg
def latest_by_group(df, group_cols, order_col):
if df.empty:
return df
df = df.sort_values(order_col, ascending=False)
return df.groupby(group_cols, as_index=False).first()
def annualize_factor(end_date: str):
if not end_date:
return 1.0
end_date = str(end_date)
suffix = end_date[-4:]
if suffix == "0331":
return 4.0
if suffix == "0630":
return 2.0
if suffix == "0930":
return 4.0 / 3.0
return 1.0
def end_date_to_annual_period(end_date: str):
if not end_date:
return None
end_date = str(end_date)
if len(end_date) < 4:
return None
year = end_date[:4]
if not year.isdigit():
return None
return f"{year}Q4"
def normalize_quarter_label(value: str):
if value is None:
return value
value = str(value)
if "Q" in value:
year = value[:4]
quarter = value[-1]
if year.isdigit() and quarter in {"1", "2", "3", "4"}:
return f"{year}Q{quarter}"
if len(value) == 8 and value.isdigit():
year = value[:4]
suffix = value[-4:]
mapping = {"0331": "1", "0630": "2", "0930": "3", "1231": "4"}
if suffix in mapping:
return f"{year}Q{mapping[suffix]}"
return value
def build_actuals(conn):
# forecast: use avg(net_profit_min, net_profit_max)
forecast = load_table_with_columns(
conn,
"forecast",
["ts_code", "end_date", "ann_date", "net_profit_min", "net_profit_max", "change_reason", "updated_at"]
)
if not forecast.empty:
forecast["actual_value"] = (
forecast["net_profit_min"].fillna(0) + forecast["net_profit_max"].fillna(0)
) / 2
forecast = latest_by_group(forecast, ["ts_code", "end_date"], "updated_at")
forecast["source"] = "forecast"
forecast["reason_hint"] = forecast["change_reason"]
# express: n_income(原单位:元,转换为万元)
express = load_table_with_columns(
conn,
"express",
["ts_code", "end_date", "ann_date", "n_income", "perf_summary", "updated_at"]
)
if not express.empty:
express = latest_by_group(express, ["ts_code", "end_date"], "updated_at")
express["actual_value"] = express["n_income"] / 10000 # 元 → 万元
express["source"] = "express"
express["reason_hint"] = express["perf_summary"]
# income: n_income(原单位:元,转换为万元)
income = load_table_with_columns(
conn,
"income",
["ts_code", "end_date", "ann_date", "f_ann_date", "n_income", "updated_at"]
)
if not income.empty:
income = latest_by_group(income, ["ts_code", "end_date"], "updated_at")
income["actual_value"] = income["n_income"] / 10000 # 元 → 万元
income["source"] = "income"
income["ann_date"] = income["f_ann_date"].fillna(income["ann_date"])
income["reason_hint"] = None
# priority: forecast > express > income
frames = []
if not forecast.empty:
frames.append(forecast[["ts_code", "end_date", "ann_date", "actual_value", "source", "reason_hint"]])
if not express.empty:
frames.append(express[["ts_code", "end_date", "ann_date", "actual_value", "source", "reason_hint"]])
if not income.empty:
frames.append(income[["ts_code", "end_date", "ann_date", "actual_value", "source", "reason_hint"]])
if not frames:
return pd.DataFrame()
combined = pd.concat(frames, ignore_index=True)
priority = {"forecast": 1, "express": 2, "income": 3}
combined["priority"] = combined["source"].map(priority)
combined = combined.sort_values("priority")
combined = combined.groupby(["ts_code", "end_date"], as_index=False).first()
combined["annualize_factor"] = combined["end_date"].apply(annualize_factor)
combined["actual_annualized"] = combined["actual_value"] * combined["annualize_factor"]
combined["period"] = combined["end_date"].apply(end_date_to_annual_period)
return combined.drop(columns=["priority"])
def build_alerts(expectations, actuals):
if expectations.empty or actuals.empty:
return pd.DataFrame()
merged = actuals.merge(expectations, on=["ts_code", "period", "ann_date"], how="left")
merged = merged.dropna(subset=["expected_max", "actual_annualized"])
# 超预期:actual > expected_max
above = merged[merged["actual_annualized"] > merged["expected_max"]].copy()
if not above.empty:
above["alert_type"] = "above"
# 低于预期:actual < expected_mean
below = merged[merged["actual_annualized"] < merged["expected_mean"]].copy()
if not below.empty:
below["alert_type"] = "below"
# 符合预期:expected_mean <= actual <= expected_max
inline = merged[
(merged["actual_annualized"] >= merged["expected_mean"]) &
(merged["actual_annualized"] <= merged["expected_max"])
].copy()
if not inline.empty:
inline["alert_type"] = "inline"
frames = [df for df in [above, below, inline] if not df.empty]
alerts = pd.concat(frames, ignore_index=True) if frames else pd.DataFrame()
if alerts.empty:
return alerts
alerts["delta_mean"] = alerts.apply(
lambda r: safe_ratio(r["actual_annualized"], r["expected_mean"]), axis=1
)
alerts["delta_median"] = alerts.apply(
lambda r: safe_ratio(r["actual_annualized"], r["expected_median"]), axis=1
)
alerts["delta_max"] = alerts.apply(
lambda r: safe_ratio(r["actual_annualized"], r["expected_max"]), axis=1
)
return alerts
def main():
parser = argparse.ArgumentParser(description="对比并生成告警")
parser.add_argument("--dry-run", action="store_true", help="只统计不入库")
args = parser.parse_args()
conn = connect()
run_migrations(conn)
actuals = build_actuals(conn)
expectations = build_expectations(conn, actuals)
alerts = build_alerts(expectations, actuals)
if not alerts.empty and not args.dry_run:
rows = []
ts_now = now_ts()
for _, row in alerts.iterrows():
payload = {
"alert_type": row.get("alert_type"), # above/below/inline
"expected_min": row.get("expected_min"),
"expected_mean": row.get("expected_mean"),
"expected_median": row.get("expected_median"),
"expected_max": row.get("expected_max"),
"expected_report_date_min": row.get("expected_report_date_min"),
"expected_report_date_max": row.get("expected_report_date_max"),
"ann_date": row.get("ann_date"),
"period": row.get("period"),
"delta_mean": row.get("delta_mean"),
"delta_median": row.get("delta_median"),
"delta_max": row.get("delta_max"),
"actual_original": row.get("actual_value"),
"annualize_factor": row.get("annualize_factor"),
"reason_hint": row.get("reason_hint"),
}
rows.append(
{
"ts_code": row["ts_code"],
"end_date": row["end_date"],
"actual_value": row["actual_annualized"],
"expected_mean": row.get("expected_mean"),
"expected_median": row.get("expected_median"),
"expected_max": row.get("expected_max"),
"delta_mean": row.get("delta_mean"),
"delta_median": row.get("delta_median"),
"delta_max": row.get("delta_max"),
"source": row["source"],
"created_at": ts_now,
"payload_json": json.dumps(payload, ensure_ascii=True, sort_keys=True),
}
)
upsert_rows(conn, "alerts", rows, ["ts_code", "end_date", "source", "actual_value"], preserve_on_update=["created_at"])
conn.close()
print(f"expectations: {len(expectations)}")
print(f"actuals: {len(actuals)}")
print(f"alerts: {len(alerts)}")
if __name__ == "__main__":
main()
#!/usr/bin/env python3
"""
SQLite 工具函数
"""
import argparse
import sqlite3
from pathlib import Path
from config import DB_PATH
SCHEMA_PATH = Path(__file__).with_name("schema.sql")
def connect(db_path: Path = DB_PATH) -> sqlite3.Connection:
"""连接数据库"""
conn = sqlite3.connect(str(db_path))
conn.row_factory = sqlite3.Row
return conn
def run_migrations(conn: sqlite3.Connection):
"""执行建表脚本"""
schema = SCHEMA_PATH.read_text(encoding="utf-8")
conn.executescript(schema)
conn.commit()
def fetch_existing(conn: sqlite3.Connection, table: str, unique_cols: list, row: dict):
"""查询已存在记录"""
where = " AND ".join([f"{col} = ?" for col in unique_cols])
sql = f"SELECT * FROM {table} WHERE {where} LIMIT 1"
values = [row.get(col) for col in unique_cols]
cur = conn.execute(sql, values)
result = cur.fetchone()
return dict(result) if result else None
def upsert_rows(conn: sqlite3.Connection, table: str, rows: list, unique_cols: list, preserve_on_update: list = None) -> int:
"""批量 upsert
Args:
preserve_on_update: 这些字段在更新时保留原值(仅 INSERT 时写入)
"""
if not rows:
return 0
preserve_on_update = preserve_on_update or []
all_cols = []
for row in rows:
for key in row.keys():
if key not in all_cols:
all_cols.append(key)
placeholders = ", ".join(["?"] * len(all_cols))
columns_sql = ", ".join(all_cols)
conflict_sql = ", ".join(unique_cols)
update_cols = [c for c in all_cols if c not in unique_cols]
# preserve_on_update 的字段保留原值,其他字段用新值覆盖
update_parts = []
for c in update_cols:
if c in preserve_on_update:
update_parts.append(f"{c} = COALESCE({table}.{c}, excluded.{c})")
else:
update_parts.append(f"{c} = excluded.{c}")
update_sql = ", ".join(update_parts)
sql = (
f"INSERT INTO {table} ({columns_sql}) "
f"VALUES ({placeholders}) "
f"ON CONFLICT({conflict_sql}) DO UPDATE SET {update_sql}"
)
values = []
for row in rows:
values.append([row.get(col) for col in all_cols])
conn.executemany(sql, values)
conn.commit()
return len(rows)
def init_db():
"""初始化数据库"""
conn = connect()
run_migrations(conn)
conn.close()
def main():
parser = argparse.ArgumentParser(description="初始化数据库")
parser.add_argument("--init", action="store_true", help="创建数据库与表")
args = parser.parse_args()
if args.init:
init_db()
print(f"数据库已初始化: {DB_PATH}")
else:
parser.print_help()
if __name__ == "__main__":
main()
#!/usr/bin/env python3
"""
数据拉取与入库
"""
import argparse
import json
import os
from datetime import datetime
from pathlib import Path
import pandas as pd
import tushare as ts
from dotenv import load_dotenv
from db import connect, run_migrations, fetch_existing, upsert_rows
def load_env(env_path: str = None):
"""加载 .env"""
if env_path:
load_dotenv(env_path)
return
for parent in Path(__file__).resolve().parents:
candidate = parent / ".env"
if candidate.exists():
load_dotenv(candidate)
return
def init_tushare(env_path: str = None):
"""初始化 Tushare"""
load_env(env_path)
api_key = os.getenv("TUSHARE_API_KEY")
if not api_key:
raise RuntimeError("未找到 TUSHARE_API_KEY")
return ts.pro_api(api_key)
def now_ts():
return datetime.now().strftime("%Y-%m-%d %H:%M:%S")
def get_periods_for_date(date_str: str) -> list:
"""根据日期自动计算需要拉取的 period 列表
规则:
- 1月1日 - 4月30日:上年Q4(1231) + 当年Q1(0331)
- 5月1日 - 8月31日:当年Q2(0630)
- 9月1日 - 10月31日:当年Q3(0930)
- 11月1日 - 12月31日:当年Q4(1231) 预披露
"""
dt = datetime.strptime(date_str, "%Y%m%d")
year = dt.year
month = dt.month
periods = []
if 1 <= month <= 4:
# 年报 + Q1 披露季
periods.append(f"{year - 1}1231") # 上年年报
periods.append(f"{year}0331") # 当年Q1
elif 5 <= month <= 8:
# 中报披露季
periods.append(f"{year}0630")
elif 9 <= month <= 10:
# Q3 披露季
periods.append(f"{year}0930")
else: # 11-12月
# 年报预披露期
periods.append(f"{year}1231")
return periods
def safe_value(value):
if pd.isna(value):
return None
return value
def to_payload(row: dict):
return json.dumps(row, ensure_ascii=True, sort_keys=True)
def build_change_log(existing: dict, new_row: dict, fields: list):
if not existing:
return "new", None
changes = {}
for field in fields:
old = existing.get(field)
new = new_row.get(field)
if old != new:
changes[field] = {"old": old, "new": new}
if not changes:
return "same", None
return "update", json.dumps(changes, ensure_ascii=True, sort_keys=True)
def normalize_rows(df: pd.DataFrame, base_map: dict, source: str, existing_fetcher=None, change_fields=None):
rows = []
ts_now = now_ts()
for _, row in df.iterrows():
row_dict = {k: safe_value(v) for k, v in row.to_dict().items()}
new_row = {col: safe_value(row_dict.get(src)) for col, src in base_map.items()}
new_row["source"] = source
new_row["payload_json"] = to_payload(row_dict)
new_row["updated_at"] = ts_now
new_row["created_at"] = ts_now
if existing_fetcher and change_fields:
existing = existing_fetcher(new_row)
change_type, change_log = build_change_log(existing, new_row, change_fields)
new_row["change_type"] = change_type
new_row["change_log"] = change_log
rows.append(new_row)
return rows
def fetch_report_rc_full(pro, report_date: str):
return pro.report_rc(report_date=report_date)
def fetch_forecast_by_periods(pro, periods: list, use_vip: bool):
"""按 period 列表拉取业绩预告"""
dfs = []
for period in periods:
if use_vip:
df = pro.forecast_vip(period=period)
else:
df = pro.forecast(period=period)
if df is not None and not df.empty:
print(f" forecast period={period}: {len(df)} 条")
dfs.append(df)
if not dfs:
return pd.DataFrame()
return pd.concat(dfs, ignore_index=True).drop_duplicates()
def fetch_express_by_periods(pro, periods: list, use_vip: bool):
"""按 period 列表拉取业绩快报"""
dfs = []
for period in periods:
if use_vip:
df = pro.express_vip(period=period)
else:
df = pro.express(period=period)
if df is not None and not df.empty:
print(f" express period={period}: {len(df)} 条")
dfs.append(df)
if not dfs:
return pd.DataFrame()
return pd.concat(dfs, ignore_index=True).drop_duplicates()
def fetch_income_by_periods(pro, periods: list, use_vip: bool):
"""按 period 列表拉取正式业绩"""
dfs = []
for period in periods:
if use_vip:
df = pro.income_vip(period=period)
else:
df = pro.income(period=period)
if df is not None and not df.empty:
print(f" income period={period}: {len(df)} 条")
dfs.append(df)
if not dfs:
return pd.DataFrame()
return pd.concat(dfs, ignore_index=True).drop_duplicates()
def fetch_disclosure_full(pro, date_str: str):
dfs = []
for key in ["ann_date", "pre_date", "actual_date"]:
df = pro.disclosure_date(**{key: date_str})
if df is not None and not df.empty:
dfs.append(df)
if not dfs:
return pd.DataFrame()
return pd.concat(dfs, ignore_index=True).drop_duplicates()
def main():
parser = argparse.ArgumentParser(description="数据拉取与入库")
parser.add_argument("--date", help="日期 (YYYYMMDD),默认今天")
parser.add_argument("--dry-run", action="store_true", help="只统计不入库")
parser.add_argument("--no-vip", action="store_true", help="不使用 VIP 接口")
parser.add_argument("--env-path", help="指定 .env 路径")
parser.add_argument("--report-rc-file", help="使用本地 report_rc CSV 文件进行模拟")
parser.add_argument("--forecast-file", help="使用本地 forecast CSV 文件进行模拟")
parser.add_argument("--express-file", help="使用本地 express CSV 文件进行模拟")
parser.add_argument("--income-file", help="使用本地 income CSV 文件进行模拟")
parser.add_argument("--disclosure-file", help="使用本地 disclosure_date CSV 文件进行模拟")
parser.add_argument("--only-report-rc", action="store_true", help="只处理 report_rc")
args = parser.parse_args()
date_str = args.date or datetime.now().strftime("%Y%m%d")
use_vip = not args.no_vip
# 根据日期自动计算需要拉取的 period
periods = get_periods_for_date(date_str)
print(f"日期: {date_str}")
print(f"自动 period: {periods}")
need_api = True
if args.only_report_rc and args.report_rc_file:
need_api = False
if args.report_rc_file and args.forecast_file and args.express_file and args.income_file and args.disclosure_file:
need_api = False
pro = None if not need_api else init_tushare(args.env_path)
conn = connect()
run_migrations(conn)
def existing_fetcher_factory(table, unique_cols):
def _fetch(row):
return fetch_existing(conn, table, unique_cols, row)
return _fetch
counts = {}
# report_rc
if args.report_rc_file:
df_report = pd.read_csv(args.report_rc_file)
if "report_date" in df_report.columns:
df_report = df_report[df_report["report_date"].astype(str) == date_str]
else:
df_report = fetch_report_rc_full(pro, date_str)
report_base = {
"ts_code": "ts_code",
"report_date": "report_date",
"org_name": "org_name",
"period": "quarter",
"np": "np",
"eps": "eps",
"quarter": "quarter",
}
report_rows = normalize_rows(
df_report,
report_base,
"report_rc",
existing_fetcher_factory("report_rc", ["ts_code", "report_date", "org_name", "period"]),
["np"]
)
counts["report_rc"] = len(report_rows)
if not args.dry_run and report_rows:
upsert_rows(conn, "report_rc", report_rows, ["ts_code", "report_date", "org_name", "period"])
if args.only_report_rc:
df_forecast = pd.DataFrame()
elif args.forecast_file:
df_forecast = pd.read_csv(args.forecast_file)
if "ann_date" in df_forecast.columns:
df_forecast = df_forecast[df_forecast["ann_date"].astype(str) == date_str]
else:
print("拉取 forecast...")
df_forecast = fetch_forecast_by_periods(pro, periods, use_vip)
forecast_base = {
"ts_code": "ts_code",
"ann_date": "ann_date",
"end_date": "end_date",
"net_profit_min": "net_profit_min",
"net_profit_max": "net_profit_max",
"type": "type",
}
forecast_rows = normalize_rows(
df_forecast,
forecast_base,
"forecast",
existing_fetcher_factory("forecast", ["ts_code", "ann_date", "end_date"]),
["net_profit_min", "net_profit_max"]
)
counts["forecast"] = len(forecast_rows)
if not args.dry_run and forecast_rows:
upsert_rows(conn, "forecast", forecast_rows, ["ts_code", "ann_date", "end_date"])
if args.only_report_rc:
df_express = pd.DataFrame()
elif args.express_file:
df_express = pd.read_csv(args.express_file)
if "ann_date" in df_express.columns:
df_express = df_express[df_express["ann_date"].astype(str) == date_str]
else:
print("拉取 express...")
df_express = fetch_express_by_periods(pro, periods, use_vip)
express_base = {
"ts_code": "ts_code",
"ann_date": "ann_date",
"end_date": "end_date",
"n_income": "n_income",
}
express_rows = normalize_rows(
df_express,
express_base,
"express",
existing_fetcher_factory("express", ["ts_code", "ann_date", "end_date"]),
["n_income"]
)
counts["express"] = len(express_rows)
if not args.dry_run and express_rows:
upsert_rows(conn, "express", express_rows, ["ts_code", "ann_date", "end_date"])
if args.only_report_rc:
df_income = pd.DataFrame()
elif args.income_file:
df_income = pd.read_csv(args.income_file)
if "ann_date" in df_income.columns:
df_income = df_income[df_income["ann_date"].astype(str) == date_str]
else:
print("拉取 income...")
df_income = fetch_income_by_periods(pro, periods, use_vip)
income_base = {
"ts_code": "ts_code",
"ann_date": "ann_date",
"end_date": "end_date",
"n_income": "n_income",
"report_type": "report_type",
"comp_type": "comp_type",
}
income_rows = normalize_rows(
df_income,
income_base,
"income",
existing_fetcher_factory("income", ["ts_code", "ann_date", "end_date"]),
["n_income"]
)
counts["income"] = len(income_rows)
if not args.dry_run and income_rows:
upsert_rows(conn, "income", income_rows, ["ts_code", "ann_date", "end_date"])
if args.only_report_rc:
df_disclosure = pd.DataFrame()
elif args.disclosure_file:
df_disclosure = pd.read_csv(args.disclosure_file)
date_cols = [c for c in ["ann_date", "pre_date", "actual_date"] if c in df_disclosure.columns]
if date_cols:
mask = False
for col in date_cols:
mask = mask | (df_disclosure[col].astype(str) == date_str)
df_disclosure = df_disclosure[mask]
else:
df_disclosure = fetch_disclosure_full(pro, date_str)
disclosure_base = {
"ts_code": "ts_code",
"end_date": "end_date",
"pre_date": "pre_date",
"ann_date": "ann_date",
"actual_date": "actual_date",
}
disclosure_rows = normalize_rows(
df_disclosure,
disclosure_base,
"disclosure_date",
existing_fetcher_factory("disclosure_date", ["ts_code", "end_date"]),
["pre_date", "ann_date", "actual_date"]
)
counts["disclosure_date"] = len(disclosure_rows)
if not args.dry_run and disclosure_rows:
upsert_rows(conn, "disclosure_date", disclosure_rows, ["ts_code", "end_date"])
conn.close()
print(f"日期: {date_str}")
print(f"使用VIP: {use_vip}")
for key, value in counts.items():
print(f"{key}: {value}")
if __name__ == "__main__":
main()
#!/usr/bin/env python3
"""
主流程入口:拉取 -> 对比 -> JSON 输出
"""
import argparse
import json
from datetime import datetime
from pathlib import Path
from db import connect, run_migrations
def now_ts():
return datetime.now().strftime("%Y-%m-%d %H:%M:%S")
def load_table(conn, table: str):
cur = conn.execute(f"SELECT * FROM {table}")
rows = cur.fetchall()
return [dict(r) for r in rows]
def change_direction(change_log: str):
if not change_log:
return None
try:
data = json.loads(change_log)
except json.JSONDecodeError:
return None
np_change = data.get("np")
if not np_change:
return None
old = np_change.get("old")
new = np_change.get("new")
if old is None or new is None:
return None
if new > old:
return "up"
if new < old:
return "down"
return "same"
def load_report_rc_changes(conn, since_ts: str):
sql = (
"SELECT ts_code, report_date, period, org_name, np, change_type, change_log, updated_at "
"FROM report_rc WHERE change_type = 'update' AND updated_at >= ?"
)
rows = conn.execute(sql, (since_ts,)).fetchall()
result = []
for row in rows:
row_dict = dict(row)
row_dict["direction"] = change_direction(row_dict.get("change_log"))
result.append(row_dict)
return result
def write_json(output_dir: Path, payload: dict):
output_dir.mkdir(exist_ok=True)
ts = datetime.now().strftime("%Y%m%d_%H%M%S")
filepath = output_dir / f"summary_{ts}.json"
filepath.write_text(json.dumps(payload, ensure_ascii=True, indent=2), encoding="utf-8")
return filepath
def cleanup_old_files(output_dir: Path, pattern: str, keep: int = 5):
"""清理旧文件,只保留最近 N 个"""
files = sorted(output_dir.glob(pattern), key=lambda f: f.stat().st_mtime, reverse=True)
for old_file in files[keep:]:
old_file.unlink()
print(f"已清理: {old_file.name}")
def main():
parser = argparse.ArgumentParser(description="主流程入口")
parser.add_argument("--summary-only", action="store_true", help="仅输出汇总 JSON")
args = parser.parse_args()
conn = connect()
run_migrations(conn)
alerts = load_table(conn, "alerts")
since_ts = datetime.now().strftime("%Y-%m-%d 00:00:00")
expectation_changes = load_report_rc_changes(conn, since_ts)
payload = {
"run_meta": {
"run_at": now_ts(),
"alerts_count": len(alerts),
},
"summary": {
"alerts": len(alerts),
},
"alerts": alerts,
"changes": expectation_changes,
}
output_dir = Path(__file__).resolve().parents[2] / "output"
filepath = write_json(output_dir, payload)
conn.close()
# 清理旧文件,只保留最近 5 个
cleanup_old_files(output_dir, "summary_*.json", keep=5)
print(f"已输出: {filepath}")
if __name__ == "__main__":
main()
PRAGMA journal_mode=WAL;
CREATE TABLE IF NOT EXISTS report_rc (
ts_code TEXT NOT NULL,
report_date TEXT NOT NULL,
org_name TEXT NOT NULL,
period TEXT,
np REAL,
eps REAL,
quarter TEXT,
source TEXT,
payload_json TEXT NOT NULL,
created_at TEXT NOT NULL,
updated_at TEXT NOT NULL,
change_type TEXT,
change_log TEXT,
UNIQUE (ts_code, report_date, org_name, period)
);
CREATE TABLE IF NOT EXISTS forecast (
ts_code TEXT NOT NULL,
ann_date TEXT,
end_date TEXT NOT NULL,
net_profit_min REAL,
net_profit_max REAL,
type TEXT,
source TEXT,
payload_json TEXT NOT NULL,
created_at TEXT NOT NULL,
updated_at TEXT NOT NULL,
change_type TEXT,
change_log TEXT,
UNIQUE (ts_code, ann_date, end_date)
);
CREATE TABLE IF NOT EXISTS express (
ts_code TEXT NOT NULL,
ann_date TEXT,
end_date TEXT NOT NULL,
n_income REAL,
source TEXT,
payload_json TEXT NOT NULL,
created_at TEXT NOT NULL,
updated_at TEXT NOT NULL,
change_type TEXT,
change_log TEXT,
UNIQUE (ts_code, ann_date, end_date)
);
CREATE TABLE IF NOT EXISTS income (
ts_code TEXT NOT NULL,
ann_date TEXT,
end_date TEXT NOT NULL,
n_income REAL,
report_type TEXT,
comp_type TEXT,
source TEXT,
payload_json TEXT NOT NULL,
created_at TEXT NOT NULL,
updated_at TEXT NOT NULL,
change_type TEXT,
change_log TEXT,
UNIQUE (ts_code, ann_date, end_date)
);
CREATE TABLE IF NOT EXISTS disclosure_date (
ts_code TEXT NOT NULL,
end_date TEXT NOT NULL,
pre_date TEXT,
ann_date TEXT,
actual_date TEXT,
source TEXT,
payload_json TEXT NOT NULL,
created_at TEXT NOT NULL,
updated_at TEXT NOT NULL,
change_type TEXT,
change_log TEXT,
UNIQUE (ts_code, end_date)
);
CREATE TABLE IF NOT EXISTS alerts (
id INTEGER PRIMARY KEY AUTOINCREMENT,
ts_code TEXT NOT NULL,
end_date TEXT NOT NULL,
actual_value REAL NOT NULL,
expected_mean REAL,
expected_median REAL,
expected_max REAL,
delta_mean REAL,
delta_median REAL,
delta_max REAL,
source TEXT NOT NULL,
created_at TEXT NOT NULL,
payload_json TEXT NOT NULL,
UNIQUE (ts_code, end_date, source, actual_value)
);
CREATE TABLE IF NOT EXISTS job_runs (
id INTEGER PRIMARY KEY AUTOINCREMENT,
run_at TEXT NOT NULL,
status TEXT NOT NULL,
error TEXT,
meta_json TEXT
);
CREATE INDEX IF NOT EXISTS idx_report_rc_ts_code ON report_rc (ts_code);
CREATE INDEX IF NOT EXISTS idx_report_rc_period ON report_rc (period);
CREATE INDEX IF NOT EXISTS idx_forecast_ts_code ON forecast (ts_code);
CREATE INDEX IF NOT EXISTS idx_express_ts_code ON express (ts_code);
CREATE INDEX IF NOT EXISTS idx_income_ts_code ON income (ts_code);
CREATE INDEX IF NOT EXISTS idx_disclosure_ts_code ON disclosure_date (ts_code);
CREATE INDEX IF NOT EXISTS idx_alerts_ts_code ON alerts (ts_code);
#!/usr/bin/env python3
"""
本地模拟流程:按日期批量 ingest -> compare -> run
"""
import argparse
from pathlib import Path
from _shared.sim_runner import (
build_ingest_cmd,
build_simple_cmd,
normalize_path,
parse_dates,
reset_db,
resolve_db_path,
run_cmd,
)
def main():
parser = argparse.ArgumentParser(description="本地模拟流程")
parser.add_argument("--dates", required=True, help="逗号分隔日期列表 (YYYYMMDD,YYYYMMDD)")
parser.add_argument("--report-rc-file", required=True, help="report_rc CSV")
parser.add_argument("--report-rc-update-file", help="report_rc 更新样本 CSV")
parser.add_argument("--update-date", help="更新样本日期 (YYYYMMDD)")
parser.add_argument("--forecast-file", help="forecast CSV")
parser.add_argument("--express-file", help="express CSV")
parser.add_argument("--income-file", help="income CSV")
parser.add_argument("--disclosure-file", help="disclosure_date CSV")
parser.add_argument("--reset-db", action="store_true", help="重建数据库")
parser.add_argument("--only-report-rc", action="store_true", help="只处理 report_rc")
parser.add_argument("--no-vip", action="store_true", help="不使用 VIP 接口")
parser.add_argument("--env-path", help="指定 .env 路径")
args = parser.parse_args()
root_dir = Path(__file__).resolve().parents[2]
ingest_path = Path(__file__).with_name("ingest.py")
compare_path = Path(__file__).with_name("compare.py")
run_path = Path(__file__).with_name("run.py")
args.report_rc_file = normalize_path(args.report_rc_file)
args.report_rc_update_file = normalize_path(args.report_rc_update_file)
args.forecast_file = normalize_path(args.forecast_file)
args.express_file = normalize_path(args.express_file)
args.income_file = normalize_path(args.income_file)
args.disclosure_file = normalize_path(args.disclosure_file)
args.env_path = normalize_path(args.env_path)
date_list = parse_dates(args.dates)
if not date_list:
raise SystemExit("未提供有效 dates")
if args.reset_db:
reset_db(resolve_db_path(root_dir))
run_cmd(build_simple_cmd(Path(__file__).with_name("db.py")) + ["--init"])
for date_str in date_list:
cmd = build_ingest_cmd(ingest_path, date_str, args)
run_cmd(cmd)
if args.report_rc_update_file:
update_date = args.update_date or date_list[-1]
args.report_rc_file = args.report_rc_update_file
cmd = build_ingest_cmd(ingest_path, update_date, args)
run_cmd(cmd)
run_cmd(build_simple_cmd(compare_path))
run_cmd(build_simple_cmd(run_path))
if __name__ == "__main__":
main()
"""财务指标数据库工具"""
import json
import sqlite3
from datetime import datetime
from pathlib import Path
from config import DB_PATH
SCHEMA_PATH = Path(__file__).parent.parent / "schema.sql"
def connect(db_path: Path = DB_PATH) -> sqlite3.Connection:
conn = sqlite3.connect(str(db_path))
conn.row_factory = sqlite3.Row
return conn
def init_table(conn: sqlite3.Connection = None):
if conn is None:
conn = connect()
should_close = True
else:
should_close = False
try:
schema = SCHEMA_PATH.read_text(encoding="utf-8")
conn.executescript(schema)
conn.commit()
finally:
if should_close:
conn.close()
def upsert_fina_indicator(conn: sqlite3.Connection, rows: list) -> int:
if not rows:
return 0
base_cols = [
"ts_code", "ann_date", "end_date",
"roe", "roa", "grossprofit_margin", "netprofit_margin",
"current_ratio", "quick_ratio", "debt_to_assets",
"assets_turn", "inv_turn", "ar_turn",
"eps", "profit_to_gr",
]
unique_cols = ["ts_code", "end_date"]
now = datetime.now().strftime("%Y-%m-%d %H:%M:%S")
db_rows = []
for row in rows:
db_row = {col: row.get(col) for col in base_cols}
if not db_row.get("ts_code") or not db_row.get("end_date"):
continue
db_row["payload_json"] = json.dumps(row, ensure_ascii=False, default=str)
db_row["source"] = "tushare"
db_row["created_at"] = now
db_row["updated_at"] = now
db_rows.append(db_row)
if not db_rows:
return 0
all_cols = list(db_rows[0].keys())
placeholders = ", ".join(["?"] * len(all_cols))
columns_sql = ", ".join(all_cols)
conflict_sql = ", ".join(unique_cols)
update_cols = [c for c in all_cols if c not in unique_cols + ["created_at"]]
update_sql = ", ".join([f"{c} = excluded.{c}" for c in update_cols])
sql = (
f"INSERT INTO fina_indicator ({columns_sql}) "
f"VALUES ({placeholders}) "
f"ON CONFLICT({conflict_sql}) DO UPDATE SET {update_sql}"
)
values = [[r.get(col) for col in all_cols] for r in db_rows]
try:
conn.executemany(sql, values)
conn.commit()
return len(db_rows)
except Exception:
conn.rollback()
raise
if __name__ == "__main__":
import argparse
parser = argparse.ArgumentParser()
parser.add_argument("--init", action="store_true", help="初始化表结构")
args = parser.parse_args()
if args.init:
init_table()
print("表结构初始化完成")
"""财务指标流程入口"""
import sys
from pathlib import Path
PROJECT_ROOT = Path(__file__).resolve().parents[2]
sys.path.insert(0, str(PROJECT_ROOT))
from analyzers.fina_indicator.search import get_fina_indicator, get_fina_indicator_history, get_field_value
from fetchers.fina_indicator import fetch_and_save
def run_fetch(**kwargs):
count = fetch_and_save(**kwargs)
print(f"保存 {count} 条记录")
return count
def run_query(ts_code: str, field: str = None, end_date: str = None):
if field:
value = get_field_value(ts_code, field, end_date)
print("未找到数据" if value is None else f"{ts_code} {field} = {value}")
else:
record = get_fina_indicator(ts_code, end_date)
print("未找到数据" if not record else record)
def run_history(ts_code: str, limit: int = 4):
records = get_fina_indicator_history(ts_code, limit=limit)
for r in records:
print(r.get("end_date"), r.get("roe"), r.get("netprofit_margin"))
if __name__ == "__main__":
import argparse
parser = argparse.ArgumentParser(description="财务指标查询")
subparsers = parser.add_subparsers(dest="command", required=True)
p_fetch = subparsers.add_parser("fetch", help="拉取数据")
p_fetch.add_argument("--ts-code")
p_fetch.add_argument("--start-date")
p_fetch.add_argument("--end-date")
p_fetch.add_argument("--period")
p_fetch.add_argument("--vip", action="store_true")
p_query = subparsers.add_parser("query", help="查询数据")
p_query.add_argument("--ts-code", required=True)
p_query.add_argument("--field")
p_query.add_argument("--end-date")
p_history = subparsers.add_parser("history", help="查询历史")
p_history.add_argument("--ts-code", required=True)
p_history.add_argument("--limit", type=int, default=4)
args = parser.parse_args()
if args.command == "fetch":
if not any([args.ts_code, args.period, args.start_date, args.end_date]):
parser.error("请至少提供一个查询条件:--ts-code / --period / --start-date / --end-date")
if args.vip and not args.period:
parser.error("使用VIP接口时必须提供 --period 参数")
run_fetch(
ts_code=args.ts_code,
start_date=args.start_date,
end_date=args.end_date,
period=args.period,
use_vip=args.vip,
)
elif args.command == "query":
run_query(args.ts_code, args.field, args.end_date)
elif args.command == "history":
run_history(args.ts_code, args.limit)
PRAGMA journal_mode=WAL;
CREATE TABLE IF NOT EXISTS fina_indicator (
ts_code TEXT NOT NULL,
ann_date TEXT,
end_date TEXT NOT NULL,
roe REAL,
roa REAL,
grossprofit_margin REAL,
netprofit_margin REAL,
current_ratio REAL,
quick_ratio REAL,
debt_to_assets REAL,
assets_turn REAL,
inv_turn REAL,
ar_turn REAL,
eps REAL,
profit_to_gr REAL,
source TEXT,
payload_json TEXT NOT NULL,
created_at TEXT NOT NULL,
updated_at TEXT NOT NULL,
UNIQUE(ts_code, end_date)
);
CREATE INDEX IF NOT EXISTS idx_fina_indicator_ts_code ON fina_indicator(ts_code);
CREATE INDEX IF NOT EXISTS idx_fina_indicator_end_date ON fina_indicator(end_date);
CREATE INDEX IF NOT EXISTS idx_fina_indicator_ts_end ON fina_indicator(ts_code, end_date);
"""财务指标查询功能"""
import json
import sys
from pathlib import Path
from typing import Optional, List, Dict
PROJECT_ROOT = Path(__file__).resolve().parents[2]
sys.path.insert(0, str(PROJECT_ROOT))
from analyzers.fina_indicator._shared.db import connect
def get_fina_indicator(ts_code: str, end_date: str = None) -> Optional[Dict]:
conn = connect()
try:
if end_date:
sql = "SELECT * FROM fina_indicator WHERE ts_code = ? AND end_date = ? ORDER BY ann_date DESC LIMIT 1"
cur = conn.execute(sql, (ts_code, end_date))
else:
sql = "SELECT * FROM fina_indicator WHERE ts_code = ? ORDER BY end_date DESC, ann_date DESC LIMIT 1"
cur = conn.execute(sql, (ts_code,))
row = cur.fetchone()
if not row:
return None
result = dict(row)
if result.get("payload_json"):
try:
result.update(json.loads(result["payload_json"]))
except (json.JSONDecodeError, TypeError):
pass
return result
finally:
conn.close()
def get_fina_indicator_history(ts_code: str, start_date: str = None, end_date: str = None, limit: int = None) -> List[Dict]:
conn = connect()
try:
conditions = ["ts_code = ?"]
params = [ts_code]
if start_date:
conditions.append("end_date >= ?")
params.append(start_date)
if end_date:
conditions.append("end_date <= ?")
params.append(end_date)
sql = f"SELECT * FROM fina_indicator WHERE {' AND '.join(conditions)} ORDER BY end_date DESC, ann_date DESC"
if limit is not None:
if not isinstance(limit, int) or limit <= 0:
raise ValueError("limit must be a positive int")
sql += " LIMIT ?"
params.append(limit)
cur = conn.execute(sql, params)
rows = cur.fetchall()
results = []
for row in rows:
result = dict(row)
if result.get("payload_json"):
try:
result.update(json.loads(result["payload_json"]))
except (json.JSONDecodeError, TypeError):
pass
results.append(result)
return results
finally:
conn.close()
def get_field_value(ts_code: str, field_name: str, end_date: str = None) -> Optional[float]:
record = get_fina_indicator(ts_code, end_date)
if not record:
return None
value = record.get(field_name)
if value is not None:
return float(value)
return None
import sqlite3
import tempfile
import unittest
from pathlib import Path
from unittest.mock import patch
from analyzers.fina_indicator import search
class TestFinaIndicatorSearch(unittest.TestCase):
@classmethod
def setUpClass(cls):
fd, db_path = tempfile.mkstemp(suffix=".db")
Path(db_path).unlink(missing_ok=True)
cls.db_path = Path(db_path)
schema_path = Path(__file__).resolve().parent / "schema.sql"
schema = schema_path.read_text(encoding="utf-8")
conn = sqlite3.connect(str(cls.db_path))
try:
conn.executescript(schema)
conn.execute(
"""
INSERT INTO fina_indicator (
ts_code, ann_date, end_date, roe, payload_json, source, created_at, updated_at
) VALUES (?, ?, ?, ?, ?, ?, ?, ?)
""",
(
"000001.SZ",
"20240101",
"20231231",
0.12,
'{"eps": 2.5}',
"unit-test",
"2024-01-01 00:00:00",
"2024-01-01 00:00:00",
),
)
conn.execute(
"""
INSERT INTO fina_indicator (
ts_code, ann_date, end_date, roe, payload_json, source, created_at, updated_at
) VALUES (?, ?, ?, ?, ?, ?, ?, ?)
""",
(
"000001.SZ",
"20240201",
"20240331",
0.15,
"{bad-json",
"unit-test",
"2024-02-01 00:00:00",
"2024-02-01 00:00:00",
),
)
conn.commit()
finally:
conn.close()
def connect_override():
conn = sqlite3.connect(str(cls.db_path))
conn.row_factory = sqlite3.Row
return conn
cls.connect_patch = patch(
"analyzers.fina_indicator.search.connect",
side_effect=connect_override,
)
cls.connect_patch.start()
@classmethod
def tearDownClass(cls):
cls.connect_patch.stop()
if cls.db_path.exists():
cls.db_path.unlink()
def test_history_limit_must_be_int(self):
with self.assertRaises(ValueError):
search.get_fina_indicator_history(
"000001.SZ",
limit="1; DROP TABLE fina_indicator",
)
def test_get_fina_indicator_skips_bad_payload(self):
record = search.get_fina_indicator("000001.SZ")
self.assertIsNotNone(record)
self.assertEqual(record["end_date"], "20240331")
self.assertEqual(record["roe"], 0.15)
def test_get_field_value_does_not_parse_payload_again(self):
with patch(
"analyzers.fina_indicator.search.get_fina_indicator",
return_value={"payload_json": "{bad", "eps": "3.5"},
), patch(
"analyzers.fina_indicator.search.json.loads",
side_effect=AssertionError("json.loads should not be called"),
):
value = search.get_field_value("000001.SZ", "eps")
self.assertEqual(value, 3.5)
"""主营业务构成数据库工具"""
import json
import sqlite3
from datetime import datetime
from pathlib import Path
from config import DB_PATH
SCHEMA_PATH = Path(__file__).parent.parent / "schema.sql"
def connect(db_path: Path = DB_PATH) -> sqlite3.Connection:
conn = sqlite3.connect(str(db_path))
conn.row_factory = sqlite3.Row
return conn
def init_table(conn: sqlite3.Connection = None):
if conn is None:
conn = connect()
should_close = True
else:
should_close = False
try:
schema = SCHEMA_PATH.read_text(encoding="utf-8")
conn.executescript(schema)
conn.commit()
finally:
if should_close:
conn.close()
def upsert_fina_mainbz(conn: sqlite3.Connection, rows: list) -> int:
if not rows:
return 0
base_cols = [
"ts_code", "end_date", "bz_item", "bz_type",
"bz_sales", "bz_profit", "bz_cost", "curr_type",
]
unique_cols = ["ts_code", "end_date", "bz_item"]
now = datetime.now().strftime("%Y-%m-%d %H:%M:%S")
db_rows = []
for row in rows:
db_row = {col: row.get(col) for col in base_cols}
if not db_row.get("ts_code") or not db_row.get("end_date") or not db_row.get("bz_item"):
continue
db_row["payload_json"] = json.dumps(row, ensure_ascii=False, default=str)
db_row["source"] = "tushare"
db_row["created_at"] = now
db_row["updated_at"] = now
db_rows.append(db_row)
if not db_rows:
return 0
all_cols = list(db_rows[0].keys())
placeholders = ", ".join(["?"] * len(all_cols))
columns_sql = ", ".join(all_cols)
conflict_sql = ", ".join(unique_cols)
update_cols = [c for c in all_cols if c not in unique_cols + ["created_at"]]
update_sql = ", ".join([f"{c} = excluded.{c}" for c in update_cols])
sql = (
f"INSERT INTO fina_mainbz ({columns_sql}) "
f"VALUES ({placeholders}) "
f"ON CONFLICT({conflict_sql}) DO UPDATE SET {update_sql}"
)
values = [[r.get(col) for col in all_cols] for r in db_rows]
try:
conn.executemany(sql, values)
conn.commit()
return len(db_rows)
except Exception:
conn.rollback()
raise
PRAGMA journal_mode=WAL;
CREATE TABLE IF NOT EXISTS fina_mainbz (
ts_code TEXT NOT NULL,
end_date TEXT NOT NULL,
bz_item TEXT NOT NULL,
bz_type TEXT,
bz_sales REAL,
bz_profit REAL,
bz_cost REAL,
curr_type TEXT,
source TEXT,
payload_json TEXT NOT NULL,
created_at TEXT NOT NULL,
updated_at TEXT NOT NULL,
UNIQUE(ts_code, end_date, bz_item)
);
CREATE INDEX IF NOT EXISTS idx_fina_mainbz_ts_code ON fina_mainbz(ts_code);
CREATE INDEX IF NOT EXISTS idx_fina_mainbz_end_date ON fina_mainbz(end_date);
CREATE INDEX IF NOT EXISTS idx_fina_mainbz_ts_end ON fina_mainbz(ts_code, end_date);
财务分析参考文档
基础指标定义
资产类指标
- 货币资金 (
money_cap): 公司持有的现金及银行存款 - 应收账款 (
accounts_receiv): 公司因销售商品、提供服务等应收取的款项 - 存货 (
inventories): 公司持有的商品、在产品、原材料等 - 流动资产合计 (
total_cur_assets): 一年内可变现的资产总和 - 固定资产 (
fix_assets): 公司长期使用的有形资产 - 总资产 (
total_assets): 公司拥有的全部资产
负债类指标
- 短期借款 (
st_borr): 一年内需要偿还的借款 - 应付账款 (
acct_payable): 公司因采购商品、接受服务等应支付的款项 - 流动负债合计 (
total_cur_liab): 一年内需要偿还的负债总和 - 总负债 (
total_liab): 公司需要偿还的全部负债
权益类指标
- 未分配利润 (
undistr_porfit): 公司累计未分配的利润 - 股东权益合计 (
total_hldr_eqy_exc_min_int): 股东拥有的净资产
常用财务比率
偿债能力指标
- 流动比率 = 流动资产合计 / 流动负债合计
- 标准值: > 2 较好,< 1 有风险
- 速动比率 = (流动资产合计 - 存货) / 流动负债合计
- 标准值: > 1 较好
- 资产负债率 = 总负债 / 总资产
- 标准值: < 60% 较好,> 80% 风险较高
资产结构指标
- 存货占比 = 存货 / 总资产
- 应收账款占比 = 应收账款 / 总资产
- 固定资产占比 = 固定资产 / 总资产
常见问题分析模板
"xx公司的存货如何"
1. 查询数据: 获取存货金额 (inventories) 2. 计算占比: 存货 / 总资产 3. 趋势分析: 对比最近4个季度的存货变化 4. 风险评估:
- 存货占比过高(>30%)可能积压
- 存货快速增长但收入未增长需警惕
"xx公司的偿债能力"
1. 查询数据: 流动资产、流动负债、总负债、总资产 2. 计算比率: 流动比率、速动比率、资产负债率 3. 评估: 对比行业标准值
"xx公司的资产结构"
1. 查询数据: 各类资产金额 2. 计算占比: 各类资产 / 总资产 3. 分析: 资产配置是否合理
数据获取方法
查询单字段
from analyzers.balancesheet.search import get_field_value
value = get_field_value('000001.SZ', 'inventories')查询完整记录
from analyzers.balancesheet.search import get_balancesheet
record = get_balancesheet('000001.SZ', end_date='20241231')查询历史数据
from analyzers.balancesheet.search import get_balancesheet_history
records = get_balancesheet_history('000001.SZ', limit=4)计算同比/环比
# 获取两年数据
records = get_balancesheet_history('000001.SZ', limit=8)
# 计算同比(去年同期)
if len(records) >= 4:
current = records[0]['inventories']
last_year = records[4]['inventories']
yoy_growth = (current - last_year) / last_year * 100注意事项
1. 公司类型差异: 银行/保险/证券的资产负债表结构不同,查询时注意 comp_type 2. 报告类型: 默认使用 report_type=1(合并报表),其他类型需明确指定 3. 数据完整性: 某些字段可能为空(如银行没有"存货"),需要容错处理 4. 时间范围: 计算同比需要至少4个季度数据,建议拉取N+1年数据
# analyzers/research
# analyzers/research/_shared
"""研报数据库工具"""
import sqlite3
from pathlib import Path
from config import DB_PATH
SCHEMA_PATH = Path(__file__).parent.parent / "schema.sql"
def connect() -> sqlite3.Connection:
conn = sqlite3.connect(str(DB_PATH))
conn.row_factory = sqlite3.Row
return conn
def init_table():
"""初始化研报表"""
conn = connect()
try:
schema = SCHEMA_PATH.read_text(encoding="utf-8")
conn.executescript(schema)
conn.commit()
finally:
conn.close()
def upsert_reports(rows: list[dict]) -> int:
"""批量插入研报元数据,返回插入条数"""
if not rows:
return 0
conn = connect()
try:
sql = """
INSERT INTO research_report
(trade_date, ts_code, name, title, abstr, report_type, author, inst_csname, ind_name, url)
VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
ON CONFLICT(url) DO UPDATE SET
trade_date = excluded.trade_date,
ts_code = excluded.ts_code,
name = excluded.name,
title = excluded.title,
abstr = excluded.abstr,
report_type = excluded.report_type,
author = excluded.author,
inst_csname = excluded.inst_csname,
ind_name = excluded.ind_name
"""
values = [
(r.get("trade_date"), r.get("ts_code"), r.get("name"), r.get("title"),
r.get("abstr"), r.get("report_type"), r.get("author"),
r.get("inst_csname"), r.get("ind_name"), r.get("url"))
for r in rows
]
conn.executemany(sql, values)
conn.commit()
return len(rows)
finally:
conn.close()
def update_local_path(report_id: int, local_path: str):
"""更新本地路径"""
conn = connect()
try:
conn.execute(
"UPDATE research_report SET local_path = ? WHERE id = ?",
(local_path, report_id)
)
conn.commit()
finally:
conn.close()
def update_parsed_at(report_id: int, parsed_at: str):
"""更新解析时间"""
conn = connect()
try:
conn.execute(
"UPDATE research_report SET parsed_at = ? WHERE id = ?",
(parsed_at, report_id)
)
conn.commit()
finally:
conn.close()
if __name__ == "__main__":
import argparse
parser = argparse.ArgumentParser(description="研报数据库工具")
parser.add_argument("--init", action="store_true", help="初始化表")
args = parser.parse_args()
if args.init:
init_table()
print(f"表已初始化: {DB_PATH}")
else:
parser.print_help()
# analyzers/research/download.py
"""下载研报 PDF"""
import re
import time
from pathlib import Path
import requests
from ._shared.db import connect, update_local_path
PROJECT_ROOT = Path(__file__).resolve().parents[2]
OUTPUT_DIR = PROJECT_ROOT / "output" / "reports"
def sanitize_filename(name: str) -> str:
"""清理文件名"""
return re.sub(r'[<>:"/\\|?*]', '_', name)[:100]
def download_report(report_id: int, force: bool = False) -> str | None:
"""
下载单个研报 PDF
Args:
report_id: 研报 ID
force: 强制重新下载
Returns:
本地文件路径,失败返回 None
"""
conn = connect()
try:
cur = conn.execute(
"SELECT id, ts_code, title, url, local_path FROM research_report WHERE id = ?",
(report_id,)
)
row = cur.fetchone()
finally:
conn.close()
if not row:
print(f"研报 {report_id} 不存在")
return None
row = dict(row)
# 已下载且不强制重下
if row["local_path"] and Path(row["local_path"]).exists() and not force:
return row["local_path"]
# 按股票代码分目录
ts_code = row["ts_code"] or "unknown"
subdir = OUTPUT_DIR / ts_code.replace(".", "_")
subdir.mkdir(parents=True, exist_ok=True)
# 文件名
filename = sanitize_filename(row["title"]) + ".pdf"
local_path = subdir / filename
# 下载
try:
resp = requests.get(row["url"], timeout=30)
resp.raise_for_status()
local_path.write_bytes(resp.content)
# 更新 DB
update_local_path(report_id, str(local_path))
print(f"已下载: {local_path}")
return str(local_path)
except Exception as e:
print(f"下载失败 [{report_id}]: {e}")
return None
def download_batch(report_ids: list[int], delay: float = 1.0) -> list[str]:
"""
批量下载
Args:
report_ids: 研报 ID 列表
delay: 下载间隔(秒)
Returns:
成功下载的本地路径列表
"""
paths = []
for i, rid in enumerate(report_ids):
path = download_report(rid)
if path:
paths.append(path)
if i < len(report_ids) - 1:
time.sleep(delay)
return paths
if __name__ == "__main__":
import argparse
parser = argparse.ArgumentParser(description="下载研报 PDF")
parser.add_argument("ids", nargs="+", type=int, help="研报 ID")
parser.add_argument("--force", action="store_true", help="强制重新下载")
args = parser.parse_args()
for rid in args.ids:
download_report(rid, force=args.force)
# analyzers/research/pipeline.py
"""研报分析流程入口"""
import json
import sys
from pathlib import Path
PROJECT_ROOT = Path(__file__).resolve().parents[2]
sys.path.insert(0, str(PROJECT_ROOT))
from fetchers.research_report import fetch_research_report
from analyzers.research.search import search_reports
from analyzers.research.download import download_report, download_batch
from interpreters.research.extract import extract_report, extract_batch
def run_search(
ts_code: str = None,
ind_name: str = None,
keyword: str = None,
inst_csname: str = None,
start_date: str = None,
end_date: str = None,
limit: int = 100,
fetch_if_empty: bool = True,
) -> list[dict]:
"""
搜索研报
Args:
fetch_if_empty: 若 DB 无数据,自动从 API 拉取
"""
results = search_reports(
ts_code=ts_code,
ind_name=ind_name,
keyword=keyword,
inst_csname=inst_csname,
start_date=start_date,
end_date=end_date,
limit=limit,
)
if not results and fetch_if_empty:
print("DB 无数据,从 API 拉取...")
fetch_research_report(
ts_code=ts_code,
ind_name=ind_name,
inst_csname=inst_csname,
start_date=start_date,
end_date=end_date,
)
results = search_reports(
ts_code=ts_code,
ind_name=ind_name,
keyword=keyword,
inst_csname=inst_csname,
start_date=start_date,
end_date=end_date,
limit=limit,
)
return results
def run_extract(report_id: int, force: bool = False) -> dict | None:
"""下载并解析单个研报"""
return extract_report(report_id, force=force)
def run_batch_extract(report_ids: list[int]) -> list[dict]:
"""批量解析"""
return extract_batch(report_ids)
if __name__ == "__main__":
import argparse
parser = argparse.ArgumentParser(description="研报分析流程")
subparsers = parser.add_subparsers(dest="command")
# 搜索
search_p = subparsers.add_parser("search", help="搜索研报")
search_p.add_argument("--ts-code", help="股票代码")
search_p.add_argument("--ind", help="行业名称")
search_p.add_argument("--keyword", help="标题关键词")
search_p.add_argument("--inst", help="券商名称")
search_p.add_argument("--start", help="开始日期")
search_p.add_argument("--end", help="结束日期")
search_p.add_argument("--limit", type=int, default=20)
# 解析
extract_p = subparsers.add_parser("extract", help="解析研报")
extract_p.add_argument("ids", nargs="+", type=int, help="研报 ID")
extract_p.add_argument("--force", action="store_true")
args = parser.parse_args()
if args.command == "search":
results = run_search(
ts_code=args.ts_code,
ind_name=args.ind,
keyword=args.keyword,
inst_csname=args.inst,
start_date=args.start,
end_date=args.end,
limit=args.limit,
)
for r in results:
print(f"[{r['id']}] [{r['trade_date']}] {r['title'][:60]}")
print(f"\n共 {len(results)} 条")
elif args.command == "extract":
for rid in args.ids:
result = run_extract(rid, force=args.force)
if result:
print(json.dumps(result, ensure_ascii=False, indent=2))
else:
parser.print_help()
-- analyzers/research/schema.sql
-- 研报元数据表
CREATE TABLE IF NOT EXISTS research_report (
id INTEGER PRIMARY KEY AUTOINCREMENT,
trade_date TEXT NOT NULL, -- 发布日期 (YYYYMMDD)
ts_code TEXT, -- 股票代码 (如 000001.SZ)
name TEXT, -- 股票名称
title TEXT NOT NULL, -- 研报标题
abstr TEXT, -- 研报摘要
report_type TEXT, -- 报告类型 (深度/点评/行业等)
author TEXT, -- 作者
inst_csname TEXT, -- 发布机构中文名
ind_name TEXT, -- 行业名称
url TEXT NOT NULL, -- 研报原始URL (唯一键)
local_path TEXT, -- 本地PDF路径
parsed_at TEXT, -- 解析完成时间
created_at TEXT DEFAULT CURRENT_TIMESTAMP, -- 记录创建时间
UNIQUE(url)
);
-- 索引:按股票代码查询
CREATE INDEX IF NOT EXISTS idx_rr_ts_code ON research_report(ts_code);
-- 索引:按发布日期查询
CREATE INDEX IF NOT EXISTS idx_rr_trade_date ON research_report(trade_date);
-- 索引:按行业查询
CREATE INDEX IF NOT EXISTS idx_rr_ind_name ON research_report(ind_name);
-- 索引:按机构查询
CREATE INDEX IF NOT EXISTS idx_rr_inst_csname ON research_report(inst_csname);
# analyzers/research/search.py
"""搜索研报"""
from ._shared.db import connect
def search_reports(
ts_code: str = None,
ind_name: str = None,
keyword: str = None,
inst_csname: str = None,
start_date: str = None,
end_date: str = None,
limit: int = 100,
) -> list[dict]:
"""
从 DB 搜索研报
Args:
ts_code: 股票代码
ind_name: 行业名称(模糊匹配)
keyword: 标题关键词(模糊匹配)
inst_csname: 券商名称(模糊匹配)
start_date/end_date: 日期范围
limit: 返回条数上限
Returns:
研报元数据列表
"""
conn = connect()
conditions = []
params = []
if ts_code:
conditions.append("ts_code = ?")
params.append(ts_code)
if ind_name:
conditions.append("ind_name LIKE ?")
params.append(f"%{ind_name}%")
if keyword:
conditions.append("title LIKE ?")
params.append(f"%{keyword}%")
if inst_csname:
conditions.append("inst_csname LIKE ?")
params.append(f"%{inst_csname}%")
if start_date:
conditions.append("trade_date >= ?")
params.append(start_date)
if end_date:
conditions.append("trade_date <= ?")
params.append(end_date)
where_clause = " AND ".join(conditions) if conditions else "1=1"
sql = f"""
SELECT id, trade_date, ts_code, name, title, abstr, report_type,
author, inst_csname, ind_name, url, local_path, parsed_at
FROM research_report
WHERE {where_clause}
ORDER BY trade_date DESC
LIMIT ?
"""
params.append(limit)
try:
cur = conn.execute(sql, params)
rows = [dict(r) for r in cur.fetchall()]
finally:
conn.close()
return rows
if __name__ == "__main__":
import argparse
parser = argparse.ArgumentParser(description="搜索研报")
parser.add_argument("--ts-code", help="股票代码")
parser.add_argument("--ind", help="行业名称")
parser.add_argument("--keyword", help="标题关键词")
parser.add_argument("--inst", help="券商名称")
parser.add_argument("--start", help="开始日期")
parser.add_argument("--end", help="结束日期")
parser.add_argument("--limit", type=int, default=20, help="返回条数")
args = parser.parse_args()
results = search_reports(
ts_code=args.ts_code,
ind_name=args.ind,
keyword=args.keyword,
inst_csname=args.inst,
start_date=args.start,
end_date=args.end,
limit=args.limit,
)
for r in results:
print(f"[{r['trade_date']}] {r['title'][:50]}...")
print(f"\n共 {len(results)} 条")
"""股票技术面因子数据库工具"""
import json
import sqlite3
from datetime import datetime
from pathlib import Path
from config import DB_PATH
SCHEMA_PATH = Path(__file__).parent.parent / "schema.sql"
# 提取到独立列的字段
EXTRACT_COLS = [
"ts_code", "trade_date",
"close", "close_hfq", "close_qfq", "pct_chg", "vol", "amount",
"turnover_rate", "pe_ttm", "pb", "total_mv", "circ_mv",
"macd_dif_qfq", "macd_dea_qfq", "macd_qfq",
"kdj_k_qfq", "kdj_d_qfq", "kdj_qfq",
"rsi_qfq_6", "rsi_qfq_12", "rsi_qfq_24",
"boll_upper_qfq", "boll_mid_qfq", "boll_lower_qfq",
"ma_qfq_5", "ma_qfq_10", "ma_qfq_20", "ma_qfq_60", "ma_qfq_250",
]
def connect(db_path: Path = DB_PATH) -> sqlite3.Connection:
conn = sqlite3.connect(str(db_path))
conn.row_factory = sqlite3.Row
return conn
def init_table(conn: sqlite3.Connection = None):
if conn is None:
conn = connect()
should_close = True
else:
should_close = False
try:
schema = SCHEMA_PATH.read_text(encoding="utf-8")
conn.executescript(schema)
conn.commit()
finally:
if should_close:
conn.close()
def upsert_stk_factor(conn: sqlite3.Connection, rows: list) -> int:
if not rows:
return 0
unique_cols = ["ts_code", "trade_date"]
now = datetime.now().strftime("%Y-%m-%d %H:%M:%S")
db_rows = []
for row in rows:
db_row = {col: row.get(col) for col in EXTRACT_COLS}
if not db_row.get("ts_code") or not db_row.get("trade_date"):
continue
db_row["payload_json"] = json.dumps(row, ensure_ascii=False, default=str)
db_row["source"] = "tushare"
db_row["created_at"] = now
db_row["updated_at"] = now
db_rows.append(db_row)
if not db_rows:
return 0
all_cols = list(db_rows[0].keys())
placeholders = ", ".join(["?"] * len(all_cols))
columns_sql = ", ".join(all_cols)
conflict_sql = ", ".join(unique_cols)
update_cols = [c for c in all_cols if c not in unique_cols + ["created_at"]]
update_sql = ", ".join([f"{c} = excluded.{c}" for c in update_cols])
sql = (
f"INSERT INTO stk_factor ({columns_sql}) "
f"VALUES ({placeholders}) "
f"ON CONFLICT({conflict_sql}) DO UPDATE SET {update_sql}"
)
values = [[r.get(col) for col in all_cols] for r in db_rows]
try:
conn.executemany(sql, values)
conn.commit()
return len(db_rows)
except Exception:
conn.rollback()
raise
PRAGMA journal_mode=WAL;
CREATE TABLE IF NOT EXISTS stk_factor (
ts_code TEXT NOT NULL,
trade_date TEXT NOT NULL,
close REAL,
close_hfq REAL,
close_qfq REAL,
pct_chg REAL,
vol REAL,
amount REAL,
turnover_rate REAL,
pe_ttm REAL,
pb REAL,
total_mv REAL,
circ_mv REAL,
macd_dif_qfq REAL,
macd_dea_qfq REAL,
macd_qfq REAL,
kdj_k_qfq REAL,
kdj_d_qfq REAL,
kdj_qfq REAL,
rsi_qfq_6 REAL,
rsi_qfq_12 REAL,
rsi_qfq_24 REAL,
boll_upper_qfq REAL,
boll_mid_qfq REAL,
boll_lower_qfq REAL,
ma_qfq_5 REAL,
ma_qfq_10 REAL,
ma_qfq_20 REAL,
ma_qfq_60 REAL,
ma_qfq_250 REAL,
source TEXT,
payload_json TEXT NOT NULL,
created_at TEXT NOT NULL,
updated_at TEXT NOT NULL,
UNIQUE(ts_code, trade_date)
);
CREATE INDEX IF NOT EXISTS idx_stk_factor_ts_code ON stk_factor(ts_code);
CREATE INDEX IF NOT EXISTS idx_stk_factor_trade_date ON stk_factor(trade_date);
"""神奇九转数据库工具"""
import json
import sqlite3
from datetime import datetime
from pathlib import Path
from config import DB_PATH
SCHEMA_PATH = Path(__file__).parent.parent / "schema.sql"
EXTRACT_COLS = [
"ts_code", "trade_date", "freq",
"open", "high", "low", "close", "vol", "amount",
"up_count", "down_count", "nine_up_turn", "nine_down_turn",
]
def connect(db_path: Path = DB_PATH) -> sqlite3.Connection:
conn = sqlite3.connect(str(db_path))
conn.row_factory = sqlite3.Row
return conn
def init_table(conn: sqlite3.Connection = None):
if conn is None:
conn = connect()
should_close = True
else:
should_close = False
try:
schema = SCHEMA_PATH.read_text(encoding="utf-8")
conn.executescript(schema)
conn.commit()
finally:
if should_close:
conn.close()
def upsert_stk_nineturn(conn: sqlite3.Connection, rows: list) -> int:
if not rows:
return 0
unique_cols = ["ts_code", "trade_date", "freq"]
now = datetime.now().strftime("%Y-%m-%d %H:%M:%S")
db_rows = []
for row in rows:
db_row = {col: row.get(col) for col in EXTRACT_COLS}
if not db_row.get("ts_code") or not db_row.get("trade_date"):
continue
if not db_row.get("freq"):
db_row["freq"] = "daily"
db_row["payload_json"] = json.dumps(row, ensure_ascii=False, default=str)
db_row["created_at"] = now
db_row["updated_at"] = now
db_rows.append(db_row)
if not db_rows:
return 0
all_cols = list(db_rows[0].keys())
placeholders = ", ".join(["?"] * len(all_cols))
columns_sql = ", ".join(all_cols)
conflict_sql = ", ".join(unique_cols)
update_cols = [c for c in all_cols if c not in unique_cols + ["created_at"]]
update_sql = ", ".join([f"{c} = excluded.{c}" for c in update_cols])
sql = (
f"INSERT INTO stk_nineturn ({columns_sql}) "
f"VALUES ({placeholders}) "
f"ON CONFLICT({conflict_sql}) DO UPDATE SET {update_sql}"
)
values = [[r.get(col) for col in all_cols] for r in db_rows]
try:
conn.executemany(sql, values)
conn.commit()
return len(db_rows)
except Exception:
conn.rollback()
raise
PRAGMA journal_mode=WAL;
CREATE TABLE IF NOT EXISTS stk_nineturn (
ts_code TEXT NOT NULL,
trade_date TEXT NOT NULL,
freq TEXT NOT NULL DEFAULT 'daily',
open REAL,
high REAL,
low REAL,
close REAL,
vol REAL,
amount REAL,
up_count INTEGER,
down_count INTEGER,
nine_up_turn TEXT,
nine_down_turn TEXT,
payload_json TEXT NOT NULL,
created_at TEXT NOT NULL,
updated_at TEXT NOT NULL,
UNIQUE(ts_code, trade_date, freq)
);
CREATE INDEX IF NOT EXISTS idx_nineturn_ts_code ON stk_nineturn(ts_code);
CREATE INDEX IF NOT EXISTS idx_nineturn_trade_date ON stk_nineturn(trade_date);
CREATE INDEX IF NOT EXISTS idx_nineturn_signal ON stk_nineturn(nine_up_turn, nine_down_turn);
"""机构调研数据库工具"""
import sqlite3
from pathlib import Path
from config import DB_PATH
SCHEMA_PATH = Path(__file__).parent.parent / "schema.sql"
def connect() -> sqlite3.Connection:
conn = sqlite3.connect(str(DB_PATH))
conn.row_factory = sqlite3.Row
return conn
def init_tables():
"""初始化调研表"""
conn = connect()
try:
schema = SCHEMA_PATH.read_text(encoding="utf-8")
conn.executescript(schema)
conn.commit()
finally:
conn.close()
def upsert_event(row: dict) -> int:
"""插入或更新调研事件,返回event_id"""
conn = connect()
try:
sql = """
INSERT INTO stk_surv_event
(ts_code, name, surv_date, rece_place, rece_mode, comp_rece, content)
VALUES (?, ?, ?, ?, ?, ?, ?)
ON CONFLICT(ts_code, surv_date) DO UPDATE SET
name = excluded.name,
rece_place = excluded.rece_place,
rece_mode = excluded.rece_mode,
comp_rece = excluded.comp_rece,
content = COALESCE(excluded.content, stk_surv_event.content)
"""
cur = conn.execute(sql, (
row.get("ts_code"), row.get("name"), row.get("surv_date"),
row.get("rece_place"), row.get("rece_mode"), row.get("comp_rece"),
row.get("content")
))
conn.commit()
# 获取event_id
cur = conn.execute(
"SELECT id FROM stk_surv_event WHERE ts_code = ? AND surv_date = ?",
(row.get("ts_code"), row.get("surv_date"))
)
return cur.fetchone()["id"]
finally:
conn.close()
def upsert_participants(event_id: int, participants: list[dict]) -> int:
"""批量插入参与人员"""
if not participants:
return 0
conn = connect()
try:
sql = """
INSERT INTO stk_surv_participant (event_id, fund_visitors, rece_org, org_type)
VALUES (?, ?, ?, ?)
ON CONFLICT(event_id, fund_visitors, rece_org) DO UPDATE SET
org_type = excluded.org_type
"""
values = [
(event_id, p.get("fund_visitors"), p.get("rece_org"), p.get("org_type"))
for p in participants
]
conn.executemany(sql, values)
conn.commit()
return len(participants)
finally:
conn.close()
if __name__ == "__main__":
import argparse
parser = argparse.ArgumentParser(description="机构调研数据库工具")
parser.add_argument("--init", action="store_true", help="初始化表")
args = parser.parse_args()
if args.init:
init_tables()
print(f"表已初始化: {DB_PATH}")
else:
parser.print_help()
# analyzers/stk_surv/pipeline.py
"""机构调研分析流程入口"""
import json
import sys
from pathlib import Path
PROJECT_ROOT = Path(__file__).resolve().parents[2]
sys.path.insert(0, str(PROJECT_ROOT))
from fetchers.stk_surv import fetch_stk_surv
from analyzers.stk_surv.search import (
search_by_company, search_by_org, search_by_person, get_survey_detail
)
def run_fetch_recent(days: int = 30) -> int:
"""拉取最近N天的所有调研数据"""
from datetime import datetime, timedelta
end_date = datetime.now()
start_date = end_date - timedelta(days=days)
start_str = start_date.strftime("%Y%m%d")
end_str = end_date.strftime("%Y%m%d")
return fetch_stk_surv(start_date=start_str, end_date=end_str, include_content=True)
def run_search(query: str, query_type: str = "company", days: int = None) -> list[dict]:
"""
统一搜索入口
Args:
query: 搜索关键词
query_type: company/org/person
days: 可选,最近N天
"""
if query_type == "company":
return search_by_company(query, days)
elif query_type == "org":
return search_by_org(query, days)
elif query_type == "person":
return search_by_person(query, days)
else:
raise ValueError(f"不支持的查询类型: {query_type}")
if __name__ == "__main__":
import argparse
parser = argparse.ArgumentParser(description="机构调研分析流程")
subparsers = parser.add_subparsers(dest="command")
# 拉取
fetch_p = subparsers.add_parser("fetch", help="拉取调研数据")
fetch_p.add_argument("--days", type=int, default=30, help="最近N天,默认30")
# 搜索
search_p = subparsers.add_parser("search", help="搜索调研")
search_p.add_argument("query", help="搜索关键词")
search_p.add_argument("--type", choices=["company", "org", "person"], default="company", help="查询类型")
search_p.add_argument("--days", type=int, help="最近N天")
# 详情
detail_p = subparsers.add_parser("detail", help="查看调研详情")
detail_p.add_argument("--event-id", type=int, required=True, help="事件ID")
args = parser.parse_args()
if args.command == "fetch":
count = run_fetch_recent(args.days)
print(f"已入库 {count} 个调研事件")
elif args.command == "search":
results = run_search(args.query, args.type, args.days)
for r in results:
print(f"[{r['id']}] [{r['surv_date']}] {r['name']} ({r['ts_code']})")
print(f"\n共 {len(results)} 条")
elif args.command == "detail":
result = get_survey_detail(args.event_id)
if result:
print(json.dumps(result, ensure_ascii=False, indent=2))
else:
print(f"未找到事件 ID {args.event_id}")
else:
parser.print_help()
-- analyzers/stk_surv/schema.sql
CREATE TABLE IF NOT EXISTS stk_surv_event (
id INTEGER PRIMARY KEY AUTOINCREMENT,
ts_code TEXT NOT NULL,
name TEXT NOT NULL,
surv_date TEXT NOT NULL,
rece_place TEXT,
rece_mode TEXT,
comp_rece TEXT,
content TEXT,
created_at TEXT DEFAULT CURRENT_TIMESTAMP,
UNIQUE(ts_code, surv_date)
);
CREATE INDEX IF NOT EXISTS idx_surv_event_ts_code ON stk_surv_event(ts_code);
CREATE INDEX IF NOT EXISTS idx_surv_event_date ON stk_surv_event(surv_date);
CREATE INDEX IF NOT EXISTS idx_surv_event_name ON stk_surv_event(name);
CREATE TABLE IF NOT EXISTS stk_surv_participant (
id INTEGER PRIMARY KEY AUTOINCREMENT,
event_id INTEGER NOT NULL,
fund_visitors TEXT,
rece_org TEXT,
org_type TEXT,
created_at TEXT DEFAULT CURRENT_TIMESTAMP,
FOREIGN KEY(event_id) REFERENCES stk_surv_event(id),
UNIQUE(event_id, fund_visitors, rece_org)
);
CREATE INDEX IF NOT EXISTS idx_surv_part_event_id ON stk_surv_participant(event_id);
CREATE INDEX IF NOT EXISTS idx_surv_part_org ON stk_surv_participant(rece_org);
CREATE INDEX IF NOT EXISTS idx_surv_part_visitor ON stk_surv_participant(fund_visitors);
# analyzers/stk_surv/search.py
"""机构调研检索功能"""
from datetime import datetime, timedelta
from ._shared.db import connect
def search_by_company(company_name: str, days: int = None) -> list[dict]:
"""
按公司名称搜索调研
Args:
company_name: 公司名称(支持模糊匹配)
days: 可选,最近N天
Returns:
调研事件列表
"""
conn = connect()
try:
conditions = ["name LIKE ?"]
params = [f"%{company_name}%"]
if days:
cutoff_date = (datetime.now() - timedelta(days=days)).strftime("%Y%m%d")
conditions.append("surv_date >= ?")
params.append(cutoff_date)
sql = f"""
SELECT id, ts_code, name, surv_date, rece_place, rece_mode, comp_rece
FROM stk_surv_event
WHERE {' AND '.join(conditions)}
ORDER BY surv_date DESC
"""
cur = conn.execute(sql, params)
return [dict(r) for r in cur.fetchall()]
finally:
conn.close()
def search_by_org(org_name: str, days: int = None) -> list[dict]:
"""
按机构名称搜索(该机构参加了哪些调研)
Args:
org_name: 机构名称(支持模糊匹配)
days: 可选,最近N天
Returns:
调研事件列表
"""
conn = connect()
try:
conditions = ["p.rece_org LIKE ?"]
params = [f"%{org_name}%"]
if days:
cutoff_date = (datetime.now() - timedelta(days=days)).strftime("%Y%m%d")
conditions.append("e.surv_date >= ?")
params.append(cutoff_date)
sql = f"""
SELECT DISTINCT e.id, e.ts_code, e.name, e.surv_date,
e.rece_place, e.rece_mode, e.comp_rece
FROM stk_surv_event e
JOIN stk_surv_participant p ON e.id = p.event_id
WHERE {' AND '.join(conditions)}
ORDER BY e.surv_date DESC
"""
cur = conn.execute(sql, params)
return [dict(r) for r in cur.fetchall()]
finally:
conn.close()
def search_by_person(person_name: str, days: int = None) -> list[dict]:
"""
按参与人员姓名搜索
Args:
person_name: 人员姓名(支持模糊匹配)
days: 可选,最近N天
Returns:
调研事件列表
"""
conn = connect()
try:
conditions = ["p.fund_visitors LIKE ?"]
params = [f"%{person_name}%"]
if days:
cutoff_date = (datetime.now() - timedelta(days=days)).strftime("%Y%m%d")
conditions.append("e.surv_date >= ?")
params.append(cutoff_date)
sql = f"""
SELECT DISTINCT e.id, e.ts_code, e.name, e.surv_date,
e.rece_place, e.rece_mode, e.comp_rece
FROM stk_surv_event e
JOIN stk_surv_participant p ON e.id = p.event_id
WHERE {' AND '.join(conditions)}
ORDER BY e.surv_date DESC
"""
cur = conn.execute(sql, params)
return [dict(r) for r in cur.fetchall()]
finally:
conn.close()
def get_survey_detail(event_id: int) -> dict:
"""
获取某次调研的详细信息
Args:
event_id: 事件ID
Returns:
{
"event": {...},
"participants": [...]
}
"""
conn = connect()
try:
# 获取事件信息
cur = conn.execute(
"SELECT * FROM stk_surv_event WHERE id = ?",
(event_id,)
)
event_row = cur.fetchone()
if not event_row:
return None
event = dict(event_row)
# 获取参与人员
cur = conn.execute(
"SELECT fund_visitors, rece_org, org_type FROM stk_surv_participant WHERE event_id = ?",
(event_id,)
)
participants = [dict(r) for r in cur.fetchall()]
return {
"event": event,
"participants": participants
}
finally:
conn.close()
"""
stock-research-group 全局配置
所有模块应从此处导入 DB_PATH,避免硬编码数据库路径。
如需更改数据库位置,只需修改此文件。
"""
from pathlib import Path
# 项目根目录
PROJECT_ROOT = Path(__file__).resolve().parent
# 数据库路径(本地项目目录)
DB_PATH = Path("/Volumes/固态 1t/openclaw-data/finance.db")
架构说明
三层职责
| 层级 | 目录 | 职责 | 输入 → 输出 |
|---|---|---|---|
| 获取层 | fetchers/ | 从外部 API 拉取原始数据 | API 参数 → DB/CSV |
| 分析层 | analyzers/ | 数据处理、计算、对比 | DB → 结构化结果 |
| 解释层 | interpreters/ | 生成人类可读报告 | JSON/DB → 文字/CSV |
数据流
Tushare API → fetchers/ → data/finance.db → analyzers/ → interpreters/ → output/分析域
| 域 | 目录 | 模块文档 |
|---|---|---|
| 财务分析 | analyzers/earnings/ balancesheet/ cashflow/ fina_indicator/ | modules/fundamental/ |
| 通用数据 | analyzers/research/ broker_recommend/ stk_surv/ | modules/data/ |
| 技术分析 | TODO | modules/technical/ |
| 舆情分析 | TODO | modules/sentiment/ |
命名规范
| 类型 | 规范 | 示例 |
|---|---|---|
| 获取脚本 | {数据类型}.py | balancesheet.py |
| 分析目录 | {业务场景}/ | earnings/, cashflow/ |
| 共用模块 | _shared/ | 下划线前缀,区分业务脚本 |
| 流程入口 | pipeline.py | 每个分析场景的统一入口 |
| 表结构 | schema.sql | 每个分析目录内 |
扩展步骤
1. fetchers/ 新建 {数据类型}.py 2. analyzers/ 对应域下新建目录,含 _shared/db.py + schema.sql + pipeline.py 3. 需要报告输出时,interpreters/ 新建对应目录 4. modules/{域}/ 新建模块文档 5. sync_registry.json 追加同步条目 6. SKILL.md 模块索引表追加一行
全量同步
full_sync.py 读取 sync_registry.json 注册表,顺序执行所有 enabled: true 的步骤。
注册表支持变量:{today} {first_period} {last_period}(由 Period 自适应规则自动计算)。
输出参数
名称 类型 默认显示 描述 ts_code str Y TS股票代码 ann_date str Y 公告日期 f_ann_date str Y 实际公告日期 end_date str Y 报告期 report_type str Y 报表类型 comp_type str Y 公司类型(1一般工商业2银行3保险4证券) end_type str Y 报告期类型 total_share float Y 期末总股本 cap_rese float Y 资本公积金 undistr_porfit float Y 未分配利润 surplus_rese float Y 盈余公积金 special_rese float Y 专项储备 money_cap float Y 货币资金 trad_asset float Y 交易性金融资产 notes_receiv float Y 应收票据 accounts_receiv float Y 应收账款 oth_receiv float Y 其他应收款 prepayment float Y 预付款项 div_receiv float Y 应收股利 int_receiv float Y 应收利息 inventories float Y 存货 amor_exp float Y 待摊费用 nca_within_1y float Y 一年内到期的非流动资产 sett_rsrv float Y 结算备付金 loanto_oth_bank_fi float Y 拆出资金 premium_receiv float Y 应收保费 reinsur_receiv float Y 应收分保账款 reinsur_res_receiv float Y 应收分保合同准备金 pur_resale_fa float Y 买入返售金融资产 oth_cur_assets float Y 其他流动资产 total_cur_assets float Y 流动资产合计 fa_avail_for_sale float Y 可供出售金融资产 htm_invest float Y 持有至到期投资 lt_eqt_invest float Y 长期股权投资 invest_real_estate float Y 投资性房地产 time_deposits float Y 定期存款 oth_assets float Y 其他资产 lt_rec float Y 长期应收款 fix_assets float Y 固定资产 cip float Y 在建工程 const_materials float Y 工程物资 fixed_assets_disp float Y 固定资产清理 produc_bio_assets float Y 生产性生物资产 oil_and_gas_assets float Y 油气资产 intan_assets float Y 无形资产 r_and_d float Y 研发支出 goodwill float Y 商誉 lt_amor_exp float Y 长期待摊费用 defer_tax_assets float Y 递延所得税资产 decr_in_disbur float Y 发放贷款及垫款 oth_nca float Y 其他非流动资产 total_nca float Y 非流动资产合计 cash_reser_cb float Y 现金及存放中央银行款项 depos_in_oth_bfi float Y 存放同业和其它金融机构款项 prec_metals float Y 贵金属 deriv_assets float Y 衍生金融资产 rr_reins_une_prem float Y 应收分保未到期责任准备金 rr_reins_outstd_cla float Y 应收分保未决赔款准备金 rr_reins_lins_liab float Y 应收分保寿险责任准备金 rr_reins_lthins_liab float Y 应收分保长期健康险责任准备金 refund_depos float Y 存出保证金 ph_pledge_loans float Y 保户质押贷款 refund_cap_depos float Y 存出资本保证金 indep_acct_assets float Y 独立账户资产 client_depos float Y 其中:客户资金存款 client_prov float Y 其中:客户备付金 transac_seat_fee float Y 其中:交易席位费 invest_as_receiv float Y 应收款项类投资 total_assets float Y 资产总计 lt_borr float Y 长期借款 st_borr float Y 短期借款 cb_borr float Y 向中央银行借款 depos_ib_deposits float Y 吸收存款及同业存放 loan_oth_bank float Y 拆入资金 trading_fl float Y 交易性金融负债 notes_payable float Y 应付票据 acct_payable float Y 应付账款 adv_receipts float Y 预收款项 sold_for_repur_fa float Y 卖出回购金融资产款 comm_payable float Y 应付手续费及佣金 payroll_payable float Y 应付职工薪酬 taxes_payable float Y 应交税费 int_payable float Y 应付利息 div_payable float Y 应付股利 oth_payable float Y 其他应付款 acc_exp float Y 预提费用 deferred_inc float Y 递延收益 st_bonds_payable float Y 应付短期债券 payable_to_reinsurer float Y 应付分保账款 rsrv_insur_cont float Y 保险合同准备金 acting_trading_sec float Y 代理买卖证券款 acting_uw_sec float Y 代理承销证券款 non_cur_liab_due_1y float Y 一年内到期的非流动负债 oth_cur_liab float Y 其他流动负债 total_cur_liab float Y 流动负债合计 bond_payable float Y 应付债券 lt_payable float Y 长期应付款 specific_payables float Y 专项应付款 estimated_liab float Y 预计负债 defer_tax_liab float Y 递延所得税负债 defer_inc_non_cur_liab float Y 递延收益-非流动负债 oth_ncl float Y 其他非流动负债 total_ncl float Y 非流动负债合计 depos_oth_bfi float Y 同业和其它金融机构存放款项 deriv_liab float Y 衍生金融负债 depos float Y 吸收存款 agency_bus_liab float Y 代理业务负债 oth_liab float Y 其他负债 prem_receiv_adva float Y 预收保费 depos_received float Y 存入保证金 ph_invest float Y 保户储金及投资款 reser_une_prem float Y 未到期责任准备金 reser_outstd_claims float Y 未决赔款准备金 reser_lins_liab float Y 寿险责任准备金 reser_lthins_liab float Y 长期健康险责任准备金 indept_acc_liab float Y 独立账户负债 pledge_borr float Y 其中:质押借款 indem_payable float Y 应付赔付款 policy_div_payable float Y 应付保单红利 total_liab float Y 负债合计 treasury_share float Y 减:库存股 ordin_risk_reser float Y 一般风险准备 forex_differ float Y 外币报表折算差额 invest_loss_unconf float Y 未确认的投资损失 minority_int float Y 少数股东权益 total_hldr_eqy_exc_min_int float Y 股东权益合计(不含少数股东权益) total_hldr_eqy_inc_min_int float Y 股东权益合计(含少数股东权益) total_liab_hldr_eqy float Y 负债及股东权益总计 lt_payroll_payable float Y 长期应付职工薪酬 oth_comp_income float Y 其他综合收益 oth_eqt_tools float Y 其他权益工具 oth_eqt_tools_p_shr float Y 其他权益工具(优先股) lending_funds float Y 融出资金 acc_receivable float Y 应收款项 st_fin_payable float Y 应付短期融资款 payables float Y 应付款项 hfs_assets float Y 持有待售的资产 hfs_sales float Y 持有待售的负债 cost_fin_assets float Y 以摊余成本计量的金融资产 fair_value_fin_assets float Y 以公允价值计量且其变动计入其他综合收益的金融资产 cip_total float Y 在建工程(合计)(元) oth_pay_total float Y 其他应付款(合计)(元) long_pay_total float Y 长期应付款(合计)(元) debt_invest float Y 债权投资(元) oth_debt_invest float Y 其他债权投资(元) oth_eq_invest float N 其他权益工具投资(元) oth_illiq_fin_assets float N 其他非流动金融资产(元) oth_eq_ppbond float N 其他权益工具:永续债(元) receiv_financing float N 应收款项融资 use_right_assets float N 使用权资产 lease_liab float N 租赁负债 contract_assets float Y 合同资产 contract_liab float Y 合同负债 accounts_receiv_bill float Y 应收票据及应收账款 accounts_pay float Y 应付票据及应付账款 oth_rcv_total float Y 其他应收款(合计)(元) fix_assets_total float Y 固定资产(合计)(元) update_flag str Y 更新标识 接口使用说明
pro = ts.pro_api()
df = pro.balancesheet(ts_code='600000.SH', start_date='20180101', end_date='20180730', fields='ts_code,ann_date,f_ann_date,end_date,report_type,comp_type,cap_rese') 获取某一季度全部股票数据
df2 = pro.balancesheet_vip(period='20181231',fields='ts_code,ann_date,f_ann_date,end_date,report_type,comp_type,cap_rese') 数据样例
ts_code ann_date f_ann_date end_date report_type comp_type \ 0 600000.SH 20180830 20180830 20180630 1 2 1 600000.SH 20180428 20180428 20180331 1 2
cap_rese 0 8.176000e+10 1 8.176000e+10 主要报表类型说明
| 代码 | 类型 | 说明 |
|---|---|---|
| 1 | 合并报表 | 上市公司最新报表(默认) |
| 2 | 单季合并 | 单一季度的合并报表 |
| 3 | 调整单季合并表 | 调整后的单季合并报表(如果有) |
| 4 | 调整合并报表 | 本年度公布上年同期的财务报表数据,报告期为上年度 |
| 5 | 调整前合并报表 | 数据发生变更,将原数据进行保留,即调整前的原数据 |
| 6 | 母公司报表 | 该公司母公司的财务报表数据 |
| 7 | 母公司单季表 | 母公司的单季度表 |
| 8 | 母公司调整单季表 | 母公司调整后的单季表 |
| 9 | 母公司调整表 | 该公司母公司的本年度公布上年同期的财务报表数据 |
| 10 | 母公司调整前报表 | 母公司调整之前的原始财务报表数据 |
| 11 | 母公司调整前合并报表 | 母公司调整之前合并报表原数据 |
| 12 | 母公司调整前报表 | 母公司报表发生变更前保留的原数据 |