初始化提交
This commit is contained in:
commit
62f7416ee6
|
|
@ -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
|
||||||
|
|
@ -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
|
||||||
|
|
@ -0,0 +1,12 @@
|
||||||
|
<?xml version="1.0" encoding="UTF-8"?>
|
||||||
|
<module type="PYTHON_MODULE" version="4">
|
||||||
|
<component name="NewModuleRootManager">
|
||||||
|
<content url="file://$MODULE_DIR$" />
|
||||||
|
<orderEntry type="inheritedJdk" />
|
||||||
|
<orderEntry type="sourceFolder" forTests="false" />
|
||||||
|
</component>
|
||||||
|
<component name="PyDocumentationSettings">
|
||||||
|
<option name="format" value="PLAIN" />
|
||||||
|
<option name="myDocStringFormat" value="Plain" />
|
||||||
|
</component>
|
||||||
|
</module>
|
||||||
|
|
@ -0,0 +1,39 @@
|
||||||
|
<component name="InspectionProjectProfileManager">
|
||||||
|
<profile version="1.0">
|
||||||
|
<option name="myName" value="Project Default" />
|
||||||
|
<inspection_tool class="PyPackageRequirementsInspection" enabled="true" level="WARNING" enabled_by_default="true">
|
||||||
|
<option name="ignoredPackages">
|
||||||
|
<list>
|
||||||
|
<option value="fastapi" />
|
||||||
|
<option value="uvicorn" />
|
||||||
|
<option value="python-multipart" />
|
||||||
|
<option value="requests" />
|
||||||
|
<option value="httpx" />
|
||||||
|
<option value="asgiref" />
|
||||||
|
<option value="chinesecalendar" />
|
||||||
|
<option value="sqlalchemy" />
|
||||||
|
<option value="pymysql" />
|
||||||
|
<option value="psycopg2-binary" />
|
||||||
|
<option value="redis" />
|
||||||
|
<option value="pymilvus" />
|
||||||
|
<option value="pymongo" />
|
||||||
|
<option value="celery" />
|
||||||
|
<option value="pydantic" />
|
||||||
|
<option value="pydantic-settings" />
|
||||||
|
<option value="python-dotenv" />
|
||||||
|
<option value="pandas" />
|
||||||
|
<option value="numpy" />
|
||||||
|
<option value="scipy" />
|
||||||
|
<option value="mplfinance" />
|
||||||
|
<option value="fastdtw" />
|
||||||
|
<option value="numba" />
|
||||||
|
<option value="PyYAML" />
|
||||||
|
<option value="pyarrow" />
|
||||||
|
<option value="APScheduler" />
|
||||||
|
<option value="duckdb" />
|
||||||
|
<option value="openpyxl" />
|
||||||
|
</list>
|
||||||
|
</option>
|
||||||
|
</inspection_tool>
|
||||||
|
</profile>
|
||||||
|
</component>
|
||||||
|
|
@ -0,0 +1,6 @@
|
||||||
|
<component name="InspectionProjectProfileManager">
|
||||||
|
<settings>
|
||||||
|
<option name="USE_PROJECT_PROFILE" value="false" />
|
||||||
|
<version value="1.0" />
|
||||||
|
</settings>
|
||||||
|
</component>
|
||||||
|
|
@ -0,0 +1,8 @@
|
||||||
|
<?xml version="1.0" encoding="UTF-8"?>
|
||||||
|
<project version="4">
|
||||||
|
<component name="ProjectModuleManager">
|
||||||
|
<modules>
|
||||||
|
<module fileurl="file://$PROJECT_DIR$/.idea/akg-factor-bridge.iml" filepath="$PROJECT_DIR$/.idea/akg-factor-bridge.iml" />
|
||||||
|
</modules>
|
||||||
|
</component>
|
||||||
|
</project>
|
||||||
|
|
@ -0,0 +1,6 @@
|
||||||
|
<?xml version="1.0" encoding="UTF-8"?>
|
||||||
|
<project version="4">
|
||||||
|
<component name="VcsDirectoryMappings">
|
||||||
|
<mapping directory="" vcs="Git" />
|
||||||
|
</component>
|
||||||
|
</project>
|
||||||
|
|
@ -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"]
|
||||||
|
|
@ -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://<akg_user>@<akg_host>: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/日),再决定是否放量。
|
||||||
|
|
@ -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') 过滤,缺列可能查不到本因子,请核实。")
|
||||||
|
|
@ -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("/")
|
||||||
|
|
@ -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
|
||||||
|
|
@ -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
|
||||||
|
|
@ -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,
|
||||||
|
}
|
||||||
|
|
@ -0,0 +1,5 @@
|
||||||
|
pandas>=2.0
|
||||||
|
numpy>=1.24
|
||||||
|
psycopg[binary]>=3.1
|
||||||
|
PyMySQL>=1.1
|
||||||
|
python-dotenv>=1.0
|
||||||
|
|
@ -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()
|
||||||
|
|
@ -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' <> '';
|
||||||
Loading…
Reference in New Issue