{"metadata":{"kernelspec":{"language":"python","display_name":"Python 3","name":"python3"},"language_info":{"name":"python","version":"3.12.13","mimetype":"text/x-python","codemirror_mode":{"name":"ipython","version":3},"pygments_lexer":"ipython3","nbconvert_exporter":"python","file_extension":".py"},"kaggle":{"accelerator":"none","dataSources":[],"dockerImageVersionId":28755,"isInternetEnabled":false,"language":"python","sourceType":"notebook","isGpuEnabled":false}},"nbformat_minor":4,"nbformat":4,"cells":[{"cell_type":"markdown","source":"# Home Credit：表内聚合前的数据排查\n\n每次查看一张逻辑表，帮助团队判断表类型，并讨论缺失字段是否有保留价值。只需修改 `TABLE_NAME`，然后运行全部单元格。","metadata":{"execution":{"iopub.status.busy":"2026-08-19T16:44:12.694368Z","iopub.execute_input":"2026-08-19T16:44:12.694671Z","iopub.status.idle":"2026-08-19T16:44:12.702440Z","shell.execute_reply.started":"2026-08-19T16:44:12.694623Z","shell.execute_reply":"2026-08-19T16:44:12.701449Z"}}},{"cell_type":"code","source":"# 1. 选择一张表\n\nTABLE_NAME = \"static_0\"\nSAMPLE_ROWS_PER_SHARD = 100\nDISPLAY_ROWS = 10\n\n# 可选：base, static_0, static_cb_0, applprev_1, applprev_2,\n# credit_bureau_a_1, credit_bureau_a_2, credit_bureau_b_1, credit_bureau_b_2,\n\n# debitcard_1, deposit_1, other_1, person_1, person_2,\n# tax_registry_a_1, tax_registry_b_1, tax_registry_c_1","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2026-08-19T16:44:21.630459Z","iopub.execute_input":"2026-08-19T16:44:21.630773Z","iopub.status.idle":"2026-08-19T16:44:21.636065Z","shell.execute_reply.started":"2026-08-19T16:44:21.630749Z","shell.execute_reply":"2026-08-19T16:44:21.635232Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"#  2. 环境和逻辑表配置\n\nfrom collections import Counter\nfrom pathlib import Path\nimport numpy as np\nimport pandas as pd\nimport polars as pl\nimport pyarrow.parquet as pq\nfrom IPython.display import display\n\npd.set_option(\"display.max_columns\", 200)\npd.set_option(\"display.max_colwidth\", 120)\n\nDATA_ROOT = Path(\"/kaggle/input/competitions/home-credit-credit-risk-model-stability\")\nTRAIN_ROOT = DATA_ROOT / \"parquet_files\" / \"train\"\nTABLE_PATTERNS = {\n    \"base\": \"train_base.parquet\",\n    \"static_0\": \"train_static_0_*.parquet\",\n    \"static_cb_0\": \"train_static_cb_0.parquet\",\n    \"applprev_1\": \"train_applprev_1_*.parquet\",\n    \"applprev_2\": \"train_applprev_2.parquet\",\n    \"credit_bureau_a_1\": \"train_credit_bureau_a_1_*.parquet\",\n    \"credit_bureau_a_2\": \"train_credit_bureau_a_2_*.parquet\",\n    \"credit_bureau_b_1\": \"train_credit_bureau_b_1.parquet\",\n    \"credit_bureau_b_2\": \"train_credit_bureau_b_2.parquet\",\n    \"debitcard_1\": \"train_debitcard_1.parquet\",\n    \"deposit_1\": \"train_deposit_1.parquet\",\n    \"other_1\": \"train_other_1.parquet\",\n    \"person_1\": \"train_person_1.parquet\",\n    \"person_2\": \"train_person_2.parquet\",\n    \"tax_registry_a_1\": \"train_tax_registry_a_1.parquet\",\n    \"tax_registry_b_1\": \"train_tax_registry_b_1.parquet\",\n    \"tax_registry_c_1\": \"train_tax_registry_c_1.parquet\",\n}\n\nassert TABLE_NAME in TABLE_PATTERNS, f\"未知表名：{TABLE_NAME}\"\nSHARDS = tuple(sorted(TRAIN_ROOT.glob(TABLE_PATTERNS[TABLE_NAME])))\nassert SHARDS, f\"没有匹配到文件：{TABLE_PATTERNS[TABLE_NAME]}\"\n\nlf = pl.scan_parquet([str(path) for path in SHARDS])\nschema = dict(lf.collect_schema())\nkeys = [c for c in [\"case_id\", \"num_group1\", \"num_group2\"] if c in schema]\nprint(\"当前表：\", TABLE_NAME)","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2026-08-19T16:44:24.289821Z","iopub.execute_input":"2026-08-19T16:44:24.290113Z","iopub.status.idle":"2026-08-19T16:44:25.348183Z","shell.execute_reply.started":"2026-08-19T16:44:24.290088Z","shell.execute_reply":"2026-08-19T16:44:25.347142Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"","metadata":{"trusted":true},"outputs":[],"execution_count":null},{"cell_type":"code","source":"# 3. 表和 shard 信息；\n# （1）这里只读取 Parquet metadata，不加载完整数据\n# （2）\n\n\nshard_rows, shard_schemas = [], []\nfor path in SHARDS:\n    meta = pq.ParquetFile(path).metadata\n    file_schema = dict(pl.read_parquet_schema(path))\n    shard_schemas.append(file_schema)\n    shard_rows.append({\n        \"shard\": path.name,\n        \"rows\": meta.num_rows,\n        \"columns\": len(file_schema),\n        \"size_mb\": round(path.stat().st_size / 1024**2, 2),\n    })\n\nshard_info = pd.DataFrame(shard_rows)\nN_ROWS = int(shard_info[\"rows\"].sum())\nschema_consistent = all(s == shard_schemas[0] for s in shard_schemas[1:])\ndisplay(shard_info)\ndisplay(pd.DataFrame([{\n    \"table\": TABLE_NAME, \"n_shards\": len(SHARDS), \"n_rows\": N_ROWS,\n    \"n_columns\": len(schema), \"keys\": \" + \".join(keys),\n    \"schema_consistent\": schema_consistent,\n}]))\nassert schema_consistent, \"不同 shard 的 schema 不一致，需要单独检查。\"","metadata":{"trusted":true},"outputs":[],"execution_count":null},{"cell_type":"code","source":"# 4. 字段名、dtype 和字段说明¶\n\nfeature_description = {}\nfeature_path = DATA_ROOT / \"feature_definitions.csv\"\nif feature_path.exists():\n    feature_def = pd.read_csv(feature_path)\n    if {\"Variable\", \"Description\"}.issubset(feature_def.columns):\n        feature_description = dict(zip(feature_def[\"Variable\"], feature_def[\"Description\"]))\n\ndef column_role(name, dtype):\n    if name in keys: return \"key\"\n    if dtype.is_temporal(): return \"date/time\"\n    if dtype.is_numeric(): return \"numeric\"\n    if str(dtype) in {\"String\", \"Utf8\", \"Categorical\", \"Enum\", \"Boolean\"}: return \"categorical\"\n    return \"other\"\n\ncolumns = pd.DataFrame([\n    {\"column\": name, \"dtype\": str(dtype), \"role\": column_role(name, dtype),\n     \"description\": feature_description.get(name)}\n    for name, dtype in schema.items()\n])\ndisplay(columns.head(68))\n","metadata":{"trusted":true},"outputs":[],"execution_count":null},{"cell_type":"code","source":"# 5. Sample rows：每个 shard 只取少量开头记录，用于理解一行代表什么\n\nsamples = pl.concat([\n    pl.read_parquet(path, n_rows=SAMPLE_ROWS_PER_SHARD)\n      .with_columns(pl.lit(path.name).alias(\"__shard\"))\n    for path in SHARDS\n])\ndisplay(samples.head(DISPLAY_ROWS).to_pandas())\ndisplay(samples.head(60).to_pandas())","metadata":{"trusted":true},"outputs":[],"execution_count":null},{"cell_type":"code","source":"# 6. 缺失、unique、范围和 flags\n# 缺失值来自全表精确扫描，unique 使用近似值。top values 来自小样本。flags 只用于讨论，不代表应直接删除字段\n\nimport gc\nimport psutil, os\nfrom collections import Counter\nimport numpy as np\nimport pandas as pd\nimport polars as pl\n\ndef mem(tag=\"\"):\n    gc.collect()\n    proc = psutil.Process(os.getpid()).memory_info().rss / 1024**3\n    avail = psutil.virtual_memory().available / 1024**3\n    print(f\"  💾 [{tag}] 进程: {proc:.2f} GB | 可用: {avail:.1f} GB\")\n\n\nBATCH_SIZE = 10  # 🔽 减到 10,更保守\n\nprofile_rows = []\nschema_items = list(schema.items())\n\nmem(\"开始分批统计\")\n\n# ============================================\n# 阶段 1: 分批算 stats(null_count / n_unique / min / max)\n# ============================================\nall_stats = {}  # 存所有列的统计结果\n\nfor batch_start in range(0, len(schema_items), BATCH_SIZE):\n    batch = schema_items[batch_start:batch_start + BATCH_SIZE]\n    batch_end = min(batch_start + BATCH_SIZE, len(schema_items))\n    print(f\"  处理第 {batch_start+1}~{batch_end} 列 ({len(batch)} 列)...\")\n    \n    exprs = []\n    for i, (name, dtype) in enumerate(batch):\n        gi = batch_start + i\n        exprs.extend([\n            pl.col(name).null_count().alias(f\"null_{gi}\"),\n            pl.col(name).drop_nulls().approx_n_unique().alias(f\"unique_{gi}\"),\n        ])\n        if dtype.is_numeric() or dtype.is_temporal():\n            exprs.extend([\n                pl.col(name).min().alias(f\"min_{gi}\"),\n                pl.col(name).max().alias(f\"max_{gi}\"),\n            ])\n    \n    # 分批 collect\n    batch_stats = lf.select(exprs).collect(engine=\"streaming\").row(0, named=True)\n    \n    # 只保存需要的数字,不保留 Polars 内部引用\n    for k, v in batch_stats.items():\n        all_stats[k] = v\n    \n    # 强力释放\n    del batch_stats, exprs\n    gc.collect()\n    mem(f\"  第 {batch_end} 列后\")\n\n# ============================================\n# 阶段 2: 一次性生成 profile(不涉及大内存操作)\n# ============================================\nprint(\"\\n生成最终 profile ...\")\nmem(\"阶段 2 开始\")\n\nfor i, (name, dtype) in enumerate(schema_items):\n    null_count = int(all_stats.get(f\"null_{i}\") or 0)\n    non_null = N_ROWS - null_count\n    n_unique = int(all_stats.get(f\"unique_{i}\") or 0)\n    null_rate = null_count / N_ROWS if N_ROWS else np.nan\n    unique_rate = n_unique / non_null if non_null else np.nan\n    \n    # 从 samples 取 top values(samples 只有几百行,内存开销可控)\n    try:\n        values = samples.get_column(name).drop_nulls().to_list()\n        top_values = Counter(map(repr, values)).most_common(5)\n        top1_rate = top_values[0][1] / len(values) if values else np.nan\n    except Exception:\n        top_values = []\n        top1_rate = np.nan\n    \n    flags = []\n    if null_rate == 1: flags.append(\"all_missing\")\n    elif null_rate >= 0.99: flags.append(\"missing>=99%\")\n    if n_unique <= 1: flags.append(\"constant_candidate\")\n    elif top1_rate >= 0.99: flags.append(\"near_constant_in_sample\")\n    if name not in keys and pd.notna(unique_rate) and unique_rate >= 0.98:\n        flags.append(\"id_like\")\n    \n    profile_rows.append({\n        \"column\": name, \"dtype\": str(dtype), \"role\": column_role(name, dtype),\n        \"description\": feature_description.get(name),\n        \"null_count\": null_count, \"null_rate\": null_rate,\n        \"approx_n_unique\": n_unique, \"approx_unique_rate\": unique_rate,\n        \"min\": all_stats.get(f\"min_{i}\"), \"max\": all_stats.get(f\"max_{i}\"),\n        \"sample_top_values\": top_values,\n        \"flags\": \" | \".join(flags),\n    })\n\ndel all_stats\ngc.collect()\nmem(\"全部完成\")\n\ncolumn_profile = pd.DataFrame(profile_rows)\nprint(f\"\\n✓ 生成 profile: {len(column_profile)} 个字段\")\ndisplay(column_profile)","metadata":{"trusted":true},"outputs":[],"execution_count":null},{"cell_type":"code","source":"# 6. 重点讨论字段\n# 集中显示高缺失、常数候选、近常数和 ID-like 字段。结合字段含义讨论是否保留\n\ndisplay(\n    column_profile[column_profile[\"flags\"] != \"\"]\n    .sort_values([\"null_rate\", \"column\"], ascending=[False, True])\n)","metadata":{"trusted":true},"outputs":[],"execution_count":null},{"cell_type":"code","source":"# 7. 每个 case 对应多少行\n# 用于区分 static 表和 repeated/history 表\n\n\nif \"case_id\" in schema:\n    case_counts = lf.select(\"case_id\").group_by(\"case_id\").len().rename({\"len\": \"rows_per_case\"})\n    repetition = case_counts.select(\n        pl.len().alias(\"n_cases\"),\n        pl.col(\"rows_per_case\").mean().alias(\"mean\"),\n        pl.col(\"rows_per_case\").median().alias(\"p50\"),\n        pl.col(\"rows_per_case\").quantile(0.90).alias(\"p90\"),\n        pl.col(\"rows_per_case\").quantile(0.95).alias(\"p95\"),\n        pl.col(\"rows_per_case\").quantile(0.99).alias(\"p99\"),\n        pl.col(\"rows_per_case\").max().alias(\"max\"),\n        (pl.col(\"rows_per_case\") == 1).cast(pl.Float64).mean().alias(\"one_row_case_rate\"),\n    ).collect(engine=\"streaming\")\n    display(repetition.to_pandas())\nelse:\n    print(\"没有 case_id，跳过。\")","metadata":{"trusted":true},"outputs":[],"execution_count":null},{"cell_type":"code","source":"# 8. 日期字段候选\n\ndate_candidates = column_profile[\n    column_profile[\"role\"].eq(\"date/time\")\n    | column_profile[\"column\"].str.lower().str.contains(\"date\")\n    | column_profile[\"column\"].str.lower().str.endswith(\"d\")\n][[\"column\", \"dtype\", \"description\", \"null_rate\", \"approx_n_unique\", \"min\", \"max\"]]\ndisplay(date_candidates)","metadata":{"trusted":true},"outputs":[],"execution_count":null},{"cell_type":"code","source":"# 9. 团队判断\n# 讨论以下问题：\n\n# 一行代表客户、合同、账户、事件、月度记录还是明细？\n# case_id 为什么重复？num_group1 / num_group2 表示什么层级？\n# 是否存在代表记录时间的字段？\n# 高缺失字段是否有业务意义，缺失本身是否有信息？\n# 常数、近常数和 ID-like 字段是否有保留价值？\n# 最终类型：static/entity、repeated entity、history/event、nested/detail 或 other。\n\n# 人工记录\n# 表类型：\n# 一行的含义：\n# 实体层级：\n# 主要时间字段：\n# 建议保留/删除的字段：\n# 后续聚合层级：\n# 备注：\n\n\nimport psutil, os\nmem = psutil.virtual_memory()\nproc = psutil.Process(os.getpid()).memory_info().rss / 1024**3\nprint(f\"当前 Python 进程占用: {proc:.2f} GB\")\nprint(f\"系统总内存: {mem.total/1024**3:.1f} GB\")\nprint(f\"已用内存: {mem.used/1024**3:.1f} GB ({mem.percent}%)\")\nprint(f\"可用内存: {mem.available/1024**3:.1f} GB\")","metadata":{"trusted":true},"outputs":[],"execution_count":null},{"cell_type":"code","source":"# 释放所有大变量\nfor name in [\"df\", \"samples\", \"lf\", \"all_results\", \"result_df\"]:\n    if name in globals():\n        del globals()[name]\n\nimport gc\ngc.collect()\n\n# 检查是否释放成功\nimport psutil, os\nproc = psutil.Process(os.getpid()).memory_info().rss / 1024**3\nprint(f\"释放后进程占用: {proc:.2f} GB\")","metadata":{"trusted":true},"outputs":[],"execution_count":null}]}