3.9 KiB
Raw Blame History

Phase 5: 数据管道优化 — 技术研究

研究日期: 2026-03-21
阶段目标: 入库去重可靠30天窗口、公司清洗流程顺畅、公司招聘信息写入 ClickHouse


1. 现状分析

1.1 DATA-0130 天窗口去重

位置: app/services/ingest/dedup.pybatch_dedup_filter()

现状(有缺口):

-- 当前去重 SQL无时间窗口
SELECT job_id FROM job_data.boss_job WHERE job_id IN {keys}

缺口: 查的是全量历史数据,同一职位 30 天内只入库一次的语义没有实现。 理论上 ClickHouse 表可能有几年的旧数据,老 job_id 永远不会重新入库。

正确逻辑: 查询条件加 AND created_at > now() - INTERVAL 30 DAY 只在近 30 天内去重,超过 30 天的同一职位可以重新入库。

修复范围: batch_dedup_filter() 的两个 SQL 查询各加一行 WHERE 条件。


1.2 DATA-02统一入库管道spiderJobs 推送)

位置: app/services/ingest/service.py(已有 IngestService.store_batch

现状(有缺口):

  • 后端 cleaning.pycompany_jobs_sync.py 已通过 IngestService.store_batch() 入库
  • 外部脚本 spiderJobs/ 是独立运行的,它们通过 HTTP 推送到服务端 API
  • 服务端 API 接收推送数据后应该走 IngestService,需要确认 API 路由是否已经接入

需确认: 查看 app/api/v1/ 里 spiderJobs 推送的接收端点,确认是否走 IngestService。

设计决策:

  • spiderJobs 推送的应该是原始 JSON 数据
  • 接收端点 → IngestService.store_batch() → ClickHouse
  • 来源字段channel/platform在 IngestService 中已经记录

1.3 DATA-03公司清洗定时任务

位置: app/services/company_cleaner.py335行

现状(基本完整):

全链路已实现:

  1. collect_pending_companies(): ClickHouse 查 30 天内有招聘的公司 ID
  2. process_pending_companies(): 遍历 MySQL 清洗队列 → 爬公司详情 → 写 MySQL
    • asyncio.to_thread(boss/qcwy/zhilian_service.get_company_detail)
  3. 同步写 company_jobs调用 company_jobs_sync.sync_company_jobs()
  4. cleanup_old_records(): 清理已处理记录

缺口: 日志链路需确认 logger.info 是否结构化(已使用 loguru

结论: DATA-03 已基本完成,无需大改。


1.4 DATA-04公司招聘信息写入 ClickHouse

位置: app/services/company_jobs_sync.py

现状(部分完成):

# 当前用 "mini" channel 和 "job" data_type 存公司职位
store_result = await router.store_batch(source, "mini", "job", jobs)

缺口: 公司职位company_jobs和普通搜索职位search_jobs使用同一个注册配置 无法区分来源,且 ClickHouse 表boss_job/qcwy_job/zhilian_job混入了不同来源。

方案:

  • 新增 channel = "company"(区分公司关联职位 vs 搜索职位)
  • ingest/configs/boss.py 等添加 channel="company" 的额外配置
  • 或者保持现有 channel="mini",但 ClickHouse 表层面通过 channel 列区分

决策(保守方案):

  • 现有表已有 channel 列,company_jobs_sync.py 改用 channel="company" 调用
  • 无需新建表,无 ClickHouse DDL 变更
  • 只需在 registry 补注册 channel="company" 的配置条目

2. 需确认spiderJobs 推送接收端点

需要查看 app/api/v1/ 是否有专门接收 spiderJobs 推送数据的端点, 以及该端点是否已经调用了 IngestService。


3. 计划分解

  • Plan 01DATA-01 修改 dedup.py 的两个查询 SQL 加 30 天窗口条件(+新增 mock 测试)
  • Plan 02DATA-02 + DATA-04 确认/修复推送 API 端点 + registry 补注册 channel="company" + company_jobs_sync 改用 channel="company"
  • DATA-03 已完成,无需单独计划)