量化数据进阶:多维实战篇(第 10 篇):板块轮动:行业资金流向与龙头股捕捉
一、前言
上一篇《自研指标:量价背离、换手异动与振幅》我们完成了「模块三·技术指标」的收官,也完成了一次转身:只拿裸字段,纯 Python 手写因子。但那一篇的视角始终停在一只股票上——价格分位、量能分位、换手 Z-Score,全部是「个股自身的时间序列」。
真实盘面里,个股很少独自上涨。资金通常先涌入某个行业,再由行业内的龙头股把涨幅传导开来,这就是所谓的板块轮动。只看个股,你永远不知道「为什么涨」;只看指数,你又抓不到「谁在涨」。
本篇进入「模块四·板块与资金」,用三个接口把视角从「个股」抬升到「行业」:
- hizj/zjh:证监会行业资金流,一次返回全市场各行业的近 3 / 5 / 10 日涨跌幅与净流额;
- hibk/zjhhy:证监会行业板块,一次返回各板块的均价、涨跌幅、流入流出资金、净流入,以及该板块的龙头股;
- higg/jlr:个股净流入榜,用于对行业龙头股做二次验证。
三者结合,我们自研「资金流向排名」「轮动强度」「多周期资金共振」「净流率背离」「龙头股捕捉」「板块→个股穿透」六组指标,把「哪只股票在涨」升级为「哪条资金链在动」。
衔接说明:本篇是「模块四·板块与资金」的开篇。上一篇解决了「个股层面的自研因子」,本篇解决「行业层面的资金与龙头」,下一篇将回到风控,补齐本系列最后一块拼图。
郑重说明:本文所有排名、打标、共振、背离、验证逻辑全部是自研因子,不是必盈接口提供的功能。三个接口只提供裸字段,所有排名、分位、象限分类、同号共振、符号背离、龙头校验均由本文代码在本地算出。
二、接口字段解读
2.1 证监会行业资金流 hizj/zjh
http://api.biyingapi.com/hizj/zjh/{LICENCE}
路径参数:
| LICENCE | 您的 licence | 鉴权串 |
返回字段(仅以下字段):
| t | string | 时间 |
| mc | string | 行业名称 |
| dm | string | 行业代码 |
| ac3 | number | 近 3 日涨跌幅(%) |
| net3 | number | 近 3 日净流额 |
| ra3 | number | 近 3 日净流率(%) |
| ac5 | number | 近 5 日涨跌幅(%) |
| net5 | number | 近 5 日净流额 |
| ra5 | number | 近 5 日净流率(%) |
| ac10 | number | 近 10 日涨跌幅(%) |
| net10 | number | 近 10 日净流额 |
| ra10 | number | 近 10 日净流率(%) |
说明:这个接口的价值在于同时给了「价」与「钱」两组三周期数据——ac3/ac5/ac10 是价格结果,net3/net5/net10 是资金原因,ra3/ra5/ra10 是把资金按流通规模归一化后的比率。三周期 × 三口径 = 九个数,多周期共振与背离就建立在这九个数之上。
2.2 证监会行业板块 hibk/zjhhy
http://api.biyingapi.com/hibk/zjhhy/{LICENCE}
返回字段(仅以下字段):
| t | string | 时间 |
| mc | string | 板块名称 |
| dm | string | 板块代码 |
| jj | number | 均价(元) |
| zdf | number | 涨跌幅(%) |
| lrzj | number | 流入资金 |
| lczj | number | 流出资金 |
| jlr | number | 净流入 |
| jlrl | number | 净流入率(%) |
| lzgmc | string | 龙头股名称 |
| lzgdm | string | 龙头股代码 |
| lzgjlrl | number | 龙头股净流入率(%) |
说明:lrzj – lczj = jlr,接口把流入、流出、净额三个数都给了,可以直接校验数据自洽性。更关键的是 lzgmc / lzgdm / lzgjlrl——接口直接告诉你每个板块当前资金的「旗手」是谁,这是本篇龙头股捕捉的数据基础。
2.3 个股净流入 higg/jlr
http://api.biyingapi.com/higg/jlr/{LICENCE}
返回字段(仅以下字段):
| t | string | 时间 |
| mc | string | 名称 |
| dm | string | 代码 |
| zxj | number | 最新价(元) |
| zdf | number | 涨跌幅(%) |
| hsl | number | 换手率(%) |
| cje | number | 成交额 |
| lczj | number | 流出资金 |
| lrzj | number | 流入资金 |
| jlr | number | 净流入 |
| jlrl | number | 净流入率(%) |
注意:三个接口都有 t、mc、dm、zdf、jlr、jlrl 等同名字段,但 hizj/zjh 的 dm 是行业代码、hibk/zjhhy 的 dm 是板块代码、higg/jlr 的 dm 是股票代码——同名不同域。本文在函数中显式区分,跨表关联时必须带上语义前缀,绝不把行业代码与股票代码混为一谈。
三、自研衍生指标 / 业务逻辑
接口只给裸字段,以下全部用 Python 手写(均为自研因子,非接口能力):
四、数据表设计
— 表1:行业资金流(对应 hizj/zjh,全字段原样落库)
CREATE TABLE IF NOT EXISTS industry_fund (
id INTEGER PRIMARY KEY AUTOINCREMENT,
t TEXT, mc TEXT, dm TEXT,
ac3 REAL, net3 REAL, ra3 REAL,
ac5 REAL, net5 REAL, ra5 REAL,
ac10 REAL, net10 REAL, ra10 REAL,
collected_at TEXT,
UNIQUE(t, dm)
);
— 表2:行业板块与龙头股(对应 hibk/zjhhy,全字段原样落库)
CREATE TABLE IF NOT EXISTS industry_board (
id INTEGER PRIMARY KEY AUTOINCREMENT,
t TEXT, mc TEXT, dm TEXT,
jj REAL, zdf REAL,
lrzj REAL, lczj REAL, jlr REAL, jlrl REAL,
lzgmc TEXT, lzgdm TEXT, lzgjlrl REAL,
collected_at TEXT,
UNIQUE(t, dm)
);
— 表3:个股净流入(对应 higg/jlr,全字段原样落库)
CREATE TABLE IF NOT EXISTS stock_netflow (
id INTEGER PRIMARY KEY AUTOINCREMENT,
t TEXT, mc TEXT, dm TEXT,
zxj REAL, zdf REAL, hsl REAL, cje REAL,
lczj REAL, lrzj REAL, jlr REAL, jlrl REAL,
collected_at TEXT,
UNIQUE(t, dm)
);
— 表4:自研轮动信号(除 t/dm/mc 与少量原值列外,其余全部自研)
CREATE TABLE IF NOT EXISTS rotation_signal (
id INTEGER PRIMARY KEY AUTOINCREMENT,
t TEXT, dm TEXT, mc TEXT,
zdf REAL, jlr REAL, jlrl REAL, — 原值,便于回看校验
ac10 REAL, net10 REAL, ra10 REAL, — 原值,背离判定输入
flow_rank INTEGER, flow_pct REAL, — 自研:资金排名与分位
price_rank INTEGER, price_pct REAL, — 自研:涨跌幅排名与分位
rotation_type TEXT, — 自研:量价齐升/逆势吸筹/价升量出/量价齐跌
lead_score REAL, — 自研:资金领先度 -100~100
resonance_cnt INTEGER, resonance_flag INTEGER, — 自研:多周期共振计数与方向
diverge_10 INTEGER, — 自研:净流率与涨跌幅背离
lzgdm TEXT, lzgmc TEXT, lzgjlrl REAL, — 原值:板块龙头股
leader_verified INTEGER, leader_stock_rank INTEGER, — 自研:龙头二次验证
leader_jlr REAL, leader_jlrl REAL, — 自研:龙头个股净流入
collected_at TEXT,
UNIQUE(t, dm)
);
说明:rotation_signal 中除 t / dm / mc / zdf / jlr / jlrl / ac10 / net10 / ra10 / lzgdm / lzgmc / lzgjlrl 外,其余列全部是本文自研衍生列,不是接口字段,官方接口不提供任何一列。前三张表是纯原始快照,保证任何时候都能从裸数据重算出指标。
五、完整可运行代码
import requests
import logging
import time
import sqlite3
from datetime import datetime
import pandas as pd
import numpy as np
# ========== 全局配置 ==========
LICENCE = "你的licence"
DB_PATH = "quant.db"
LOG_FILE = "quant_collect.log"
API_BASE = "http://api.biyingapi.com"
# ———- 常量:接口返回字段(严格对照官方文档,不增不减) ———-
FUND_FIELDS = ["t", "mc", "dm", "ac3", "net3", "ra3",
"ac5", "net5", "ra5", "ac10", "net10", "ra10"]
FUND_NUMERIC = ["ac3", "net3", "ra3", "ac5", "net5", "ra5",
"ac10", "net10", "ra10"]
BOARD_FIELDS = ["t", "mc", "dm", "jj", "zdf", "lrzj", "lczj", "jlr", "jlrl",
"lzgmc", "lzgdm", "lzgjlrl"]
BOARD_NUMERIC = ["jj", "zdf", "lrzj", "lczj", "jlr", "jlrl", "lzgjlrl"]
STOCK_FIELDS = ["t", "mc", "dm", "zxj", "zdf", "hsl", "cje",
"lczj", "lrzj", "jlr", "jlrl"]
STOCK_NUMERIC = ["zxj", "zdf", "hsl", "cje", "lczj", "lrzj", "jlr", "jlrl"]
# ———- 常量:自研指标参数 ———-
TOP_N = 5 # 资金流入榜 Top N
BOTTOM_N = 5 # 资金流出榜 Bottom N
PENETRATE_TOP = 10 # 板块穿透时的个股候选池大小
RESONANCE_PERIODS = 3 # 共振要求的同号周期数(net3/net5/net10)
FLOW_EPS = 1e-9 # 零值判定下限,防除零与符号误判
# ———- 落库列 ———-
FUND_COLUMNS = FUND_FIELDS + ["collected_at"]
BOARD_COLUMNS = BOARD_FIELDS + ["collected_at"]
STOCK_COLUMNS = STOCK_FIELDS + ["collected_at"]
SIGNAL_COLUMNS = [
"t", "dm", "mc",
"zdf", "jlr", "jlrl", "ac10", "net10", "ra10",
"flow_rank", "flow_pct", "price_rank", "price_pct",
"rotation_type", "lead_score",
"resonance_cnt", "resonance_flag", "diverge_10",
"lzgdm", "lzgmc", "lzgjlrl",
"leader_verified", "leader_stock_rank", "leader_jlr", "leader_jlrl",
"collected_at",
]
# ———- 日志初始化 ———-
logging.basicConfig(
filename=LOG_FILE,
level=logging.INFO,
format="%(asctime)s %(levelname)s %(message)s",
datefmt="%Y-%m-%d %H:%M:%S",
filemode="a"
)
logger = logging.getLogger(__name__)
# ———- 带重试 HTTP 请求(复用) ———-
def biying_api_get_retry(full_url, timeout=15, max_retry=3):
for attempt in range(1, max_retry + 1):
try:
resp = requests.get(full_url, timeout=timeout)
if resp.status_code == 200:
return resp.json()
logger.warning(f"HTTP状态码异常:{resp.status_code},第{attempt}次重试")
except Exception as e:
logger.warning(f"网络请求异常,第{attempt}次重试,错误信息:{str(e)}")
time.sleep(2)
logger.error("达到最大重试次数,接口请求失败")
return []
# ———- 代码补市场后缀 ———-
def with_exchange_suffix(code):
code = str(code).strip()
if code.startswith("60") or code.startswith("68") or code.startswith("9"):
return code + ".SH"
if code.startswith("00") or code.startswith("30") or code.startswith("2"):
return code + ".SZ"
if code.startswith("8") or code.startswith("4"):
return code + ".BJ"
return code + ".SH"
# ———- 自研工具:代码归一化(行业代码 / 板块代码 / 股票代码通用) ———-
def normalize_code(code):
"""去掉可能已有的交易所后缀,再用 with_exchange_suffix 补齐,保证两侧可比"""
if code is None or (not isinstance(code, str) and pd.isna(code)):
return None
raw = str(code).strip().split(".")[0]
if not raw:
return None
return with_exchange_suffix(raw)
# ———- 抓取:行业资金流 hizj/zjh ———-
def fetch_industry_fund():
url = f"{API_BASE}/hizj/zjh/{LICENCE}"
data = biying_api_get_retry(url)
if isinstance(data, dict):
data = data.get("data", [])
if not isinstance(data, list):
return []
return [row for row in data if isinstance(row, dict)]
# ———- 抓取:行业板块 hibk/zjhhy ———-
def fetch_industry_board():
url = f"{API_BASE}/hibk/zjhhy/{LICENCE}"
data = biying_api_get_retry(url)
if isinstance(data, dict):
data = data.get("data", [])
if not isinstance(data, list):
return []
return [row for row in data if isinstance(row, dict)]
# ———- 抓取:个股净流入 higg/jlr ———-
def fetch_stock_netflow():
url = f"{API_BASE}/higg/jlr/{LICENCE}"
data = biying_api_get_retry(url)
if isinstance(data, dict):
data = data.get("data", [])
if not isinstance(data, list):
return []
return [row for row in data if isinstance(row, dict)]
# ———- pandas 清洗链:行业资金流 ———-
def clean_fund_frame(rows):
df = pd.DataFrame(rows)
if df.empty:
return df
for col in FUND_FIELDS:
if col not in df.columns:
df[col] = np.nan
df = df.reindex(columns=FUND_FIELDS)
df["t"] = df["t"].astype(str)
df["dm"] = df["dm"].astype(str)
for col in FUND_NUMERIC:
df[col] = pd.to_numeric(df[col], errors="coerce")
df = df.dropna(subset=["t", "dm"])
df = df.drop_duplicates(subset=["t", "dm"], keep="last")
df = df.reset_index(drop=True)
return df
# ———- pandas 清洗链:行业板块 ———-
def clean_board_frame(rows):
df = pd.DataFrame(rows)
if df.empty:
return df
for col in BOARD_FIELDS:
if col not in df.columns:
df[col] = np.nan
df = df.reindex(columns=BOARD_FIELDS)
df["t"] = df["t"].astype(str)
df["dm"] = df["dm"].astype(str)
for col in BOARD_NUMERIC:
df[col] = pd.to_numeric(df[col], errors="coerce")
df = df.dropna(subset=["t", "dm"])
df = df.drop_duplicates(subset=["t", "dm"], keep="last")
df = df.reset_index(drop=True)
return df
# ———- pandas 清洗链:个股净流入 ———-
def clean_stock_frame(rows):
df = pd.DataFrame(rows)
if df.empty:
return df
for col in STOCK_FIELDS:
if col not in df.columns:
df[col] = np.nan
df = df.reindex(columns=STOCK_FIELDS)
df["t"] = df["t"].astype(str)
df["dm"] = df["dm"].astype(str)
for col in STOCK_NUMERIC:
df[col] = pd.to_numeric(df[col], errors="coerce")
df = df.dropna(subset=["t", "dm"])
df = df.drop_duplicates(subset=["t", "dm"], keep="last")
df = df.reset_index(drop=True)
return df
# ———- 自研指标 1:降序排名 + 分位 ———-
def add_rank_columns(df, value_col, rank_col, pct_col):
"""排名 1 = 最大值;分位 = (N – rank) / (N – 1) * 100,取值 0~100"""
if df.empty or value_col not in df.columns:
return df
valid = df[value_col].notna()
cnt = int(valid.sum())
df[rank_col] = df[value_col].rank(ascending=False, method="min")
if cnt > 1:
df[pct_col] = ((cnt – df[rank_col]) / (cnt – 1) * 100).round(2)
else:
df[pct_col] = np.nan
df.loc[~valid, [rank_col, pct_col]] = np.nan
return df
# ———- 自研指标 2:Top N / Bottom N 榜单 ———-
def top_bottom(df, value_col, name_col="mc", top_n=TOP_N, bottom_n=BOTTOM_N):
if df.empty or value_col not in df.columns:
return pd.DataFrame(), pd.DataFrame()
sub = df.dropna(subset=[value_col]).sort_values(value_col, ascending=False)
return sub.head(top_n), sub.tail(bottom_n).iloc[::–1]
# ———- 自研指标 3:轮动类型四象限打标 ———-
def label_rotation(zdf, jlr):
"""zdf 与 jlr 构造四象限;任一为 0 或缺失时不硬塞进象限"""
if pd.isna(zdf) or pd.isna(jlr):
return None
if abs(zdf) < FLOW_EPS or abs(jlr) < FLOW_EPS:
return "无可判定"
if zdf > 0 and jlr > 0:
return "量价齐升"
if zdf <= 0 and jlr > 0:
return "逆势吸筹"
if zdf > 0 and jlr <= 0:
return "价升量出"
return "量价齐跌"
# ———- 自研指标 4:多周期资金共振 ———-
def multi_period_resonance(net3, net5, net10):
"""net3/net5/net10 三者同号 -> 共振;返回 (同号个数, 方向标记)"""
vals = [net3, net5, net10]
if any(pd.isna(v) for v in vals):
return (None, None)
pos = sum(1 for v in vals if v > 0)
neg = sum(1 for v in vals if v < 0)
cnt = max(pos, neg)
if cnt == RESONANCE_PERIODS:
return (cnt, 1 if pos == RESONANCE_PERIODS else –1)
return (cnt, 0)
# ———- 自研指标 5:净流率与涨跌幅背离 ———-
def calc_divergence(ra10, ac10):
"""ra10 与 ac10 符号相反记 1,同号记 0,缺失记 None"""
if pd.isna(ra10) or pd.isna(ac10):
return None
if abs(ra10) < FLOW_EPS or abs(ac10) < FLOW_EPS:
return 0
return 1 if ra10 * ac10 < 0 else 0
# ———- 自研指标 6:龙头股二次验证 ———-
def match_leader(lzgdm, stock_df):
"""板块接口认定的龙头,个股净流入榜里认不认?
返回 (是否命中, 个股榜排名, jlr, jlrl);代码缺失返回全 None。
"""
if stock_df is None or stock_df.empty or "dm" not in stock_df.columns:
return (None, None, None, None)
target = normalize_code(lzgdm)
if target is None:
return (None, None, None, None)
codes = stock_df["dm"].map(normalize_code)
hit = stock_df[codes == target]
if hit.empty:
return (0, None, None, None)
row = hit.iloc[0]
return (1, row.get("stock_rank"), row.get("jlr"), row.get("jlrl"))
# ———- 自研指标 7:板块 -> 个股穿透 ———-
def penetrate_board(board_row, stock_df, top_n=PENETRATE_TOP):
"""以板块龙头股 lzgdm 为锚点穿透到个股,并输出全市场净流入候选池。
注意:higg/jlr 不返回个股所属行业,候选池个股无法声称归属该板块。
"""
info = {
"dm": board_row.get("dm"),
"mc": board_row.get("mc"),
"lzgmc": board_row.get("lzgmc"),
"lzgdm": board_row.get("lzgdm"),
"leader_verified": None,
"leader_stock_rank": None,
"leader_jlr": None,
"leader_jlrl": None,
"candidates": [],
}
verified, rank, jlr, jlrl = match_leader(board_row.get("lzgdm"), stock_df)
info["leader_verified"] = verified
info["leader_stock_rank"] = rank
info["leader_jlr"] = jlr
info["leader_jlrl"] = jlrl
if stock_df is None or stock_df.empty:
return info
pool = stock_df.dropna(subset=["jlr"]).sort_values("jlr", ascending=False)
cols = [c for c in ["dm", "mc", "zxj", "zdf", "jlr", "jlrl", "stock_rank"]
if c in pool.columns]
info["candidates"] = pool.head(top_n)[cols].to_dict("records")
return info
# ———- 汇总:生成自研轮动信号表 ———-
def build_rotation_frame(board_df, fund_df, stock_df):
if board_df.empty:
return pd.DataFrame(columns=SIGNAL_COLUMNS)
df = board_df.copy()
fund_part = ["t", "dm", "ac3", "net3", "ra3", "ac5", "net5", "ra5",
"ac10", "net10", "ra10"]
if not fund_df.empty:
df = df.merge(fund_df[fund_part], on=["t", "dm"], how="left")
else:
for col in fund_part:
if col not in df.columns:
df[col] = np.nan
df = add_rank_columns(df, "jlr", "flow_rank", "flow_pct")
df = add_rank_columns(df, "zdf", "price_rank", "price_pct")
df["rotation_type"] = [label_rotation(z, j) for z, j in zip(df["zdf"], df["jlr"])]
df["lead_score"] = (df["flow_pct"] – df["price_pct"]).round(2)
res = [multi_period_resonance(a, b, c)
for a, b, c in zip(df["net3"], df["net5"], df["net10"])]
df["resonance_cnt"] = [r[0] for r in res]
df["resonance_flag"] = [r[1] for r in res]
df["diverge_10"] = [calc_divergence(r, a) for r, a in zip(df["ra10"], df["ac10"])]
lead = [match_leader(d, stock_df) for d in df["lzgdm"]]
df["leader_verified"] = [x[0] for x in lead]
df["leader_stock_rank"] = [x[1] for x in lead]
df["leader_jlr"] = [x[2] for x in lead]
df["leader_jlrl"] = [x[3] for x in lead]
df["collected_at"] = datetime.now().strftime("%Y-%m-%d %H:%M:%S")
return df.reindex(columns=SIGNAL_COLUMNS)
# ———- 建表 ———-
def init_db():
conn = sqlite3.connect(DB_PATH)
try:
conn.execute("""
CREATE TABLE IF NOT EXISTS industry_fund (
id INTEGER PRIMARY KEY AUTOINCREMENT,
t TEXT, mc TEXT, dm TEXT,
ac3 REAL, net3 REAL, ra3 REAL,
ac5 REAL, net5 REAL, ra5 REAL,
ac10 REAL, net10 REAL, ra10 REAL,
collected_at TEXT,
UNIQUE(t, dm)
)
""")
conn.execute("""
CREATE TABLE IF NOT EXISTS industry_board (
id INTEGER PRIMARY KEY AUTOINCREMENT,
t TEXT, mc TEXT, dm TEXT,
jj REAL, zdf REAL,
lrzj REAL, lczj REAL, jlr REAL, jlrl REAL,
lzgmc TEXT, lzgdm TEXT, lzgjlrl REAL,
collected_at TEXT,
UNIQUE(t, dm)
)
""")
conn.execute("""
CREATE TABLE IF NOT EXISTS stock_netflow (
id INTEGER PRIMARY KEY AUTOINCREMENT,
t TEXT, mc TEXT, dm TEXT,
zxj REAL, zdf REAL, hsl REAL, cje REAL,
lczj REAL, lrzj REAL, jlr REAL, jlrl REAL,
collected_at TEXT,
UNIQUE(t, dm)
)
""")
conn.execute("""
CREATE TABLE IF NOT EXISTS rotation_signal (
id INTEGER PRIMARY KEY AUTOINCREMENT,
t TEXT, dm TEXT, mc TEXT,
zdf REAL, jlr REAL, jlrl REAL,
ac10 REAL, net10 REAL, ra10 REAL,
flow_rank INTEGER, flow_pct REAL,
price_rank INTEGER, price_pct REAL,
rotation_type TEXT,
lead_score REAL,
resonance_cnt INTEGER, resonance_flag INTEGER,
diverge_10 INTEGER,
lzgdm TEXT, lzgmc TEXT, lzgjlrl REAL,
leader_verified INTEGER, leader_stock_rank INTEGER,
leader_jlr REAL, leader_jlrl REAL,
collected_at TEXT,
UNIQUE(t, dm)
)
""")
conn.commit()
finally:
conn.close()
# ———- 落库(INSERT OR REPLACE 幂等写入) ———-
def save_frame(table, columns, df):
if df is None or df.empty:
return
conn = sqlite3.connect(DB_PATH)
try:
data = df.copy()
for col in columns:
if col not in data.columns:
data[col] = np.nan
data = data[columns]
data.to_sql("tmp_frame", conn, if_exists="replace", index=False)
cols = ", ".join(columns)
conn.execute(
f"INSERT OR REPLACE INTO {table} ({cols}) SELECT {cols} FROM tmp_frame"
)
conn.execute("DROP TABLE IF EXISTS tmp_frame")
conn.commit()
logger.info(f"落库 {table} 共 {len(data)} 条")
except Exception as e:
logger.error(f"落库 {table} 失败:{str(e)}")
finally:
conn.close()
# ———- 打印辅助:榜单 ———-
def print_board(title, frame, value_col):
print(f"\\n— {title} —")
if frame.empty:
print("(无有效数据)")
return
for _, r in frame.iterrows():
print(f"{r.get('mc')} | 代码 {r.get('dm')} | {value_col} {r.get(value_col)} "
f"| 涨跌幅 {r.get('zdf')} | 净流入 {r.get('jlr')} | 净流入率 {r.get('jlrl')}")
# ———- 演示入口 ———-
def demo():
init_db()
fund_df = clean_fund_frame(fetch_industry_fund())
board_df = clean_board_frame(fetch_industry_board())
stock_df = clean_stock_frame(fetch_stock_netflow())
print(f"行业资金流 {len(fund_df)} 条 | 行业板块 {len(board_df)} 条 | 个股净流入 {len(stock_df)} 条")
# 个股榜:先排名,供龙头股二次验证使用
if not stock_df.empty:
stock_df = add_rank_columns(stock_df, "jlr", "stock_rank", "stock_pct")
else:
stock_df["stock_rank"] = np.nan
stock_df["stock_pct"] = np.nan
if fund_df.empty and board_df.empty and stock_df.empty:
logger.error("三个接口均无有效数据,请检查 licence 与网络")
print("三个接口均无有效数据,请检查 licence 与网络")
return
# ———- 1. 原始三表落库 ———-
now = datetime.now().strftime("%Y-%m-%d %H:%M:%S")
for frame in (fund_df, board_df, stock_df):
if not frame.empty:
frame["collected_at"] = now
save_frame("industry_fund", FUND_COLUMNS, fund_df)
save_frame("industry_board", BOARD_COLUMNS, board_df)
save_frame("stock_netflow", STOCK_COLUMNS, stock_df)
# ———- 2. 自研信号 ———-
signal = build_rotation_frame(board_df, fund_df, stock_df)
save_frame("rotation_signal", SIGNAL_COLUMNS, signal)
# ———- 3. 资金流向排名 ———-
top_df, bottom_df = top_bottom(signal, "jlr", top_n=TOP_N, bottom_n=BOTTOM_N)
print_board(f"行业板块净流入 Top {TOP_N}", top_df, "jlr")
print_board(f"行业板块净流出 Bottom {BOTTOM_N}", bottom_df, "jlr")
if not fund_df.empty:
f_top, f_bottom = top_bottom(fund_df, "net10", top_n=TOP_N, bottom_n=BOTTOM_N)
print(f"\\n— 行业资金流近10日净流额 Top {TOP_N} —")
for _, r in f_top.iterrows():
print(f"{r.get('mc')} | 代码 {r.get('dm')} | net10 {r.get('net10')} "
f"| ra10 {r.get('ra10')} | ac10 {r.get('ac10')}")
print(f"\\n— 行业资金流近10日净流额 Bottom {BOTTOM_N} —")
for _, r in f_bottom.iterrows():
print(f"{r.get('mc')} | 代码 {r.get('dm')} | net10 {r.get('net10')} "
f"| ra10 {r.get('ra10')} | ac10 {r.get('ac10')}")
# ———- 4. 轮动类型分布 ———-
if not signal.empty:
print("\\n— 轮动类型分布 —")
for name, cnt in signal["rotation_type"].value_counts().items():
print(f"{name}:{int(cnt)} 个板块")
strong = signal[signal["rotation_type"].isin(["量价齐升", "逆势吸筹"])]
strong = strong.dropna(subset=["lead_score"]).sort_values("lead_score", ascending=False)
print(f"\\n— 强势板块资金领先度 Top {TOP_N}(量价齐升 / 逆势吸筹)—")
for _, r in strong.head(TOP_N).iterrows():
print(f"{r.get('mc')} | 类型 {r.get('rotation_type')} "
f"| 资金分位 {r.get('flow_pct')} | 价格分位 {r.get('price_pct')} "
f"| 领先度 {r.get('lead_score')}")
# ———- 5. 共振与背离 ———-
print("\\n— 多周期资金共振与背离统计 —")
print(f"共振流入(resonance_flag=1):{int((signal['resonance_flag'] == 1).sum())} 个板块")
print(f"共振流出(resonance_flag=-1):{int((signal['resonance_flag'] == –1).sum())} 个板块")
print(f"未共振(resonance_flag=0):{int((signal['resonance_flag'] == 0).sum())} 个板块")
print(f"近10日净流率与涨跌幅背离:{int(signal['diverge_10'].fillna(0).sum())} 个板块")
res_hit = signal[(signal["resonance_flag"] == 1) & (signal["diverge_10"] == 1)]
print(f"既共振流入又背离(资金持续进、价格没跟上):{len(res_hit)} 个板块")
for _, r in res_hit.head(TOP_N).iterrows():
print(f" {r.get('mc')} | net10 {r.get('net10')} | ra10 {r.get('ra10')} "
f"| ac10 {r.get('ac10')} | 领先度 {r.get('lead_score')}")
# ———- 6. 龙头股捕捉与板块穿透 ———-
if not signal.empty:
verified = signal[signal["leader_verified"] == 1]
print(f"\\n— 龙头股二次验证:{len(verified)} / {len(signal)} 个板块龙头出现在个股净流入榜 —")
if not verified.empty:
v = verified.dropna(subset=["leader_stock_rank"]).sort_values("leader_stock_rank")
for _, r in v.head(TOP_N).iterrows():
print(f"{r.get('mc')} -> {r.get('lzgmc')}({r.get('lzgdm')})"
f" | 个股榜排名 {r.get('leader_stock_rank')} "
f"| 个股净流入 {r.get('leader_jlr')} | 龙头净流入率 {r.get('lzgjlrl')}")
print(f"\\n— 板块 -> 个股穿透(取净流入 Top {TOP_N} 板块)—")
for _, r in top_df.iterrows():
info = penetrate_board(r, stock_df, top_n=PENETRATE_TOP)
print(f"\\n【{info['mc']}】龙头 {info['lzgmc']}({info['lzgdm']})"
f" 验证 {info['leader_verified']} 个股榜排名 {info['leader_stock_rank']}")
for c in info["candidates"][:3]:
print(f" 候选 {c.get('mc')}({c.get('dm')})| 最新价 {c.get('zxj')} "
f"| 涨跌幅 {c.get('zdf')} | 净流入 {c.get('jlr')}")
print(f" (候选池来自全市场个股净流入榜前 {PENETRATE_TOP} 只,"
f"接口不返回个股所属行业,不做行业归属断言)")
print("\\n【仅为数据演示,不构成投资建议】已采集并落库板块轮动数据")
if __name__ == "__main__":
demo()
代码运行后,会依次打印:三个接口各自返回的数据条数;行业板块净流入 Top / Bottom 榜单(含涨跌幅、净流入、净流入率);行业资金流近 10 日净流额 Top / Bottom(含 ra10 与 ac10);轮动四象限的类型分布;强势板块的资金分位、价格分位与资金领先度;多周期共振与背离的板块计数,以及「既共振流入又背离」的板块清单;龙头股二次验证的命中率与命中板块的个股榜排名;最后是板块 → 个股穿透结果。所有数值均由实际接口返回计算得出,本文不预设任何具体结果。
六、业务关键点
七、拓展练习
八、下篇预告
下一篇《历史涨跌停价:涨跌停序列与止损线测算》将用 hsstock/stopprice/history 拿到每只股票历史每日的涨停价 h 与跌停价 l,并自研「涨跌停价序列分析」「跌停价偏离度止损线」「波动率自适应止损」等风控指标——这是本系列在「进攻」之外补上的最后一块「防守」拼图。
九、免责申明
免责申明:文中所有数据处理逻辑仅为编程演示,仅为数据演示,不构成投资建议。市场有风险,投资需谨慎。
网硕互联帮助中心





评论前必须登录!
注册