流水线运行是否正随时间变得更加集中?

文章来源声明: 原文作者:Discipline1029; 来源站点:掘金; 原文链接:https://juejin.cn/post/7691132455894499370; 本文基于上述来源整理/加工,觅优补充点评,仅供技术学习交流。版权归原作者所有。
觅优短评

该分析把集中度趋势转化为可检验指标,并强调FDR校正与效应量。适合平台治理、编排组件选型及生态集中度监测,但结论仅适用于本数据集。

流水线运行是否正随时间变得更加集中? ------------------

数据与完整代码来源:

<span>https:</span>/<span>/mbd.pub/o</span><span>/bread/</span><span>YZaVmJhp</span>ZQ

数据集简介

围绕“流水线运行是否正随时间变得更加集中”这一问题,这里使用了一份流水线运行记录数据。数据以逐条运行的形式组织,每条记录对应一次流水线执行,涵盖运行发生的时间以及运行所属的主体(如项目或执行者)等维度信息,可用于刻画流水线活动在不同主体之间的分布状况。

从分析角度看,这类数据的价值在于把“集中度”变成可度量的对象:按时间切片统计各主体的运行份额,再借助基尼系数、赫芬达尔指数或头部占比等指标,就能观察流水线运行是否正在向少数主体聚集,以及这种趋势随时间如何演变。数据覆盖的时间跨度越长,越有助于区分短期波动与长期结构性变化。

适用场景方面,它既适合做集中度与不平等程度的时序分析,也适合通过可视化展示头部主体份额的变化轨迹,还可以进一步探讨集中度升降背后的驱动因素,例如生态扩张、平台策略调整或自动化程度的提升。

后文将基于这份数据,从指标构建到图表呈现,完整演示如何回答标题中的问题。

数据集

Data Engineering & Modern AI Pipelines 2026

研究问题

流水线运行是否正随时间日益集中于越来越少的编排组件集合?

核心要点

在本分析中,我们发现以下方面不存在统计上显著的集中度漂移:

  • 编排平台,以及
  • 执行引擎。

也就是说,分布总体上保持稳定,而非发生结构性漂移。

方法概述

  1. 使用 pipeline_execution_telemetry 作为主要运营表。
  2. 构建按小时分桶、实体级运行计数的面板数据。
  3. 按小时计算集中度指标:
    • Gini(值越高 = 越集中)
    • HHI(值越高 = 越集中)
    • Normalized Entropy(值越低 = 越集中)
  4. 使用以下方法检验时间趋势:
    • 针对 Spearman 趋势的块置换检验
    • Mann-Kendall 趋势检验
    • Theil-Sen 稳健斜率
    • 跨指标的 FDR 校正
  5. 评估在注入噪声下的稳健性。
<span>from</span> dataclasses <span>import</span> dataclass, asdict, field
<span>from</span> pathlib <span>import</span> Path
<span>from</span> datetime <span>import</span> datetime
<span>from</span> typing <span>import</span> <span>Dict</span>, <span>Any</span>, <span>List</span>, <span>Optional</span>
<span>import</span> hashlib
<span>import</span> json

<span>import</span> numpy <span>as</span> np
<span>import</span> pandas <span>as</span> pd
<span>import</span> matplotlib.pyplot <span>as</span> plt
<span>from</span> scipy <span>import</span> stats

plt.rcParams.update({
    <span>"font.family"</span>: <span>"sans-serif"</span>,
    <span>"figure.facecolor"</span>: <span>"#0d1117"</span>,
    <span>"axes.facecolor"</span>: <span>"#161b22"</span>,
    <span>"axes.edgecolor"</span>: <span>"#30363d"</span>,
    <span>"axes.labelcolor"</span>: <span>"#c9d1d9"</span>,
    <span>"text.color"</span>: <span>"#c9d1d9"</span>,
    <span>"xtick.color"</span>: <span>"#8b949e"</span>,
    <span>"ytick.color"</span>: <span>"#8b949e"</span>,
    <span>"grid.color"</span>: <span>"#21262d"</span>,
    <span>"grid.alpha"</span>: <span>0.6</span>,
})

