云计算百科
云计算领域专业知识百科平台

量化数据进阶:多维实战篇(第 10 篇):板块轮动:行业资金流向与龙头股捕捉

量化数据进阶:多维实战篇(第 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 手写(均为自研因子,非接口能力):

  • 资金流向排名 flow_rank(自研):按净流入 jlr 降序排名,rank(ascending=False, method="min"),排名 1 表示净流入最大。缺失值不参与排名,排名结果记为 NaN。
  • 资金分位 flow_pct(自研):(N – flow_rank) / (N – 1) * 100,取值 0~100,把排名归一化成可跨日期比较的分位。样本数 N <= 1 时无法计算,返回 NaN。
  • Top N / Bottom N 榜单(自研):排序后 head(TOP_N) 与 tail(BOTTOM_N),Bottom 再 iloc[::-1] 反转,保证输出是从「流出最多」向下排列,而不是倒序的尾部。
  • 涨跌幅排名 price_rank 与价格分位 price_pct(自研):与资金排名同构,但按 zdf 降序。有了「资金分位」与「价格分位」两套坐标,才能比较钱的位置与价的位置是否匹配。
  • 轮动类型 rotation_type(自研,四象限打标):这是「轮动强度」的核心。以 zdf(价格方向)与 jlr(资金方向)构造四象限:zdf > 0 且 jlr > 0 → 「量价齐升」;zdf <= 0 且 jlr > 0 → 「逆势吸筹」;zdf > 0 且 jlr <= 0 → 「价升量出」;zdf <= 0 且 jlr <= 0 → 「量价齐跌」。两者任一为 0 或缺失时记「无可判定」,不硬塞进象限。
  • 资金领先度 lead_score(自研):flow_pct – price_pct,取值 -100~100。显著为正 = 资金分位远高于价格分位 = 钱已经进去了、价格还没跟上,是轮动的早期特征;显著为负 = 价格已在高位但资金分位靠后,是轮动的末期特征。它回答的是「资金比价格领先多少」,与四象限标签互补。
  • 多周期共振计数 resonance_cnt(自研):取同一行业的 net3、net5、net10,统计其中符号相同的最大个数(正数个数与负数个数取较大者),取值 1~3。三个周期中任一缺失时返回 None,不做部分判定。
  • 多周期共振标记 resonance_flag(自研):resonance_cnt == 3 时,若三个净流额全为正记 1(资金共振流入),全为负记 -1(资金共振流出);未达三条同号记 0;数据缺失记 None。短、中、长三个周期同向,说明资金行为不是单日脉冲,而是持续性的。
  • 净流率与涨跌幅背离 diverge_10(自研):ra10 * ac10 < 0 时记 1,表示近 10 日「资金净流率方向」与「涨跌幅方向」相反——资金在流入但价格在跌,或资金在流出但价格在涨,两者必有一个不可持续。同号记 0;ra10 或 ac10 缺失记 None;任一恰为 0 时同样记 0,避免把零值误判为「相反」。
  • 代码归一化 normalize_code()(自研工具):龙头股代码 lzgdm 与个股榜 dm 可能一个带交易所后缀、一个不带。统一 split(".")[0] 去掉后缀,再用系列工具函数 with_exchange_suffix() 按代码前缀补回,两侧都归一化后再比对,避免「600519」与「600519.SH」匹配不上。
  • 龙头股二次验证 leader_verified(自研):板块接口说它是龙头,个股净流入榜里能不能找到它?找得到记 1,找不到记 0,龙头股代码缺失记 None。接口给的龙头是一回事,真金白银的净流入榜认不认是另一回事,两者都成立才算验证通过。
  • 龙头股榜内排名 leader_stock_rank(自研):龙头股在全市场个股净流入榜中的排名。排名越靠前,说明资金在个股层面与板块层面是同一批钱,龙头地位越扎实。
  • 板块 → 个股穿透 penetrate_board()(自研):先按 flow_rank 与 rotation_type 定强势板块,再以该板块的 lzgdm 为锚点在个股榜中定位,同时输出全市场净流入榜前 PENETRATE_TOP 只作为资金候选池。必须诚实说明:higg/jlr 不返回个股所属行业字段,因此候选池个股无法声称归属于某个板块,穿透只能做到「龙头锚点 + 全市场资金候选池」这一步,个股的行业归属需要自建映射表才能闭环。
  • null 值防御贯穿全程:以上每一项在计算前统一 pd.isna 判空,分母(样本数 N-1、净流额零值)单独判零,绝不让 None 参与比较与四则运算;「数据不足无法判定」一律返回 None,与「已计算但无信号」的 0 严格区分。
  • 四、数据表设计

    — 表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);轮动四象限的类型分布;强势板块的资金分位、价格分位与资金领先度;多周期共振与背离的板块计数,以及「既共振流入又背离」的板块清单;龙头股二次验证的命中率与命中板块的个股榜排名;最后是板块 → 个股穿透结果。所有数值均由实际接口返回计算得出,本文不预设任何具体结果。

    六、业务关键点

  • 三个接口的 dm 是三种不同的东西:hizj/zjh 的 dm 是行业代码,hibk/zjhhy 的 dm 是板块代码,higg/jlr 的 dm 是股票代码。本文在 build_rotation_frame() 里只敢用 (t, dm) 去关联前两张表(同为证监会行业分类口径),绝不把行业代码与股票代码放在一起 join。跨表关联前务必确认两侧代码同域,这是这类多接口组合最容易出静默错误的地方——关联不上时 merge 不会报错,只会给你一片 NaN。
  • lrzj – lczj = jlr 是一个免费的数据自洽校验:行业板块接口把流入、流出、净额三个数都给了。落库后可以跑一条 SELECT COUNT(*) FROM industry_board WHERE ABS(lrzj – lczj – jlr) > 阈值 检查接口是否有异常行。同理 higg/jlr 也有 lrzj / lczj / jlr 三件套,同样可校验。这类「接口自己给自己出的校验题」成本极低,建议每次接新接口都先找一遍。
  • 「轮动强度」不能用单一数字概括,本文拆成三个互补视角:四象限 rotation_type 回答「价与钱是否同向」;lead_score 回答「资金比价格领先多少」;resonance_flag 回答「这种同向持续了几个周期」。三者合起来才构成可解释的判断链——只看 jlr 排名,你会把「连续三天流入但今天已开始下跌」和「今天刚开始流入」混为一谈。
  • 零值必须单独判,不能直接看符号:zdf 为 0(平盘)或 jlr 为 0(资金持平)时,zdf > 0 与 zdf <= 0 的分支会把它误划进「逆势吸筹」或「量价齐跌」。本文用 FLOW_EPS 拦在前面,零值一律记「无可判定」。同理背离判定里 abs(ra10) < FLOW_EPS or abs(ac10) < FLOW_EPS 直接记 0——零既不是正也不是负,拿它去乘另一个数判符号相反,得到的永远是错的结论。
  • 共振要求「三条同号」是有意为之:net3、net5、net10 三取三同号,条件相当严格,命中板块数量通常很少,这是刻意的——共振的意义在于排除单日脉冲。如果你想放宽,把 RESONANCE_PERIODS 改成 2 即可,代码里 resonance_cnt 已经把「几條同号」记下来,不需要改判定逻辑。但要注意:resonance_cnt 取的是「正数个数与负数个数的较大者」,若三个数里有一个恰好为 0,它既不计入 pos 也不计入 neg,此时 cnt 最大只能是 2,不会误判为共振。
  • 龙头股代码必须归一化后再比对:lzgdm 与 higg/jlr 的 dm 谁带后缀、谁不带,不同批次可能不一致。本文 normalize_code() 先 split(".")[0] 去掉后缀,再用系列工具函数 with_exchange_suffix() 按前缀补回,两侧都走同一条路径,杜绝「600519」与「600519.SH」匹配不上的假阴性。
  • 龙头股「验证不通过」不等于「不是龙头」:higg/jlr 返回的是净流入榜,本质是一个截断的榜单,不是全市场股票清单。龙头股没出现在榜上,可能是资金没进它,也可能只是它排在榜外。所以 leader_verified = 0 只能解读为「个股净流入榜未覆盖」,不能解读为「龙头失格」。这也是本文把它命名为「验证」而不是「证伪」的原因。
  • 板块 → 个股穿透在本接口体系下无法完全闭环:任务描述里说「在该板块的个股净流入榜中取靠前标的」,但 higg/jlr 的字段表里没有任何行业归属字段——它返回的是一份跨行业的个股净流入榜。本文的穿透因此退化为「以板块龙头股为锚点 + 全市场净流入候选池」两步,并在输出里明确标注「不做行业归属断言」。要实现真正的板块成分股穿透,必须自己维护一份「股票代码 → 行业代码」映射表(可从板块成分股接口或本地行业分类表构建),这是本篇留下的扩展口子,不要假装接口给了你没给的东西。
  • 三个接口的 t 未必严格一致:资金流、板块、个股三个接口各自返回自己的更新时间,落库后可能出现 t 相差一两分钟的情况。本文用 (t, dm) 关联前两张表,若两侧 t 不完全相同,merge 会大面积失配。生产环境建议先按「日期 + 最近时点」做一层时间对齐(例如统一取当日最新的那个 t 作为关联键),或单独维护一张「交易时点」维表。
  • rotation_signal 是可重算的中间产物,不是原始数据:前三张表 industry_fund / industry_board / stock_netflow 是纯原始快照,任何时候都能从裸数据重算出全部指标。若哪天你改了 TOP_N、把 RESONANCE_PERIODS 从 3 调到 2、或者换了四象限的划分阈值,直接从原始表重跑即可,不用重新采集。自研因子的全部代价与自由,都落在「原始层与派生层必须严格分离」这一条上。
  • 七、拓展练习

  • 把 TOP_N 改成 10、RESONANCE_PERIODS 改成 2 分别重跑,统计 rotation_signal 表中「既共振流入又背离」的板块数量变化,体会阈值松紧对信号数量的影响——信号变多不等于变好,关键是看它们在下一时点是否仍然成立。
  • 用 rotation_signal 按 dm 做自关联:取当日 rotation_type = '逆势吸筹' 且 lead_score > 30 的板块,回看它们在后续若干采集时点中 rotation_type 是否切换为「量价齐升」,检验「资金领先度」是否真的领先价格。
  • 自建一张 stock_industry 映射表(字段自拟,至少含股票代码与行业代码),把 higg/jlr 的个股按行业分组后再做净流入排名,与本文「全市场候选池」的做法对比,观察龙头股在「板块内排名」与「全市场排名」两个口径下的差异。
  • 八、下篇预告

    下一篇《历史涨跌停价:涨跌停序列与止损线测算》将用 hsstock/stopprice/history 拿到每只股票历史每日的涨停价 h 与跌停价 l,并自研「涨跌停价序列分析」「跌停价偏离度止损线」「波动率自适应止损」等风控指标——这是本系列在「进攻」之外补上的最后一块「防守」拼图。

    九、免责申明

    免责申明:文中所有数据处理逻辑仅为编程演示,仅为数据演示,不构成投资建议。市场有风险,投资需谨慎。

    赞(0)
    未经允许不得转载:网硕互联帮助中心 » 量化数据进阶:多维实战篇(第 10 篇):板块轮动:行业资金流向与龙头股捕捉
    分享到: 更多 (0)

    评论 抢沙发

    评论前必须登录!