From 62f7416ee60e9110b62df7edc999af1e12d6b8bf Mon Sep 17 00:00:00 2001 From: zlt Date: Fri, 24 Jul 2026 14:20:54 +0800 Subject: [PATCH] =?UTF-8?q?=E5=88=9D=E5=A7=8B=E5=8C=96=E6=8F=90=E4=BA=A4?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .env.example | 36 ++++ .idea/.gitignore | 10 ++ .idea/akg-factor-bridge.iml | 12 ++ .idea/inspectionProfiles/Project_Default.xml | 39 +++++ .../inspectionProfiles/profiles_settings.xml | 6 + .idea/modules.xml | 8 + .idea/vcs.xml | 6 + Dockerfile | 10 ++ README.md | 75 ++++++++ common.py | 107 ++++++++++++ config.py | 65 +++++++ db.py | 52 ++++++ docker-compose.yml | 12 ++ factors.py | 165 ++++++++++++++++++ requirements.txt | 5 + run.py | 90 ++++++++++ sql/astock_kg_slot_views.sql | 60 +++++++ 17 files changed, 758 insertions(+) create mode 100644 .env.example create mode 100644 .idea/.gitignore create mode 100644 .idea/akg-factor-bridge.iml create mode 100644 .idea/inspectionProfiles/Project_Default.xml create mode 100644 .idea/inspectionProfiles/profiles_settings.xml create mode 100644 .idea/modules.xml create mode 100644 .idea/vcs.xml create mode 100644 Dockerfile create mode 100644 README.md create mode 100644 common.py create mode 100644 config.py create mode 100644 db.py create mode 100644 docker-compose.yml create mode 100644 factors.py create mode 100644 requirements.txt create mode 100644 run.py create mode 100644 sql/astock_kg_slot_views.sql diff --git a/.env.example b/.env.example new file mode 100644 index 0000000..3000400 --- /dev/null +++ b/.env.example @@ -0,0 +1,36 @@ +# akg-factor-bridge 连接配置。复制为 .env 后填写。桥不硬编码任何主机。 +# 三处连接:读基座视图(PG) · 读热度(153) · 写因子+注册(平台MySQL)。部署在哪台 +# 服务器由"能同时连通这三处"决定。 + +# --- ① astock-kg 基座 PostgreSQL(读四个只读视图,只读账号即可)--- +AKG_PG_HOST= +AKG_PG_PORT=5432 +AKG_PG_USER= +AKG_PG_PASSWORD= +AKG_PG_DB=akg + +# --- ② 153 代理 MySQL(读 stock_fund_heat_scores 热度;即 astock-kg 的 MASTER_MYSQL_*)--- +HEAT_MYSQL_HOST= +HEAT_MYSQL_PORT=3306 +HEAT_MYSQL_USER= +HEAT_MYSQL_PASSWORD= +HEAT_MYSQL_DB= + +# --- ③ 平台因子库 MySQL(PROXY_DB_URL 指向的库:写 t_factor_* + factor_metadata)--- +FACTOR_MYSQL_HOST= +FACTOR_MYSQL_PORT=3306 +FACTOR_MYSQL_USER= +FACTOR_MYSQL_PASSWORD= +FACTOR_MYSQL_DB= + +# --- 现价来源(upside 用;默认复用 ③ 同实例的 gp_day_data,可单独指向)--- +# PRICE_MYSQL_HOST= +# PRICE_MYSQL_PORT=3306 +# PRICE_MYSQL_USER= +# PRICE_MYSQL_PASSWORD= +# PRICE_MYSQL_DB= +# gp_day_data 的股票代码列名(待实机核实:可能是 ts_code 或 symbol) +PRICE_CODE_COL=ts_code + +# 可选:用平台 REST 注册因子时填(留空=直连 ③ 写 factor_metadata) +# FACTOR_API_BASE=http://192.168.16.155:8000 diff --git a/.idea/.gitignore b/.idea/.gitignore new file mode 100644 index 0000000..30cf57e --- /dev/null +++ b/.idea/.gitignore @@ -0,0 +1,10 @@ +# Default ignored files +/shelf/ +/workspace.xml +# Editor-based HTTP Client requests +/httpRequests/ +# Ignored default folder with query files +/queries/ +# Datasource local storage ignored files +/dataSources/ +/dataSources.local.xml diff --git a/.idea/akg-factor-bridge.iml b/.idea/akg-factor-bridge.iml new file mode 100644 index 0000000..8b8c395 --- /dev/null +++ b/.idea/akg-factor-bridge.iml @@ -0,0 +1,12 @@ + + + + + + + + + + \ No newline at end of file diff --git a/.idea/inspectionProfiles/Project_Default.xml b/.idea/inspectionProfiles/Project_Default.xml new file mode 100644 index 0000000..c492292 --- /dev/null +++ b/.idea/inspectionProfiles/Project_Default.xml @@ -0,0 +1,39 @@ + + + + \ No newline at end of file diff --git a/.idea/inspectionProfiles/profiles_settings.xml b/.idea/inspectionProfiles/profiles_settings.xml new file mode 100644 index 0000000..105ce2d --- /dev/null +++ b/.idea/inspectionProfiles/profiles_settings.xml @@ -0,0 +1,6 @@ + + + + \ No newline at end of file diff --git a/.idea/modules.xml b/.idea/modules.xml new file mode 100644 index 0000000..8bb1b52 --- /dev/null +++ b/.idea/modules.xml @@ -0,0 +1,8 @@ + + + + + + + + \ No newline at end of file diff --git a/.idea/vcs.xml b/.idea/vcs.xml new file mode 100644 index 0000000..35eb1dd --- /dev/null +++ b/.idea/vcs.xml @@ -0,0 +1,6 @@ + + + + + + \ No newline at end of file diff --git a/Dockerfile b/Dockerfile new file mode 100644 index 0000000..cd4a99c --- /dev/null +++ b/Dockerfile @@ -0,0 +1,10 @@ +FROM python:3.11-slim + +WORKDIR /app +COPY requirements.txt . +RUN pip install --no-cache-dir -r requirements.txt + +COPY . . + +# 默认空跑连通性自检;生产由 cron / 平台 XXL-JOB / docker exec 触发 build。 +CMD ["python", "run.py", "views"] diff --git a/README.md b/README.md new file mode 100644 index 0000000..7098a46 --- /dev/null +++ b/README.md @@ -0,0 +1,75 @@ +# akg-factor-bridge + +astock-kg(知识图谱基座)↔ `quant_factor_service`(通用因子平台)之间的**因子导出桥**。 +把基座的四路产出(分析师预期空间 / 热度 / 利好利空事件 / 板块传导)做成符合平台规范的 +子因子,写进平台因子库,供平台合成为选股因子。 + +> 设计与决策依据见 astock-kg `docs/量化因子导出与合成设计.md`。本工程**不 import** +> 基座或平台任何代码,只靠 `.env` 里三处数据库连接工作,可单独部署于任意能连通三库的服务器。 + +## 架构(基座出视图,桥算变换,平台算合成) + +``` +astock-kg 基座 PG ── 只读视图(sql/astock_kg_slot_views.sql) + v_factor_universe / v_factor_consensus / v_factor_events / v_factor_transmission + │ (热度不经基座:桥直连 153 读 stock_fund_heat_scores) + ▼ +akg-factor-bridge:读视图+热度 → 算四路日截面(极性/衰减/打分) → 转前缀码 SH600000 + → 写平台 t_factor_akg_* + 注册 factor_metadata + ▼ +平台 quant_factor_service:把四子因子当普通 single 因子 → 合成/回测/调度 +``` + +- **基座只暴露数据、不算因子**;桥承载因子建模(极性表、衰减、打分——见 `factors.py`); + 跨信号 alpha 组合在平台侧。 +- universe = KG 覆盖池(`v_factor_universe` = 全部 `industry_pools` 成员并集),四路都限制其内。 + +## 四个子因子 + +| factor_code | 表 | 口径 | 缺失 | +|---|---|---|---| +| `akg_upside` | t_factor_akg_upside | 目标价中枢/现价−1(as-of) | NaN(无覆盖不出行) | +| `akg_heat` | t_factor_akg_heat | 热度分 0~1(最新批次) | NaN | +| `akg_event` | t_factor_akg_event | Σ 事件极性×时间衰减 | 0(无事件=中性) | +| `akg_transmission` | t_factor_akg_transmission | 路径数×(1−已动比例) | 0 | + +## 用法 + +```bash +# 0) 先把基座视图建好(在 astock-kg 的 PG 上执行一次) +psql "postgresql://@:5432/akg" -f sql/astock_kg_slot_views.sql + +# 1) 配置连接 +cp .env.example .env && vim .env # 填三处连接 + gp_day_data 代码列 + +# 2) 连通性自检(四视图 / 热度 / gp_day_data / factor_metadata 行数) +pip install -r requirements.txt +python run.py views + +# 3) 注册四子因子 +python run.py register + +# 4) 历史回填 / 每日增量(幂等,可重跑) +python run.py build all --mode history --start 2024-01-01 --end 2026-07-24 +python run.py build all --mode daily --date 2026-07-24 +``` + +容器化:`docker compose up -d`(常驻),宿主 cron/平台 XXL-JOB 以 +`docker exec akg_factor_bridge python run.py build all --mode daily` 触发。 + +## ⚠️ 待实机核实项(本工程 DB 细节以线上为准,跑不通按此排查) + +1. **三库连通性**:`python run.py views` 六项全 ✅ 才算通。任一 ❌ 先解决网络/账号 + (尤其桥所在服务器到基座 PG、153、平台 MySQL 的可达性)。 +2. **`gp_day_data` 代码列与形态**:`upside` 现价来自它。列名可能是 `ts_code` 或 + `symbol`(`.env` 的 `PRICE_CODE_COL`);代码形态(`600000.SH` / `SH600000` / `600000`) + 两边已统一折前缀式再 join——若 `upside` 出行为 0,多半是形态没对上,在此调 join 口径。 +3. **`factor_metadata` 列**:`register` 自适应实际列写入;若无 `factor_type` 列,平台 + `/mining/factors/all` 可能查不到本因子(会打印告警),需与平台侧确认。 +4. **事件极性/方向(§8-3 开放问题)**:`factors.EVENT_POLARITY / EVENT_DIR` 是草案; + `增减持` 的增/减方向若 `qualifiers.direction` 里没有(当前视图取 direction),会落 0, + 需确认基座 EVENT 的方向到底存在哪个 qualifier 键。半衰期 10 交易日 / 窗口 60 交易日可调。 +5. **事件交易日龄近似**:v1 用自然日×(5/7) 折算交易日龄,非精确交易日历——够用,后续可 + 换真实交易日历向量化。 +6. **universe 覆盖面(F0 前置)**:先用 astock-kg 的 `factor_coverage_probe.py` 确认池内 + 四信号日截面覆盖数(尤其热度 ≥30~50/日),再决定是否放量。 diff --git a/common.py b/common.py new file mode 100644 index 0000000..d942791 --- /dev/null +++ b/common.py @@ -0,0 +1,107 @@ +"""共用:股票码规范化、覆盖池、幂等写因子表、注册 factor_metadata(自适应列)。""" +import json + +import pandas as pd + +import db + + +def to_prefix(ts_code: str) -> str: + """600000.SH -> SH600000;已是前缀式(SZ002625)或纯代码则原样返回。""" + s = str(ts_code).strip() + if "." in s: + num, exch = s.split(".", 1) + return f"{exch.upper()}{num}" + return s + + +def load_universe() -> set[str]: + """KG 覆盖池 ts_code 集合(600000.SH 形态),来自基座视图 v_factor_universe。""" + df = db.read_pg("SELECT ts_code FROM v_factor_universe") + return set(df["ts_code"].astype(str).str.strip()) + + +def trading_days(start: str, end: str) -> list: + """目标区间交易日历——用热度表(153,日频、覆盖广)的 distinct trade_date 近似。""" + df = db.read_mysql("heat", + "SELECT DISTINCT trade_date FROM stock_fund_heat_scores " + "WHERE trade_date BETWEEN %s AND %s ORDER BY trade_date", + (start, end)) + return list(pd.to_datetime(df["trade_date"])) + + +_CREATE_FACTOR_TABLE = """ +CREATE TABLE IF NOT EXISTS {t} ( + trade_date DATE NOT NULL, + stock_code VARCHAR(15) NOT NULL, + factor_value DOUBLE, + created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP, + PRIMARY KEY (trade_date, stock_code), + KEY idx_stock (stock_code) +) ENGINE=InnoDB +""" + + +def ensure_table(table: str) -> None: + with db.factor_conn() as conn: + with conn.cursor() as cur: + cur.execute(_CREATE_FACTOR_TABLE.format(t=table)) + conn.commit() + + +def write_factor(table: str, df: pd.DataFrame, mode: str) -> None: + """df 列 [trade_date, stock_code, factor_value] → 幂等写因子表。 + daily/history 都是「删涉及日期区间 → 批插」。stock_code 统一转前缀式。""" + if df is None or df.empty: + print(f" {table}: 无数据(跳过)") + return + df = df.dropna(subset=["trade_date", "stock_code", "factor_value"]).copy() + df["stock_code"] = df["stock_code"].map(to_prefix) + df["trade_date"] = pd.to_datetime(df["trade_date"]).dt.date + df = df.drop_duplicates(["trade_date", "stock_code"], keep="last") + if df.empty: + print(f" {table}: 清洗后无数据") + return + dmin, dmax = df["trade_date"].min(), df["trade_date"].max() + ensure_table(table) + rows = list(df[["trade_date", "stock_code", "factor_value"]] + .itertuples(index=False, name=None)) + with db.factor_conn() as conn: + with conn.cursor() as cur: + cur.execute(f"DELETE FROM {table} WHERE trade_date BETWEEN %s AND %s", (dmin, dmax)) + cur.executemany( + f"INSERT INTO {table} (trade_date, stock_code, factor_value) VALUES (%s,%s,%s)", + rows) + conn.commit() + print(f" {table}: 写入 {len(rows)} 行, 日期 {dmin}~{dmax}") + + +def register(factor_code, display_name, table, category, desc, + author="akg-factor-bridge") -> None: + """注册 factor_metadata。自适应实际存在的列(避免猜死 schema)—— + 列名以线上 factor_metadata 为准,本函数只写它有的列。""" + values = { + "factor_code": factor_code, "display_name": display_name, + "target_ds_name": "ds_a", "target_table_name": table, + "author": author, "description": desc, + "category": json.dumps(category, ensure_ascii=False), + "frequency": "daily", "status": "active", "factor_type": "single", + } + with db.factor_conn() as conn: + with conn.cursor() as cur: + cur.execute("SELECT column_name FROM information_schema.columns " + "WHERE table_schema=DATABASE() AND table_name='factor_metadata'") + cols = {r[0] for r in cur.fetchall()} + use = [k for k in values if k in cols] + if "factor_code" not in use: + raise RuntimeError("factor_metadata 无 factor_code 列?请核实线上 schema") + ph = ",".join(["%s"] * len(use)) + upd = ",".join(f"{k}=VALUES({k})" for k in use if k != "factor_code") + cur.execute( + f"INSERT INTO factor_metadata ({','.join(use)}) VALUES ({ph}) " + f"ON DUPLICATE KEY UPDATE {upd}", [values[k] for k in use]) + conn.commit() + print(f" 注册 {factor_code} -> {table}(写入列: {sorted(use)})") + if "factor_type" not in cols: + print(" ⚠️ factor_metadata 无 factor_type 列——平台 /mining/factors/all 按 " + "factor_type IN('single','multiple') 过滤,缺列可能查不到本因子,请核实。") diff --git a/config.py b/config.py new file mode 100644 index 0000000..2459dfa --- /dev/null +++ b/config.py @@ -0,0 +1,65 @@ +"""连接配置:全部从环境变量/.env 读取,不硬编码任何主机。 + +三处连接: + AKG_PG_* —— astock-kg 基座 PostgreSQL(读四个只读视图,只读账号即可) + HEAT_MYSQL_* —— 153 代理 MySQL(读 stock_fund_heat_scores 热度) + FACTOR_MYSQL_* —— 平台 PROXY_DB_URL 指向的 MySQL(写 t_factor_* + factor_metadata) +现价来源 PRICE_MYSQL_*(upside 用)默认复用 FACTOR_MYSQL_* 同实例。 +""" +import os +from dataclasses import dataclass + +try: + from dotenv import load_dotenv + load_dotenv() +except Exception: + pass + + +def _req(k: str) -> str: + v = os.environ.get(k) + if not v: + raise RuntimeError(f".env 缺配置: {k}") + return v + + +def _opt(k: str, fallback: str) -> str: + return os.environ.get(k) or fallback + + +@dataclass(frozen=True) +class Conn: + host: str + port: int + user: str + password: str + db: str + + +def akg_pg() -> Conn: + return Conn(_req("AKG_PG_HOST"), int(os.environ.get("AKG_PG_PORT", 5432)), + _req("AKG_PG_USER"), os.environ.get("AKG_PG_PASSWORD", ""), _req("AKG_PG_DB")) + + +def heat_mysql() -> Conn: + return Conn(_req("HEAT_MYSQL_HOST"), int(os.environ.get("HEAT_MYSQL_PORT", 3306)), + _req("HEAT_MYSQL_USER"), os.environ.get("HEAT_MYSQL_PASSWORD", ""), _req("HEAT_MYSQL_DB")) + + +def factor_mysql() -> Conn: + return Conn(_req("FACTOR_MYSQL_HOST"), int(os.environ.get("FACTOR_MYSQL_PORT", 3306)), + _req("FACTOR_MYSQL_USER"), os.environ.get("FACTOR_MYSQL_PASSWORD", ""), _req("FACTOR_MYSQL_DB")) + + +def price_mysql() -> Conn: + """现价(gp_day_data)来源,默认与因子库同实例。""" + return Conn(_opt("PRICE_MYSQL_HOST", _req("FACTOR_MYSQL_HOST")), + int(_opt("PRICE_MYSQL_PORT", os.environ.get("FACTOR_MYSQL_PORT", "3306"))), + _opt("PRICE_MYSQL_USER", _req("FACTOR_MYSQL_USER")), + _opt("PRICE_MYSQL_PASSWORD", os.environ.get("FACTOR_MYSQL_PASSWORD", "")), + _opt("PRICE_MYSQL_DB", _req("FACTOR_MYSQL_DB"))) + + +# gp_day_data 股票代码列名(待实机核实:ts_code 或 symbol) +PRICE_CODE_COL = os.environ.get("PRICE_CODE_COL", "ts_code") +FACTOR_API_BASE = os.environ.get("FACTOR_API_BASE", "").rstrip("/") diff --git a/db.py b/db.py new file mode 100644 index 0000000..83e922e --- /dev/null +++ b/db.py @@ -0,0 +1,52 @@ +"""连接与读写 IO。PG 用 psycopg(3),MySQL 用 pymysql。桥不 import 基座/平台代码。""" +from contextlib import contextmanager + +import pandas as pd +import psycopg +import pymysql + +import config + + +@contextmanager +def akg_pg_conn(): + c = config.akg_pg() + conn = psycopg.connect(host=c.host, port=c.port, user=c.user, + password=c.password, dbname=c.db) + try: + yield conn + finally: + conn.close() + + +def _mysql(c: "config.Conn"): + return pymysql.connect(host=c.host, port=c.port, user=c.user, password=c.password, + database=c.db, charset="utf8mb4", read_timeout=180) + + +@contextmanager +def _mysql_cm(which: str): + c = {"heat": config.heat_mysql, "factor": config.factor_mysql, + "price": config.price_mysql}[which]() + conn = _mysql(c) + try: + yield conn + finally: + conn.close() + + +def read_pg(sql: str, params=None) -> pd.DataFrame: + with akg_pg_conn() as conn: + return pd.read_sql(sql, conn, params=params) + + +def read_mysql(which: str, sql: str, params=None) -> pd.DataFrame: + with _mysql_cm(which) as conn: + return pd.read_sql(sql, conn, params=params) + + +@contextmanager +def factor_conn(): + """写因子库用(需要游标提交)。""" + with _mysql_cm("factor") as conn: + yield conn diff --git a/docker-compose.yml b/docker-compose.yml new file mode 100644 index 0000000..cecb832 --- /dev/null +++ b/docker-compose.yml @@ -0,0 +1,12 @@ +# akg-factor-bridge:独立部署单元,可落在任意能同时连通「基座PG/153/平台MySQL」的服务器。 +# 容器常驻(sleep infinity),由宿主 cron 或平台 XXL-JOB 以 docker exec 触发 build; +# 也可改 command 为一次性任务由外部调度拉起。 +services: + akg-factor-bridge: + build: . + container_name: akg_factor_bridge + env_file: .env + restart: unless-stopped + command: sleep infinity + # 触发示例(宿主 crontab,每日 18:40): + # 40 18 * * 1-5 docker exec akg_factor_bridge python run.py build all --mode daily diff --git a/factors.py b/factors.py new file mode 100644 index 0000000..9344f72 --- /dev/null +++ b/factors.py @@ -0,0 +1,165 @@ +"""四路子因子构造。输入日期区间 [start, end](YYYY-MM-DD),输出 +DataFrame[trade_date, stock_code, factor_value],全部限制在 KG 覆盖池 universe 内。 +stock_code 输出形态不限,common.write_factor 统一转前缀式。 + +建模参数(EVENT_POLARITY / *_HALF_LIFE / *_WINDOW)是**因子决策**,见设计文档 +§3.3 与 §8-3,此处取草案默认,待用户确认后调。方向(direction)统一在合成配置里声明, +子因子表只存原始值——但 event 天然带符号,故此处保留极性。 +""" +import numpy as np +import pandas as pd + +import config +import common +import db + +FACTORS = { + "akg_upside": "t_factor_akg_upside", + "akg_heat": "t_factor_akg_heat", + "akg_event": "t_factor_akg_event", + "akg_transmission": "t_factor_akg_transmission", +} + +# ---- 事件极性草案(§3.3,待 §8-3 确认)---- +EVENT_POLARITY = { + "股份回购": 1.0, "重大合同中标": 1.0, "股权激励授予": 0.5, + "诉讼仲裁": -1.0, "行政处罚": -1.0, "股权质押": -0.5, + "发行上市": 0.0, "并购交割": 0.0, "other": 0.0, + # 增减持 / 业绩预告:符号取决于 direction(下 EVENT_DIR) + "增减持": 0.0, "业绩预告": 0.0, +} +EVENT_DIR = {"预增": 1.0, "预减": -1.0, "增持": 1.0, "减持": -1.0} +EVENT_HALF_LIFE = 10 # 交易日 +EVENT_WINDOW = 60 # 交易日(超窗不计) +_EMPTY = pd.DataFrame(columns=["trade_date", "stock_code", "factor_value"]) + + +# ---------------------------------------------------------------- 热度 +def build_heat(start, end): + """热度 = stock_fund_heat_scores 最新批次 score(0~1)。stock_code 已前缀式。""" + uni = {common.to_prefix(x) for x in common.load_universe()} + df = db.read_mysql("heat", + """SELECT s.trade_date, s.stock_code, s.score AS factor_value + FROM stock_fund_heat_scores s + JOIN (SELECT trade_date, MAX(batch_no) bn FROM stock_fund_heat_scores + WHERE trade_date BETWEEN %s AND %s GROUP BY trade_date) m + ON m.trade_date = s.trade_date AND m.bn = s.batch_no""", + (start, end)) + if df.empty: + return _EMPTY + df["stock_code"] = df["stock_code"].astype(str).str.strip() + df["factor_value"] = pd.to_numeric(df["factor_value"], errors="coerce") + df = df[df["stock_code"].isin(uni)] + return df[["trade_date", "stock_code", "factor_value"]] + + +# ---------------------------------------------------------------- 预期空间 +def build_upside(start, end): + """upside = 一致预期目标价中枢 / 当日现价 − 1(as-of:现价日取 asof<=当日最新一致预期)。 + 现价来自平台 gp_day_data(代码列 config.PRICE_CODE_COL,待实机核实)。""" + uni = common.load_universe() + cons = db.read_pg( + "SELECT ts_code, asof_date, target_mid_avg FROM v_factor_consensus WHERE asof_date <= %s", + (end,)) + cons = cons[cons["ts_code"].isin(uni)].copy() + if cons.empty: + return _EMPTY + col = config.PRICE_CODE_COL + price = db.read_mysql("price", + f"SELECT `timestamp` AS trade_date, `{col}` AS ts_code, close " + f"FROM gp_day_data WHERE `timestamp` BETWEEN %s AND %s", (start, end)) + if price.empty: + return _EMPTY + price["close"] = pd.to_numeric(price["close"], errors="coerce") + # 归一到前缀式两边对齐(gp_day_data 代码形态不定 → 都折前缀式后 join) + price["k"] = price["ts_code"].map(common.to_prefix) + price = price[(price["close"] > 0)].dropna(subset=["close"]) + cons["k"] = cons["ts_code"].map(common.to_prefix) + cons = cons[cons["k"].isin(set(price["k"]))] + if cons.empty: + return _EMPTY + cons["asof_date"] = pd.to_datetime(cons["asof_date"]) + price["trade_date"] = pd.to_datetime(price["trade_date"]) + out = [] + cons_sorted = cons.sort_values("asof_date") + for k, pg in price.groupby("k"): + cg = cons_sorted[cons_sorted["k"] == k] + if cg.empty: + continue + m = pd.merge_asof(pg.sort_values("trade_date"), + cg[["asof_date", "target_mid_avg"]], + left_on="trade_date", right_on="asof_date", direction="backward") + m = m.dropna(subset=["target_mid_avg"]) + if m.empty: + continue + m["factor_value"] = m["target_mid_avg"].astype(float) / m["close"] - 1.0 + m["stock_code"] = k + out.append(m[["trade_date", "stock_code", "factor_value"]]) + return pd.concat(out) if out else _EMPTY + + +# ---------------------------------------------------------------- 事件 +def _polarity(event_type, direction): + d = (direction or "").strip() + if d in EVENT_DIR: + return EVENT_DIR[d] + return EVENT_POLARITY.get(event_type, 0.0) + + +def build_event(start, end): + """事件分 = Σ 近窗口内事件 极性 × 时间衰减(exp(-交易日龄·ln2/半衰期))。 + ts_code 取文档锚(v_factor_events 已解析)。无事件的股当天不出行(= 缺 → 合成侧填 0)。""" + uni = common.load_universe() + look = (pd.Timestamp(start) - pd.Timedelta(days=EVENT_WINDOW * 2)).date() + ev = db.read_pg( + "SELECT ts_code, disclosure_date, event_type, direction " + "FROM v_factor_events WHERE disclosure_date BETWEEN %s AND %s", (look, end)) + ev = ev[ev["ts_code"].isin(uni)].copy() + if ev.empty: + return _EMPTY + ev["pol"] = [_polarity(t, d) for t, d in zip(ev["event_type"], ev["direction"])] + ev = ev[ev["pol"] != 0.0] + if ev.empty: + return _EMPTY + cal = common.trading_days(start, end) + if not cal: + return _EMPTY + cal = pd.DatetimeIndex(cal) + decay = np.log(2) / EVENT_HALF_LIFE + rows = [] + for ts, g in ev.groupby("ts_code"): + disc = pd.to_datetime(g["disclosure_date"]).values.astype("datetime64[ns]") + pol = g["pol"].to_numpy(dtype=float) + for d in cal: + # 自然日龄 → 交易日龄近似 ×(5/7)(v1 近似,见 README 待优化项) + age_td = ((d.value - disc.astype("int64")) / 86_400e9) * (5.0 / 7.0) + mask = (age_td >= 0) & (age_td <= EVENT_WINDOW) + if not mask.any(): + continue + val = float((pol[mask] * np.exp(-age_td[mask] * decay)).sum()) + if val != 0.0: + rows.append((d.date(), ts, val)) + return pd.DataFrame(rows, columns=["trade_date", "stock_code", "factor_value"]) if rows else _EMPTY + + +# ---------------------------------------------------------------- 传导 +def build_transmission(start, end): + """传导分 = 指向该股所在环节的路径数 ×(1 − 已动比例);同股同日多候选取最大。""" + uni = common.load_universe() + tr = db.read_pg( + "SELECT scan_date, ts_code, n_paths, moved_ratio " + "FROM v_factor_transmission WHERE scan_date BETWEEN %s AND %s", (start, end)) + tr = tr[tr["ts_code"].isin(uni)].copy() + if tr.empty: + return _EMPTY + tr["factor_value"] = (tr["n_paths"].astype(float) + * (1.0 - pd.to_numeric(tr["moved_ratio"], errors="coerce").fillna(0.0))) + g = (tr.groupby(["scan_date", "ts_code"])["factor_value"].max().reset_index() + .rename(columns={"scan_date": "trade_date", "ts_code": "stock_code"})) + return g[["trade_date", "stock_code", "factor_value"]] + + +BUILDERS = { + "akg_upside": build_upside, "akg_heat": build_heat, + "akg_event": build_event, "akg_transmission": build_transmission, +} diff --git a/requirements.txt b/requirements.txt new file mode 100644 index 0000000..f583d31 --- /dev/null +++ b/requirements.txt @@ -0,0 +1,5 @@ +pandas>=2.0 +numpy>=1.24 +psycopg[binary]>=3.1 +PyMySQL>=1.1 +python-dotenv>=1.0 diff --git a/run.py b/run.py new file mode 100644 index 0000000..6bad855 --- /dev/null +++ b/run.py @@ -0,0 +1,90 @@ +"""akg-factor-bridge CLI。 + + python run.py views # 连通性自检:打印视图/表行数 + python run.py register # 注册四子因子到 factor_metadata + python run.py build all --mode history --start 2024-01-01 --end 2025-12-31 + python run.py build akg_heat --mode daily --date 2026-07-24 + python run.py build akg_event --mode history --start 2025-01-01 --end 2026-07-24 + +daily 模式:不给 --date 则取今天;start=end=date。history 模式:需 --start/--end。 +所有写入幂等(删涉及日期区间再插),可安全重跑。 +""" +import argparse +import datetime as dt + +import common +import db +import factors + + +def cmd_views(): + checks = [ + ("PG v_factor_universe", "pg", "SELECT count(*) FROM v_factor_universe"), + ("PG v_factor_consensus", "pg", "SELECT count(*) FROM v_factor_consensus"), + ("PG v_factor_events", "pg", "SELECT count(*) FROM v_factor_events"), + ("PG v_factor_transmission", "pg", "SELECT count(*) FROM v_factor_transmission"), + ("153 stock_fund_heat_scores","heat", "SELECT count(*) FROM stock_fund_heat_scores"), + ("平台 gp_day_data", "price", "SELECT count(*) FROM gp_day_data"), + ("平台 factor_metadata", "factor", "SELECT count(*) FROM factor_metadata"), + ] + print("连通性自检:") + for name, src, sql in checks: + try: + df = db.read_pg(sql) if src == "pg" else db.read_mysql(src, sql) + print(f" ✅ {name}: {int(df.iloc[0, 0])}") + except Exception as e: # noqa: BLE001 + print(f" ❌ {name}: {e!r}") + + +_META = { + "akg_upside": ("astock-kg 预期空间", "分析师一致预期目标价隐含收益率(target_mid/price-1)"), + "akg_heat": ("astock-kg 热度", "生态日频资金热度分(0~1)"), + "akg_event": ("astock-kg 事件", "利好利空事件时间衰减加权分"), + "akg_transmission": ("astock-kg 传导", "板块传导未动成员传导强度(路径数×(1-已动比例))"), +} + + +def cmd_register(): + print("注册四子因子:") + for code, (name, desc) in _META.items(): + common.register(code, name, factors.FACTORS[code], ["astock-kg", code.split("_", 1)[1]], desc) + + +def cmd_build(which, mode, start, end, date): + if mode == "daily": + d = date or dt.date.today().isoformat() + start = end = d + if not start or not end: + raise SystemExit("history 模式需要 --start 与 --end") + codes = list(factors.FACTORS) if which == "all" else [which] + for code in codes: + if code not in factors.BUILDERS: + raise SystemExit(f"未知因子: {code}(可选: {list(factors.FACTORS)} 或 all)") + print(f"[{code}] {mode} {start} ~ {end}") + df = factors.BUILDERS[code](start, end) + common.write_factor(factors.FACTORS[code], df, mode) + + +def main(): + ap = argparse.ArgumentParser(description="akg-factor-bridge") + sub = ap.add_subparsers(dest="cmd", required=True) + sub.add_parser("views") + sub.add_parser("register") + b = sub.add_parser("build") + b.add_argument("factor", help="akg_upside|akg_heat|akg_event|akg_transmission|all") + b.add_argument("--mode", choices=["daily", "history"], default="daily") + b.add_argument("--start") + b.add_argument("--end") + b.add_argument("--date") + a = ap.parse_args() + + if a.cmd == "views": + cmd_views() + elif a.cmd == "register": + cmd_register() + elif a.cmd == "build": + cmd_build(a.factor, a.mode, a.start, a.end, a.date) + + +if __name__ == "__main__": + main() diff --git a/sql/astock_kg_slot_views.sql b/sql/astock_kg_slot_views.sql new file mode 100644 index 0000000..6ca1289 --- /dev/null +++ b/sql/astock_kg_slot_views.sql @@ -0,0 +1,60 @@ +-- ============================================================================ +-- astock-kg 因子插槽接口:只读视图(供 akg-factor-bridge 消费) +-- ---------------------------------------------------------------------------- +-- 应用到 astock-kg 的 PostgreSQL(akg 库): psql "$AKG_PG_DSN" -f astock_kg_slot_views.sql +-- 只读投影、零新增计算;基座内部表可自由重构,只要这四个视图的列不变,桥不受影响。 +-- 幂等:CREATE OR REPLACE。列/JSON 键均已从 astock-kg 代码核准(见每条注释出处)。 +-- ============================================================================ + +-- 0) 覆盖池 universe = 全部 industry_pools 成员 ts_code 并集 +-- (对齐 market_snapshot._pool_ts_codes;members 为 JSONB 数组,元素含 ts_code/name) +CREATE OR REPLACE VIEW v_factor_universe AS +SELECT DISTINCT m->>'ts_code' AS ts_code +FROM industry_pools p +CROSS JOIN LATERAL jsonb_array_elements(COALESCE(p.members, '[]'::jsonb)) AS m +WHERE m->>'ts_code' IS NOT NULL AND m->>'ts_code' <> ''; + +-- 1) 一致预期(供 upside):consensus_daily 全 asof 行投影,桥按 trade_date 做 as-of。 +-- 列源:market_snapshot._CONSENSUS_DDL。 +CREATE OR REPLACE VIEW v_factor_consensus AS +SELECT ts_code, + asof_date, + target_mid_avg, + eps_med, + np_med, + n_orgs, + n_reports_90d +FROM consensus_daily +WHERE target_mid_avg IS NOT NULL; + +-- 2) 利好利空事件(供 event):EVENT 断言投影。 +-- ts_code 取【文档锚】documents.meta->>'company_ts_code'(公告单公司、100% 可靠; +-- claim_store 多处以此为公司锚),不用 claims.subject_id(抽取原始名/码、未解析)。 +-- event_type / direction / scope 从 qualifiers 取,键名核自 pipeline._event_to_claim: +-- qualifiers.event_type(EVENT_TYPES 或 other) +-- qualifiers.direction = 预增/预减(业绩预告必填;其余多为空) +-- qualifiers.scope = 累计/单次(回购、增减持) +CREATE OR REPLACE VIEW v_factor_events AS +SELECT d.meta->>'company_ts_code' AS ts_code, + c.disclosure_date, + c.qualifiers->>'event_type' AS event_type, + c.qualifiers->>'direction' AS direction, + c.qualifiers->>'scope' AS scope, + c.confidence, + c.dedup_key +FROM claims c +JOIN documents d ON d.doc_id = c.doc_id +WHERE c.predicate = 'EVENT' + AND d.meta->>'company_ts_code' IS NOT NULL; + +-- 3) 板块传导(供 transmission):transmission_candidates 的未动成员(quiet)摊平成每股一行。 +-- 一只股当日可能出现在多个候选(属多个目标环节)→ 多行,桥侧聚合(取最大 n_paths)。 +-- 列源:transmission._DDL(paths/quiet 均 JSONB,moved_ratio NUMERIC)。 +CREATE OR REPLACE VIEW v_factor_transmission AS +SELECT t.scan_date, + q->>'ts_code' AS ts_code, + jsonb_array_length(COALESCE(t.paths, '[]'::jsonb)) AS n_paths, -- 指向该环节的传导路径数 + t.moved_ratio -- 该环节已动成员比例 +FROM transmission_candidates t +CROSS JOIN LATERAL jsonb_array_elements(COALESCE(t.quiet, '[]'::jsonb)) AS q +WHERE q->>'ts_code' IS NOT NULL AND q->>'ts_code' <> '';