<span>@dataclass</span>
<span>class</span> <span>Config</span>:
    random_seed: <span>int</span> = <span>42</span>

    <span># Data discovery</span>
    dataset_dir: <span>Optional</span>[<span>str</span>] = <span>None</span>
    expected_file: <span>str</span> = <span>"pipeline_execution_telemetry.parquet"</span>

    <span># Core schema</span>
    time_col: <span>str</span> = <span>"execution_timestamp"</span>
    freq: <span>str</span> = <span>"1h"</span>

    <span># Filtering</span>
    min_total_per_bucket: <span>int</span> = <span>5</span>

    <span># Scenario entities</span>
    entity_candidates: <span>List</span>[<span>str</span>] = field(default_factory=<span>lambda</span>: [<span>"orchestrator"</span>, <span>"execution_engine"</span>])

    <span># Statistical controls</span>
    n_permutations: <span>int</span> = <span>3000</span>
    alpha: <span>float</span> = <span>0.05</span>
    noise_sigmas: <span>List</span>[<span>float</span>] = field(default_factory=<span>lambda</span>: [<span>0.0</span>, <span>0.02</span>, <span>0.05</span>, <span>0.1</span>, <span>0.2</span>])
    robustness_repeats: <span>int</span> = <span>100</span>

    <span># Output</span>
    output_dir: <span>str</span> = <span>"kaggle_outputs_pipeline_concentration_stability"</span>

cfg = Config()
rng = np.random.default_rng(cfg.random_seed)
Path(cfg.output_dir).mkdir(parents=<span>True</span>, exist_ok=<span>True</span>)
cfg

Config(random_seed=42, dataset_dir=None, expected_file='pipeline_execution_telemetry.parquet', time_col='execution_timestamp', freq='1h', min_total_per_bucket=5, entity_candidates=['orchestrator', 'execution_engine'], n_permutations=3000, alpha=0.05, noise_sigmas=[0.0, 0.02, 0.05, 0.1, 0.2], robustness_repeats=100, output_dir='kaggle_outputs_pipeline_concentration_stability')

<span>def</span> <span>discover_file</span>(<span>cfg: Config</span>) -> Path:
    roots = []
    <span>if</span> cfg.dataset_dir:
        roots.append(Path(cfg.dataset_dir))
    roots.extend([Path(<span>'.'</span>), Path(<span>'/kaggle/input'</span>)])

    <span>for</span> root <span>in</span> roots:
        <span>if</span> <span>not</span> root.exists():
            <span>continue</span>

        p = root / cfg.expected_file
        <span>if</span> p.exists():
            <span>return</span> p

        hits = <span>list</span>(root.rglob(cfg.expected_file))
        <span>if</span> hits:
            <span>return</span> hits[<span>0</span>]

        csv_name = cfg.expected_file.replace(<span>'.parquet'</span>, <span>'.csv'</span>)
        p_csv = root / csv_name
        <span>if</span> p_csv.exists():
            <span>return</span> p_csv
        hits_csv = <span>list</span>(root.rglob(csv_name))
        <span>if</span> hits_csv:
            <span>return</span> hits_csv[<span>0</span>]

    <span>raise</span> FileNotFoundError(<span>"Could not locate pipeline_execution_telemetry.parquet/.csv"</span>)

<span>def</span> <span>load_data</span>(<span>cfg: Config</span>) -> pd.DataFrame:
    file_path = discover_file(cfg)
    <span>print</span>(<span>"Using file:"</span>, file_path)

    <span>if</span> file_path.suffix.lower() == <span>".parquet"</span>:
        df = pd.read_parquet(file_path)
    <span>else</span>:
        df = pd.read_csv(file_path, low_memory=<span>False</span>)

    required = [cfg.time_col, <span>"orchestrator"</span>, <span>"execution_engine"</span>]
    miss = [c <span>for</span> c <span>in</span> required <span>if</span> c <span>not</span> <span>in</span> df.columns]
    <span>if</span> miss:
        <span>raise</span> ValueError(<span>f"Missing required columns: <span>{miss}</span>"</span>)

    df[cfg.time_col] = pd.to_datetime(df[cfg.time_col], errors=<span>"coerce"</span>, utc=<span>True</span>)
    df = df.dropna(subset=[cfg.time_col]).copy()
    <span>return</span> df

