# JD电商搜索与转化数据仓库 **Repository Path**: akatini/SearchFlow ## Basic Information - **Project Name**: JD电商搜索与转化数据仓库 - **Description**: 这是一个离线数据仓库,将京东提供的匿名用户历史行为、末次搜索请求、候选商品曝光和商品元数据统一建模,并支持搜索漏斗、商品转化、店铺分类表现、用户画像、长尾商品分析和数据质量监控的离线分析型数据仓库。 - **Primary Language**: Unknown - **License**: Not specified - **Default Branch**: master - **Homepage**: None - **GVP Project**: No ## Statistics - **Stars**: 0 - **Forks**: 0 - **Created**: 2026-07-20 - **Last Updated**: 2026-07-31 ## Categories & Tags **Categories**: Uncategorized **Tags**: None ## README # SearchFlow 数据仓库项目 ## 1. 项目业务背景 SearchFlow 是一个面向 **JDsearch 公开搜索日志**的小型数据仓库项目,专注于把匿名化的搜索行为数据加工成可供分析使用的层次化模型。 项目的核心业务目标: - **理解查询模式**:什么样的查询词、什么样的查询词组合,会带来什么样的商品曝光与交互 - **发现排序问题**:在搜索结果中频繁曝光但几乎没人交互的"靠前低交互"商品 - **识别强匹配**:在多个搜索请求中稳定获得高交互率的商品 - **支撑运营诊断**:把诊断结论沉淀到 ADS 层表,供运营/算法直接查询 ## 2. 原始数据说明 数据源来自 [JDsearch 公开数据集](https://github.com/rec-research/JDsearch)。一份原始数据包含两个文本文件: ### `product_meta_data.txt` 商品元数据,每行一个商品,字段以 `\x18` 分隔: ``` product_id product_name_terms brand_id brand_name_terms category_1_id category_1_name_terms category_2_id category_2_name_terms category_3_id category_3_name_terms category_4_id category_4_name_terms shop_id ``` 所有名称字段都是**匿名化词项 ID 列表**(不是真实文本),例如 `product_name_terms = ["63995226", "60257419"]`。 ### `user_behavior_data.txt` 用户行为日志,每行表示一次搜索请求。字段以 `\t` 分隔: ``` query candidate_wid_list candidate_label_list history_qry_list history_wid_list history_type_list history_time_list ``` 字段说明: - `query`:本次搜索的查询词项 ID 列表 - `candidate_wid_list` / `candidate_label_list`:候选商品 ID 与交互标签(`_` 分隔的等长数组) - `history_qry_list` / `history_wid_list` / `history_type_list` / `history_time_list`:用户历史交互(同样为等长数组) ### 字段值编码 | 字段 | 取值 | 含义 | |:-----|:-----|:-----| | `candidate_label` | `"0.0"` ~ `"3.0"` | NONE / CLICK / CART / PURCHASE | | `history_type` | `CLICK`/`CART`/`ORD`/`FLW` | 历史行为类型 | | `history_qry` 中 `-1` | — | 表示"该历史行为无关联搜索" | ### 数据集的关键限制(重要) 为了避免误读项目,必须明确以下边界: - **没有真实 `user_id`**:数据集只提供匿名快照键,无法跨搜索请求关联用户 - **没有订单金额 / 商品价格**:仅能反映行为层级,不能做收益分析 - **没有绝对时间戳**:所有时间字段都是相对于"测试搜索"的相对偏移量 - **`-1` 不代表特定入口**:仅表示"该历史行为没有关联的搜索",具体是收藏夹、推荐页还是直接访问均无法判断 - **`candidate_label` 是最高交互层级**:例如某商品被同时点击和加购,标签只记录为 CART(更高层级) ## 3. 分层架构 项目按数据加工阶段分为 6 层: ``` ┌─────────────────────┐ │ Source Audit │ ← 数据源审计 │ (源文件统计) │ └──────────┬──────────┘ │ ┌──────────▼──────────┐ │ Landing │ ← Parquet 文件落地 │ (4 个数据集目录) │ └──────────┬──────────┘ │ ┌──────────▼──────────┐ │ ODS │ ← 贴源层 │ (4 张表) │ └──────────┬──────────┘ │ ┌──────────▼──────────┐ │ DIM │ ← 维度层 │ (3 张维表) │ └──────────┬──────────┘ │ ┌──────────▼──────────┐ │ DWD │ ← 明细事实层 │ (3 张事实表) │ └──────────┬──────────┘ │ ┌──────────▼──────────┐ │ DWS │ ← 汇总层 │ (3 张宽表) │ └──────────┬──────────┘ │ ┌──────────▼──────────┐ │ ADS │ ← 应用层 │ (2 张诊断表) │ └─────────────────────┘ ``` 每层都有独立的 DDL/加载逻辑、对账机制、批次追踪。 ## 4. 各表粒度与职责 ### Landing(Parquet) | 目录 | 粒度 | 说明 | |:-----|:-----|:-----| | `product_snapshot/` | 一个源行为行对应一个商品 | 商品快照 | | `search_request/` | 一个源行为行 | 搜索请求 | | `candidate_exposure/` | 一个源行内的一个候选位置 | 候选商品曝光 | | `history_interaction/` | 一个源行内的一条历史交互 | 历史交互 | ### ODS(贴源) | 表 | 粒度 | 主键 | |:---|:-----|:-----| | `ods_product_snapshot` | 一个源行 | `(dataset_version, source_file, source_row_number)` | | `ods_search_request` | 一次搜索 | `(dataset_version, search_request_key)` | | `ods_candidate_exposure` | 一个候选位置 | `(dataset_version, candidate_exposure_key)` | | `ods_history_interaction` | 一条历史交互 | `(dataset_version, history_interaction_key)` | ### DIM(维度) | 表 | 粒度 | 主键 | |:---|:-----|:-----| | `dim_product` | 一个 `(dataset_version, product_id)` | 同左 | | `dim_query` | 一个查询词项序列 | `(dataset_version, query_key)` | | `dim_interaction_type` | 一个交互类型 | `(interaction_domain, source_code)` | 每张表都有特殊成员(UNKNOWN / NO_QUERY)作为兜底行。 ### DWD(明细事实) | 表 | 粒度 | 主键 | |:---|:-----|:-----| | `dwd_search_request` | 一次搜索 | `(dataset_version, search_request_key)` | | `dwd_candidate_exposure` | 一个候选曝光 | `(dataset_version, candidate_exposure_key)` | | `dwd_history_interaction` | 一条历史交互 | `(dataset_version, history_interaction_key)` | 每张表通过代理键关联到维度表。 ### DWS(汇总) | 表 | 粒度 | 用途 | |:---|:-----|:-----| | `dws_query_performance` | 一个 query_sk | 查询维度的整体表现 | | `dws_query_product_performance` | 一个 `(query_sk, product_id)` | 查询 × 商品粒度的曝光与交互 | | `dws_search_context_summary` | 一个 search_request_key | 单次搜索请求的候选曝光 + 历史行为摘要 | ### ADS(应用层) | 表 | 粒度 | 用途 | |:---|:-----|:-----| | `ads_query_performance_dashboard` | 一个 query_sk | 反范式化宽表,附 6 项 DENSE_RANK 与漏斗状态 | | `ads_query_product_opportunity` | 一个 `(query_sk, product_id)` | 6 类机会分类诊断(INSUFFICIENT_SUPPORT / QUERY_WIDE_NO_INTERACTION / FRONT_POSITION_LOW_INTERACTION / HIGH_EXPOSURE_LOW_INTERACTION / STRONG_MATCH / NORMAL) | ## 5. 完整数据流 ``` 源文件 (.txt) │ Source Audit(行数预估、字段校验) ▼ Landing Parquet │ DELETE + INSERT(按版本覆盖) ▼ ODS DuckDB │ UPSERT(按业务键合并) ▼ DIM DuckDB │ DELETE + INSERT(按版本覆盖) ▼ DWD DuckDB │ GROUP BY(按维度聚合) ▼ DWS DuckDB │ 反范式化 + 排名 + 状态派生 ▼ ADS DuckDB ``` 每个 ETL 阶段都通过 `control.etl_batch` 与 `control.etl_table_load` 两张表记录运行状态,支持: - 失败重试(FAILED → RUNNING → SUCCESS) - 批次审计(每张表的源行数 / 加载行数 / 状态) - 历史回溯(按 `pipeline_name + dataset_version + status` 查询) ## 6. 环境安装方法 ### 系统要求 - Linux / macOS - Python 3.11+ - Conda(推荐 miniconda 或 mamba) ### 安装步骤 ```bash # 创建并激活 conda 环境 conda create -n searchflow python=3.11 -y conda activate searchflow # 安装项目依赖 pip install -e . # 验证安装 searchflow --help ``` ### 数据准备 把原始数据集放到指定路径(路径在配置文件中可修改): ``` data/raw/jdsearch/v1/archive/JDsearch.tar.gz ← 原始压缩包 data/raw/jdsearch/v1/extracted/JDsearch/ ← 解压后的目录 ├── product_meta_data.txt └── user_behavior_data.txt ``` 如果只想跑测试,可以直接用项目提供的 sample 数据,无需完整数据集。 ## 7. 全链路执行命令 ### 一键构建(推荐) ```bash searchflow build-all --dataset jdsearch --version v1 --overwrite ``` 依次执行:Landing → ODS → DIM → DWD → DWS → ADS,**任意一层失败立即停止**。 ### 分层构建 ```bash searchflow build-landing --dataset jdsearch --version v1 --overwrite searchflow build-ods --dataset jdsearch --version v1 --overwrite searchflow build-dim --dataset jdsearch --version v1 --overwrite searchflow build-dwd --dataset jdsearch --version v1 --overwrite searchflow build-dws --dataset jdsearch --version v1 --overwrite searchflow build-ads --dataset jdsearch --version v1 --overwrite ``` ### 检查每层数据 ```bash searchflow inspect-ods --dataset jdsearch --version v1 searchflow inspect-dim --dataset jdsearch --version v1 searchflow inspect-dwd --dataset jdsearch --version v1 searchflow inspect-dws --dataset jdsearch --version v1 searchflow inspect-ads --dataset jdsearch --version v1 ``` ### 审计数据源 ```bash searchflow audit-source --dataset jdsearch --version v1 ``` ## 8. 数据质量与对账机制 ### 对账层次 每个加载流程都有对应的对账模块: | 层 | 对账内容 | |:---|:---------| | Landing | 源文件行数 vs 写入 Parquet 行数 | | ODS | Parquet 行数 vs DuckDB 表行数 | | DIM | ODS 行数 vs 维表行数;属性完整度评分 | | DWD | 漏斗层级约束、计数与布尔一致性、维度关联成功率 | | DWS | 指标完整复制、查询覆盖率、跨 DWS 一致性 | | ADS | 反范式化字段复制、排名合法性、状态一致性、机会分类规则 | ### 对账失败处理 - **同事务回滚**:对账失败时,事务立即回滚,已写入的数据不会被 commit - **批次状态标记**:批次表记录 FAILED 状态与错误信息 - **幂等重试**:使用 `--overwrite` 重新执行,可清除旧数据后重建 ### 关键 CHECK 约束 所有 ADS 表都通过 CHECK 约束做硬一致性保护: - `has_purchase = (purchase_count > 0)` - `query_result_status` 与 `has_*` 一致性 - 漏斗层级 `purchase ≤ cart ≤ interaction ≤ exposure` - 比率字段 `[0, 1]` 范围 ## 9. 幂等和事务机制 ### 幂等 每张表的加载都遵循"先 DELETE 当前版本 + 再 INSERT"模式,所以: - 不带 `--overwrite`:如果版本已存在,抛出 `VersionConflictError` - 带 `--overwrite`:删除旧数据后重建,结果与第一次完全一致 - 行数 / 指标 / 排名在重复执行后保持稳定;批次 ID 与时间戳会更新 ### 事务 每个 ETL 流程都包裹在事务中: ``` BEGIN DELETE 旧版本数据 INSERT 新数据 对账检查(FAIL → ROLLBACK) 记录批次审计 COMMIT ``` 对账失败时立即回滚,不会留下半完成数据。 ### 批次追踪 `control.etl_batch` 记录每个 pipeline 的运行: - `pipeline_name`(如 `jdsearch_ods_build`) - `dataset_name` / `dataset_version` - `status`(RUNNING / SUCCESS / FAILED) - `started_at` / `finished_at` - `error_message`(失败时填充) `control.etl_table_load` 记录每张表的加载情况: - `target_table` - `source_row_count` / `loaded_row_count` - `target_version_row_count` - `status` ## 10. 项目目录说明 ``` SearchFlow/ ├── src/searchflow/ │ ├── common/ # 配置、日志、异常、标识符生成 │ ├── ingestion/ # 源数据审计 │ ├── parsing/ # 解析器(行为、商品、词项、相对时间) │ ├── landing/ # Landing 层(Parquet 写入、流水线) │ ├── warehouse/ # ODS/DIM/DWD/DWS/ADS 流水线与对账 │ ├── quality/ # 异常检测与报告 │ └── cli.py # 命令行入口 │ ├── sql/ │ ├── ddl/ # CREATE TABLE 定义 │ │ ├── control/ # 批次与表加载审计表 │ │ ├── ods/ # 贴源层 DDL │ │ ├── dim/ # 维度层 DDL │ │ ├── dwd/ # 明细层 DDL │ │ ├── dws/ # 汇总层 DDL │ │ ├── ads/ # 应用层 DDL │ │ └── common/ # 共享宏(如 query_key 生成) │ └── load/ # 加载 SQL(每张表一个) │ ├── dim/ │ ├── dws/ │ └── ads/ │ ├── configs/ # 配置文件与数据合同 ├── data/ # 原始数据 + sample 数据 ├── tests/ # 单元测试 + 集成测试 │ ├── unit/ # 单文件单函数级别测试 │ └── integration/ # 端到端流程测试 ├── reports/ # 报告输出(异常、对账) ├── pyproject.toml # 项目配置与依赖 └── README.md # 本文档 ``` ## 11. 测试 ```bash # 全部测试 pytest tests/ # 单元测试 pytest tests/unit/ # 集成测试 pytest tests/integration/ # 单独跑某个层 pytest tests/integration/test_ads_pipeline.py -v ``` 测试覆盖: - Parsing:行为解析、商品解析、词项解析、相对时间、键生成 - Landing:Parquet 写入、行数对账、批次审计 - Warehouse:每个 pipeline 的端到端流程、对账逻辑、批次审计 - Quality:异常检测、报告生成 ## 12. 主要分析结果(机会分类示例) ADS 层的 `ads_query_product_opportunity` 表对每个 `(query_sk, product_id)` 组合做 6 类分类诊断: | 分类 | 业务含义 | |:-----|:---------| | `INSUFFICIENT_SUPPORT` | 该 (query, product) 组合涉及的搜索请求数过少,无法得出统计意义结论 | | `QUERY_WIDE_NO_INTERACTION` | 整个查询都没有用户交互,可能是长尾查询 | | `FRONT_POSITION_LOW_INTERACTION` | 商品在搜索结果靠前位置曝光但交互率偏低,疑似排序问题 | | `HIGH_EXPOSURE_LOW_INTERACTION` | 商品被频繁曝光但交互率明显低于同查询 P25,可能是错配 | | `STRONG_MATCH` | 高曝光 + 高交互率,符合用户搜索意图的优质商品 | | `NORMAL` | 其他正常组合 | 分类逻辑使用同查询内的分位基准(P25 / P75),避免跨查询比较带来的偏差。 ## 13. 技术栈 - **Python 3.11** + **DuckDB**(嵌入式 OLAP 引擎) - **PyArrow** + **Parquet**(Landing 层文件格式) - **Typer** + **Rich**(CLI) - **Pytest**(测试框架) 选 DuckDB 的原因: - 嵌入式部署,无需独立服务 - 与 Parquet 生态深度集成,读写性能优秀 - 支持窗口函数 / 分位数 / CHECK 约束等高级特性 - 适合中等规模(千万级行)数据分析场景