<span>def</span> <span>audit_data</span>(<span>df: pd.DataFrame, cfg: Config</span>) -> <span>Dict</span>[<span>str</span>, <span>Any</span>]:
    <span>return</span> {
        <span>"shape"</span>: [<span>int</span>(df.shape[<span>0</span>]), <span>int</span>(df.shape[<span>1</span>])],
        <span>"data_hash"</span>: hashlib.sha256(pd.util.hash_pandas_object(df, index=<span>True</span>).values).hexdigest()[:<span>16</span>],
        <span>"min_timestamp"</span>: <span>str</span>(df[cfg.time_col].<span>min</span>()),
        <span>"max_timestamp"</span>: <span>str</span>(df[cfg.time_col].<span>max</span>()),
        <span>"orchestrator_nunique"</span>: <span>int</span>(df[<span>"orchestrator"</span>].nunique(dropna=<span>True</span>)),
        <span>"execution_engine_nunique"</span>: <span>int</span>(df[<span>"execution_engine"</span>].nunique(dropna=<span>True</span>)),
    }

raw_df = load_data(cfg)
audit = audit_data(raw_df, cfg)
audit

Using file: /kaggle/input/datasets/dianatofficial/data-engineering-ai-pipelines/pipeline_execution_telemetry.parquet

{'shape': [150000, 24],
 'data_hash': '11fc4fa4d9ca3c1b',
 'min_timestamp': '2025-01-01 00:00:58+00:00',
 'max_timestamp': '2026-03-26 23:59:11+00:00',
 'orchestrator_nunique': 7,
 'execution_engine_nunique': 7}

<span>def</span> <span>gini</span>(<span>x: np.ndarray</span>) -> <span>float</span>:
    x = np.asarray(x, dtype=<span>float</span>)
    x = x[x >= <span>0</span>]
    <span>if</span> x.size == <span>0</span> <span>or</span> x.<span>sum</span>() == <span>0</span>:
        <span>return</span> <span>0.0</span>
    x = np.sort(x)
    n = x.size
    idx = np.arange(<span>1</span>, n + <span>1</span>)
    <span>return</span> <span>float</span>(((<span>2</span> * idx - n - <span>1</span>) * x).<span>sum</span>() / (n * x.<span>sum</span>()))

<span>def</span> <span>entropy_norm</span>(<span>p: np.ndarray</span>) -> <span>float</span>:
    p = np.asarray(p, dtype=<span>float</span>)
    p = p[p > <span>0</span>]
    n = <span>max</span>(<span>2</span>, <span>len</span>(p))
    <span>return</span> <span>float</span>(-(p * np.log(p)).<span>sum</span>() / np.log(n))

<span>def</span> <span>compute_metrics</span>(<span>shares: pd.DataFrame</span>) -> pd.DataFrame:
    rows = []
    <span>for</span> t, row <span>in</span> shares.iterrows():
        v = row.values.astype(<span>float</span>)
        rows.append({
            <span>"time_bucket"</span>: t,
            <span>"gini"</span>: gini(v),
            <span>"hhi"</span>: <span>float</span>(np.square(v).<span>sum</span>()),
            <span>"normalized_entropy"</span>: entropy_norm(v),
        })
    <span>return</span> pd.DataFrame(rows).set_index(<span>"time_bucket"</span>)

<span>def</span> <span>mann_kendall</span>(<span>y: np.ndarray</span>) -> <span>Dict</span>[<span>str</span>, <span>float</span>]:
    y = np.asarray(y, dtype=<span>float</span>)
    n = <span>len</span>(y)
    s = <span>0</span>
    <span>for</span> i <span>in</span> <span>range</span>(n - <span>1</span>):
        s += np.sign(y[i + <span>1</span>:] - y[i]).<span>sum</span>()

    _, counts = np.unique(y, return_counts=<span>True</span>)
    tie_term = np.<span>sum</span>(counts * (counts - <span>1</span>) * (<span>2</span> * counts + <span>5</span>))
    var_s = (n * (n - <span>1</span>) * (<span>2</span> * n + <span>5</span>) - tie_term) / <span>18.0</span>

    <span>if</span> s > <span>0</span>:
        z = (s - <span>1</span>) / np.sqrt(var_s)
    <span>elif</span> s < <span>0</span>:
        z = (s + <span>1</span>) / np.sqrt(var_s)
    <span>else</span>:
        z = <span>0.0</span>

    p = <span>2</span> * (<span>1</span> - stats.norm.cdf(<span>abs</span>(z)))
    <span>return</span> {<span>"S"</span>: <span>float</span>(s), <span>"z"</span>: <span>float</span>(z), <span>"p_value"</span>: <span>float</span>(p)}

<span>def</span> <span>block_permutation</span>(<span>series: np.ndarray, n_perms: <span>int</span>, block_size: <span>int</span>, seed: <span>int</span></span>):
    rng_local = np.random.default_rng(seed)
    x = np.asarray(series, dtype=<span>float</span>)
    t = np.arange(<span>len</span>(x), dtype=<span>float</span>)
    rho_obs, _ = stats.spearmanr(t, x)
    rho_obs = <span>float</span>(rho_obs)

    blocks = [x[i:i + block_size] <span>for</span> i <span>in</span> <span>range</span>(<span>0</span>, <span>len</span>(x), block_size)]
    null = np.empty(n_perms, dtype=<span>float</span>)

    <span>for</span> i <span>in</span> <span>range</span>(n_perms):
        perm = rng_local.permutation(<span>len</span>(blocks))
        xp = np.concatenate([blocks[j] <span>for</span> j <span>in</span> perm])[:<span>len</span>(x)]
        rho, _ = stats.spearmanr(t, xp)
        null[i] = <span>float</span>(rho)

    k = <span>int</span>(np.<span>sum</span>(np.<span>abs</span>(null) >= <span>abs</span>(rho_obs)))
    p = (k + <span>1</span>) / (n_perms + <span>1</span>)
    <span>return</span> rho_obs, <span>float</span>(p), null

<span>def</span> <span>fdr_bh</span>(<span>pvals: <span>Dict</span>[<span>str</span>, <span>float</span>]</span>) -> <span>Dict</span>[<span>str</span>, <span>float</span>]:
    ordered = <span>sorted</span>(pvals.items(), key=<span>lambda</span> kv: kv[<span>1</span>])
    m = <span>len</span>(ordered)
    out = {}
    prev = <span>1.0</span>
    <span>for</span> rank, (name, p) <span>in</span> <span>enumerate</span>(<span>reversed</span>(ordered), start=<span>1</span>):
        i = m - rank + <span>1</span>
        q = <span>min</span>(prev, p * m / i)
        out[name] = q
        prev = q
    <span>return</span> {k: <span>float</span>(out[k]) <span>for</span> k, _ <span>in</span> ordered}

<span>def</span> <span>robustness_mc</span>(<span>shares: pd.DataFrame, cfg: Config</span>) -> pd.DataFrame:
    base = compute_metrics(shares)
    b = {
        <span>"gini"</span>: base[<span>"gini"</span>].values,
        <span>"hhi"</span>: base[<span>"hhi"</span>].values,
        <span>"entropy"</span>: base[<span>"normalized_entropy"</span>].values,
    }

    rows = []
    rng_local = np.random.default_rng(cfg.random_seed)

    <span>for</span> sigma <span>in</span> cfg.noise_sigmas:
        corr = {<span>"gini"</span>: [], <span>"hhi"</span>: [], <span>"entropy"</span>: []}
        <span>for</span> _ <span>in</span> <span>range</span>(cfg.robustness_repeats):
            noisy = shares.values * (<span>1</span> + rng_local.normal(<span>0</span>, sigma, size=shares.shape))
            noisy = np.clip(noisy, <span>0</span>, <span>None</span>)
            s = noisy.<span>sum</span>(axis=<span>1</span>, keepdims=<span>True</span>)
            s[s == <span>0</span>] = <span>1.0</span>
            noisy = noisy / s

            nm = compute_metrics(pd.DataFrame(noisy, index=shares.index, columns=shares.columns))
            corr[<span>"gini"</span>].append(<span>float</span>(stats.spearmanr(b[<span>"gini"</span>], nm[<span>"gini"</span>].values).statistic))
            corr[<span>"hhi"</span>].append(<span>float</span>(stats.spearmanr(b[<span>"hhi"</span>], nm[<span>"hhi"</span>].values).statistic))
            corr[<span>"entropy"</span>].append(<span>float</span>(stats.spearmanr(b[<span>"entropy"</span>], nm[<span>"normalized_entropy"</span>].values).statistic))

        <span>for</span> m <span>in</span> [<span>"gini"</span>, <span>"hhi"</span>, <span>"entropy"</span>]:
            v = np.array(corr[m], dtype=<span>float</span>)
            rows.append({
                <span>"sigma"</span>: <span>float</span>(sigma),
                <span>"metric"</span>: m,
                <span>"mean_rank_corr"</span>: <span>float</span>(v.mean()),
                <span>"ci_low"</span>: <span>float</span>(np.quantile(v, <span>0.025</span>)),
                <span>"ci_high"</span>: <span>float</span>(np.quantile(v, <span>0.975</span>)),
            })

    <span>return</span> pd.DataFrame(rows)

<span>def</span> <span>run_scenario</span>(<span>df: pd.DataFrame, cfg: Config, entity_col: <span>str</span></span>) -> <span>Dict</span>[<span>str</span>, <span>Any</span>]:
    x = df[[cfg.time_col, entity_col]].dropna().copy()
    x[<span>"time_bucket"</span>] = x[cfg.time_col].dt.floor(cfg.freq)

    counts = x.groupby([<span>"time_bucket"</span>, entity_col], as_index=<span>False</span>).size().rename(columns={<span>"size"</span>: <span>"value"</span>})
    panel = counts.pivot(index=<span>"time_bucket"</span>, columns=entity_col, values=<span>"value"</span>).fillna(<span>0.0</span>).sort_index()

    full_idx = pd.date_range(panel.index.<span>min</span>(), panel.index.<span>max</span>(), freq=cfg.freq)
    panel = panel.reindex(full_idx, fill_value=<span>0.0</span>)

    totals = panel.<span>sum</span>(axis=<span>1</span>)
    panel = panel.loc[totals >= cfg.min_total_per_bucket].copy()
    totals = panel.<span>sum</span>(axis=<span>1</span>)

    shares = panel.div(totals.replace(<span>0</span>, np.nan), axis=<span>0</span>).fillna(<span>0.0</span>)
    metrics = compute_metrics(shares)

    trend = {}
    raw_p = {}
    null_gini = <span>None</span>
    t = np.arange(<span>len</span>(metrics), dtype=<span>float</span>)

    <span>for</span> m <span>in</span> [<span>"gini"</span>, <span>"hhi"</span>, <span>"normalized_entropy"</span>]:
        y = metrics[m].values.astype(<span>float</span>)
        rho, p, null = block_permutation(y, cfg.n_permutations, block_size=<span>24</span>, seed=cfg.random_seed)
        mk = mann_kendall(y)
        slope, intercept, lo, hi = stats.theilslopes(y, t, <span>0.95</span>)

        trend[m] = {
            <span>"block_permutation"</span>: {<span>"rho"</span>: <span>float</span>(rho), <span>"p"</span>: <span>float</span>(p)},
            <span>"mann_kendall"</span>: mk,
            <span>"theil_sen"</span>: {
                <span>"slope"</span>: <span>float</span>(slope),
                <span>"slope_ci_low"</span>: <span>float</span>(lo),
                <span>"slope_ci_high"</span>: <span>float</span>(hi),
                <span>"intercept"</span>: <span>float</span>(intercept),
            },
        }
        raw_p[m] = <span>float</span>(p)
        <span>if</span> m == <span>"gini"</span>:
            null_gini = null

    q = fdr_bh(raw_p)
    <span>for</span> m <span>in</span> trend:
        trend[m][<span>"block_permutation"</span>][<span>"q_fdr"</span>] = <span>float</span>(q[m])
        trend[m][<span>"block_permutation"</span>][<span>"significant"</span>] = <span>bool</span>(q[m] < cfg.alpha)

    robust = robustness_mc(shares, cfg)

    diag = {
        <span>"rows"</span>: <span>int</span>(panel.shape[<span>0</span>]),
        <span>"entities"</span>: <span>int</span>(panel.shape[<span>1</span>]),
        <span>"nonzero_ratio"</span>: <span>float</span>((panel.values > <span>0</span>).mean()),
        <span>"active_entities_median"</span>: <span>float</span>((panel > <span>0</span>).<span>sum</span>(axis=<span>1</span>).median()),
        <span>"active_entities_p90"</span>: <span>float</span>((panel > <span>0</span>).<span>sum</span>(axis=<span>1</span>).quantile(<span>0.9</span>)),
    }

    summary = {
        <span>"entity_col"</span>: entity_col,
        <span>"rows"</span>: diag[<span>"rows"</span>],
        <span>"entities"</span>: diag[<span>"entities"</span>],
        <span>"nonzero_ratio"</span>: diag[<span>"nonzero_ratio"</span>],
        <span>"gini_current"</span>: <span>float</span>(metrics[<span>"gini"</span>].iloc[-<span>1</span>]),
        <span>"hhi_current"</span>: <span>float</span>(metrics[<span>"hhi"</span>].iloc[-<span>1</span>]),
        <span>"entropy_current"</span>: <span>float</span>(metrics[<span>"normalized_entropy"</span>].iloc[-<span>1</span>]),
        <span>"gini_rho"</span>: <span>float</span>(trend[<span>"gini"</span>][<span>"block_permutation"</span>][<span>"rho"</span>]),
        <span>"gini_q"</span>: <span>float</span>(trend[<span>"gini"</span>][<span>"block_permutation"</span>][<span>"q_fdr"</span>]),
        <span>"entropy_q"</span>: <span>float</span>(trend[<span>"normalized_entropy"</span>][<span>"block_permutation"</span>][<span>"q_fdr"</span>]),
    }

    <span>return</span> {
        <span>"diag"</span>: diag,
        <span>"metrics"</span>: metrics,
        <span>"trend"</span>: trend,
        <span>"robustness"</span>: robust,
        <span>"summary"</span>: summary,
        <span>"null_gini"</span>: null_gini,
    }

results = {}
<span>for</span> entity <span>in</span> cfg.entity_candidates:
    results[entity] = run_scenario(raw_df, cfg, entity)

summary_df = pd.DataFrame([results[e][<span>"summary"</span>] <span>for</span> e <span>in</span> results]).sort_values(<span>"entity_col"</span>)
summary_df

entity_col   rows  entities  nonzero_ratio  gini_current  \
1  execution_engine  10783         7       0.765928      0.609524   
0      orchestrator  10783         7       0.778752      0.400000   

   hhi_current  entropy_current  gini_rho    gini_q  entropy_q  
1     0.333333         0.873233 -0.007772  0.417194   0.417194  
0     0.217778         0.915138  0.002089  0.830057   0.830057

<span>def</span> <span>plot_result</span>(<span>entity: <span>str</span>, result: <span>Dict</span>[<span>str</span>, <span>Any</span>]</span>):
    metrics = result[<span>"metrics"</span>]
    trend = result[<span>"trend"</span>]
    null_gini = result[<span>"null_gini"</span>]
    diag = result[<span>"diag"</span>]

    red, amber, green = <span>"#ff7b72"</span>, <span>"#d29922"</span>, <span>"#3fb950"</span>
    fig, axes = plt.subplots(<span>2</span>, <span>2</span>, figsize=(<span>14</span>, <span>9</span>))
    t = np.arange(<span>len</span>(metrics), dtype=<span>float</span>)

    ax = axes[<span>0</span>, <span>0</span>]
    ax.plot(metrics.index, metrics[<span>"gini"</span>], color=red, label=<span>"Gini"</span>, linewidth=<span>2</span>)
    ax.plot(metrics.index, metrics[<span>"hhi"</span>], color=amber, label=<span>"HHI"</span>, linewidth=<span>2</span>)
    ax.plot(metrics.index, metrics[<span>"normalized_entropy"</span>], color=green, label=<span>"Entropy"</span>, linewidth=<span>2</span>)
    ax.set_title(<span>"Concentration metrics over time"</span>, fontweight=<span>"bold"</span>)
    ax.legend(framealpha=<span>0.3</span>)
    ax.grid(<span>True</span>)

    ax = axes[<span>0</span>, <span>1</span>]
    y = metrics[<span>"gini"</span>].values
    ts = trend[<span>"gini"</span>][<span>"theil_sen"</span>]
    yhat = ts[<span>"intercept"</span>] + ts[<span>"slope"</span>] * t
    ax.scatter(t, y, s=<span>10</span>, alpha=<span>0.7</span>, color=<span>"#8b949e"</span>)
    ax.plot(t, yhat, color=red, linewidth=<span>2.5</span>)
    ax.set_title(<span>"Gini trend (Theil-Sen)"</span>, fontweight=<span>"bold"</span>)
    ax.grid(<span>True</span>)

    ax = axes[<span>1</span>, <span>0</span>]
    rho = trend[<span>"gini"</span>][<span>"block_permutation"</span>][<span>"rho"</span>]
    p = trend[<span>"gini"</span>][<span>"block_permutation"</span>][<span>"p"</span>]
    ax.hist(null_gini, bins=<span>40</span>, density=<span>True</span>, color=<span>"#30363d"</span>, edgecolor=<span>"#8b949e"</span>, alpha=<span>0.85</span>)
    ax.axvline(rho, color=red, linestyle=<span>"--"</span>, linewidth=<span>2.5</span>, label=<span>f"rho=<span>{rho:<span>.3</span>f}</span>, p=<span>{p:<span>.4</span>g}</span>"</span>)
    ax.set_title(<span>"Block permutation null (Gini)"</span>, fontweight=<span>"bold"</span>)
    ax.legend(framealpha=<span>0.3</span>)
    ax.grid(<span>True</span>)

    ax = axes[<span>1</span>, <span>1</span>]
    ax.axis(<span>"off"</span>)
    txt = [
        <span>f"entity: <span>{entity}</span>"</span>,
        <span>f"rows/entities: <span>{diag[<span>'rows'</span>]}</span>/<span>{diag[<span>'entities'</span>]}</span>"</span>,
        <span>f"nonzero_ratio: <span>{diag[<span>'nonzero_ratio'</span>]:<span>.4</span>f}</span>"</span>,
        <span>f"active_entities_median: <span>{diag[<span>'active_entities_median'</span>]:<span>.1</span>f}</span>"</span>,
        <span>f"gini q(FDR): <span>{trend[<span>'gini'</span>][<span>'block_permutation'</span>][<span>'q_fdr'</span>]:<span>.4</span>g}</span>"</span>,
        <span>f"entropy q(FDR): <span>{trend[<span>'normalized_entropy'</span>][<span>'block_permutation'</span>][<span>'q_fdr'</span>]:<span>.4</span>g}</span>"</span>,
    ]
    ax.text(<span>0.02</span>, <span>0.98</span>, <span>"\n"</span>.join(txt), va=<span>"top"</span>, fontsize=<span>11</span>)

    fig.suptitle(<span>f"<span>{entity}</span> scenario"</span>, fontweight=<span>"bold"</span>)
    fig.tight_layout(rect=[<span>0</span>, <span>0</span>, <span>1</span>, <span>0.96</span>])
    plt.show()

<span>for</span> entity <span>in</span> results:
    plot_result(entity, results[entity])

fig, axes = plt.subplots(<span>1</span>, <span>2</span>, figsize=(<span>12</span>, <span>4.5</span>))
axes[<span>0</span>].bar(summary_df[<span>"entity_col"</span>], summary_df[<span>"gini_rho"</span>], color=<span>"#58a6ff"</span>)
axes[<span>0</span>].set_title(<span>"Observed Gini trend rho"</span>)
axes[<span>0</span>].grid(<span>True</span>, axis=<span>"y"</span>)

axes[<span>1</span>].bar(summary_df[<span>"entity_col"</span>], summary_df[<span>"gini_q"</span>], color=<span>"#ff7b72"</span>)
axes[<span>1</span>].axhline(cfg.alpha, color=<span>"#8b949e"</span>, linestyle=<span>"--"</span>, label=<span>f"alpha=<span>{cfg.alpha}</span>"</span>)
axes[<span>1</span>].set_title(<span>"Gini q-values (FDR)"</span>)
axes[<span>1</span>].legend()
axes[<span>1</span>].grid(<span>True</span>, axis=<span>"y"</span>)

plt.tight_layout()
plt.show()

图1

图2

图3

结果解读

如果集中度正在上升,我们预期:

  • Gini/HHI 趋势 > 0 且显著,
  • 熵趋势 < 0 且显著。

在本次数据集运行中,两种情形都应通过 FDR 校正后的 q 值和效应量(rho/斜率)来解读,而不仅仅是原始 p 值。

report = {
    <span>"metadata"</span>: {
        <span>"created_at"</span>: datetime.now().isoformat(),
        <span>"dataset_name"</span>: <span>"Data Engineering & Modern AI Pipelines 2026"</span>,
        <span>"config"</span>: asdict(cfg),
    },
    <span>"audit"</span>: audit,
    <span>"scenario_summary"</span>: summary_df.to_dict(orient=<span>"records"</span>),
    <span>"scenarios"</span>: {},
}

<span>for</span> entity, res <span>in</span> results.items():
    report[<span>"scenarios"</span>][entity] = {
        <span>"diag"</span>: res[<span>"diag"</span>],
        <span>"trend"</span>: res[<span>"trend"</span>],
        <span>"robustness"</span>: res[<span>"robustness"</span>].to_dict(orient=<span>"records"</span>),
        <span>"current_metrics"</span>: {
            <span>"gini"</span>: <span>float</span>(res[<span>"metrics"</span>][<span>"gini"</span>].iloc[-<span>1</span>]),
            <span>"hhi"</span>: <span>float</span>(res[<span>"metrics"</span>][<span>"hhi"</span>].iloc[-<span>1</span>]),
            <span>"entropy"</span>: <span>float</span>(res[<span>"metrics"</span>][<span>"normalized_entropy"</span>].iloc[-<span>1</span>]),
        },
    }

out_path = Path(cfg.output_dir) / <span>"kaggle_pipeline_concentration_stability_report.json"</span>
<span>with</span> <span>open</span>(out_path, <span>"w"</span>, encoding=<span>"utf-8"</span>) <span>as</span> f:
    json.dump(report, f, indent=<span>2</span>, default=<span>float</span>)

<span>print</span>(<span>"Saved report:"</span>, out_path)
summary_df

Saved report: kaggle_outputs_pipeline_concentration_stability/kaggle_pipeline_concentration_stability_report.json

entity_col   rows  entities  nonzero_ratio  gini_current  \
1  execution_engine  10783         7       0.765928      0.609524   
0      orchestrator  10783         7       0.778752      0.400000   

   hhi_current  entropy_current  gini_rho    gini_q  entropy_q  
1     0.333333         0.873233 -0.007772  0.417194   0.417194  
0     0.217778         0.915138  0.002089  0.830057   0.830057

数据与完整代码来源:

<span>https:</span>/<span>/mbd.pub/o</span><span>/bread/</span><span>YZaVmJhp</span>ZQ