Files
InboundVerify/docs/2026-07-29-应到提交导出去重实施计划.md
Misaka_Company f710622a3d feat: dedup expected data by handover_no before export submit
- store.get_existing_handover_nos: query PG expected_record.handover_no

- inject dedup skip before submitting export in zto/yunda/anneng/shunxin

- shunxin reads RTS handover_no from waybill-list view (method 1)

- force-redownload switch threaded via task_spec -> dispatch -> impl

- schema: add idx_expected_handover index

Co-Authored-By: Claude <noreply@anthropic.com>
2026-07-29 10:57:07 +08:00

18 KiB
Raw Blame History

应到数据「提交导出任务前」去重 实施计划 v2

执行约定: 本仓库无 pytest验证靠实跑站点流程 + DB 核对(见第 7 节)。 本仓库约定 不自动提交;所有改动落地后等用户明确说"提交"再 commit/push。 本次会话额外约定:未获用户明确指示前不动代码、不提交。

Goal 在周期性自动落库场景下,于「提交导出任务」之前,按交接单号(顺心=运单列表界面里的交接单号)判断该批应到数据是否已落库,已落库则跳过提交导出任务,从源头消除重复下载与重复落库;并提供一个"强制重下"开关(默认关)兜底。

Architecture

  • 去重数据源 = PostgreSQL expected_record.handover_no(已存在字段,权威);store.py 新增 get_existing_handover_nos(site)(含 cpolar 降级)。
  • 在 4 站应到下载循环内、提交导出动作之前注入"命中已落库则 continue"。
  • force 开关经 task_specdispatch_task → handler → 各站 download(page, force)impl(page, force) 透传;force=True 时跳过去重。周期 job 默认不 force。

Tech Stack Python 3.10+ / Playwright网页 3 站)/ 裸 CDP安能/ psycopg / SQLite前端 Next.js 16 + React 19。

Global Constraints

  • Python 一律 .venv;改完任何 .py 必须跑 .venv/Scripts/python.exe -m black inbound_verify
  • 改完跑 .venv/Scripts/python.exe -m py_compile inbound_verify 自检。
  • 不自动提交/推送(覆盖全局 auto-push 默认)。
  • 不动 export_times 时间容差≤40s匹配机制CLAUDE.md 约定)。
  • 安能 CDP 驱动,绝不用 Page.reload
  • config.yaml 已 gitignore不提交真实凭据。
  • 前端是 Next.js 16有 breaking changes,写前端代码前先查 node_modules/next/dist/docs/
  • 行号基于 2026-07-29 快照,实现时以当前代码为准、就近定位。

1. 已定决策v1 审核反馈)

决策 结论
A. 去重数据源 查 PostgreSQL expected_record.handover_no
B. 范围 本次只做应到expected;实到/百世不动
C. 强制开关 加 force 开关,默认不强制重下
顺心标识 方式1点"运单列表"后、点导出前,从运单列表界面读交接单号(与其他 3 站统一用交接单号去重)

2. 背景与问题根源

周期链路:fetch_schedule(IntervalTrigger) → task_queue → worker → dispatch_task → handler → 提交导出+下载 → _persist_to_db(UPSERT)

  • DB 已幂等(expected_record(site, waybill_no) UPSERT
  • dispatch_task 调 handler 前无"是否需要下载"判断,周期触发重复"提交导出→下载→解析"。
  • 本方案在「提交导出任务」前按交接单号去重,从源头省掉重复下载。

3. 各站探索结论(注入点)

站点 文件 提交导出位置 提交前标识 来源
中通 sites/zto.py zto_expected_download_impl L264 已有 handover_noL251 主表行 td.nth(3) 正则 18 位
韵达 sites/yunda.py yunda_expected_download_impl L342 已有 raw_noL291 列表行 td.nth(1)
安能 sites/anneng.py 主循环 L927逐条 已有 ewbs_noL907-913 CDP 复选框 ewbsListNo= 正则 19 位
顺心 sites/shunxin.py L319 点导出 🆕 方式1L316 后读运单列表界面交接单号 DOM 待实勘

顺心方式1关键事实已探明raw 里「班次号」「交接单号」都有,且一个交接单 = 一个班次 = 多条运单;交接单号即入库 handover_no,与另 3 站同键。

4. 文件结构

文件 改动
schema.sql expected_record(site, handover_no) 索引
inbound_verify/store.py 新增 get_existing_handover_nos(site)
inbound_verify/cli/server.py TaskRequestforcecreate_task 透传 force周期不 force
inbound_verify/runtime.py dispatch_task 读 force 传 handler所有 handler 加 force 形参并透传到 download_func
inbound_verify/sites/zto.py download/implforce;应到循环注入去重 + 空兜底
inbound_verify/sites/yunda.py 同上(空兜底已存在)
inbound_verify/sites/anneng.py 同上
inbound_verify/sites/shunxin.py download/implforce方式1点运单列表后读交接单号去重 + 退回 + 空兜底
dashboard/app/page.tsx forceRedownload state + checkboxtrigger/triggerPrimary 透传 force

4 站的 actual实到download 入口也统一加 force=False 形参(接收但不用,仅让 _web_handler 的统一调用成立actual impl 不改。

5. 任务分解

Task 1schema.sql 加索引

Files: Modify schema.sqlidx_expected_site_date 之后)

CREATE INDEX IF NOT EXISTS idx_expected_handover ON expected_record (site, handover_no);

验证: .venv/Scripts/python.exe -m inbound_verify.store init(幂等)。


Task 2store.py 新增查已落库交接单号集合

Files: Modify inbound_verify/store.pyingest_task 之后)

Produces: get_existing_handover_nos(site: str) -> set[str]

def get_existing_handover_nos(site):
    """查该站点已落库的交接单号集合expected_record.handover_no    供"提交导出任务前"去重:已落库的不再重复提交导出。
    PG 不可用cpolar 抖动等)时返回空集 + 告警,调用方按"未确认存在"处理
    继续提交导出UPSERT 兜底,绝不因去重查询失败而漏数据)。"""
    try:
        with _connect(_load_pg_config()["dbname"]) as conn:
            with conn.cursor() as cur:
                cur.execute(
                    "SELECT handover_no FROM expected_record "
                    "WHERE site=%s AND handover_no IS NOT NULL AND handover_no <> ''",
                    (site,),
                )
                return {str(r[0]).strip() for r in cur.fetchall()}
    except Exception as e:
        print(f">> [去重] 查询已落库交接单号失败({site}),本次不去重: {e}")
        return set()

验证: .venv/Scripts/python.exe -c "from inbound_verify import store; print(len(store.get_existing_handover_nos('中通')))" 不抛异常即过。


Task 3force 开关后端骨架server + runtime

forcetask_spec 一路透传到各站 download(page, force)with_retryflow 是零参 lambdaforce 经闭包捕获,with_retry 不动

(a) server.py TaskRequest + create_task

class TaskRequest(BaseModel):
    site: str
    kind: str
    force: bool = False          # 新增:强制重下(忽略已落库去重),默认关
# create_task 内
task_queue.put((task_id, {"site": req.site, "kind": req.kind, "force": req.force}))

_enqueue_fetchL133 周期投递)保持不变(不带 force → 默认 False

(b) runtime.py dispatch_task 透传 force

def dispatch_task(ctx, task_spec):
    site = task_spec.get("site")
    kind = task_spec.get("kind")
    force = bool(task_spec.get("force", False))     # 新增
    ...
    try:
        ret = handler(ctx, force)                    # 改:原 handler(ctx)

(c) runtime.py 所有 handler 加 force 形参:

_web_handler

    def handler(ctx, force=False):
        pg = ctx.pages_map[site]
        if isinstance(pg, list):                                     # 顺心双账号
            return download_func(pg, foreground=ctx.foreground, force=force)
        if ctx.foreground:
            pg.bring_to_front()
        return download_func(pg, force=force)

_site_undelivered_handler

    def handler(ctx, force=False):
        exp_ok = TASK_HANDLERS[(site, "expected")](ctx, force) is not False
        act_ok = (TASK_HANDLERS[(site, "actual")](ctx, force) is not False) if exp_ok else False
        ...

安能 expected/actual 与 compare签名兼容即可

("安能", "expected"): lambda ctx, force=False: anneng.anneng_expected_download(force=force),
("安能", "actual"):   lambda ctx, force=False: anneng.anneng_actual_download(force=force),
("__compare__", "compare"): lambda ctx, force=False: (compare.main() or True),

验证: py_compile 通过;服务重启后 POST /tasks {site,kind,force:true} 不报 TypeError此时各站 download 的 force 形参由 Task 4-7 补齐,连续实施)。


Task 4中通 zto.pydownload/impl 加 force + 去重)

Files: Modify inbound_verify/sites/zto.py

(a) 入口与 impl 加 force闭包透传with_retry 不动):

def zto_expected_download(page, force=False):
    return with_retry(
        "中通", "应到",
        lambda: zto_expected_download_impl(page, force=force),
        lambda: zto_reset(page),
    )

def zto_expected_download_impl(page, force=False):
    ...
# zto_actual_download / zto_actual_download_impl 同样加 force=False 形参actual 不用 force仅兼容

(b) 循环前加载已落库集合L241 print 之后、L243 for 之前):

        # 【去重】加载本站已落库交接单号force=True 或查询失败时 existing=空集(不去重)
        if force:
            existing = set()
            print(">> [去重] 强制重下,跳过去重。")
        else:
            try:
                from inbound_verify import store
                existing = store.get_existing_handover_nos("中通")
            except Exception as _e:
                existing = set()
                print(f">> [去重] 加载失败,本次不去重: {_e}")

(c) 循环内命中跳过L252 print 之后、L254 row.dblclick() 之前):

            print(f"      -> 当前交接单号:{handover_no}")
            if handover_no in existing:
                print(f"      ⏭️ 交接单号 {handover_no} 已落库,跳过提交导出。")
                continue
            row.dblclick()

(d) 空列表兜底L305 循环后、进入 _zto_poll_and_download_tasks 之前):

        if not export_times:
            print(">> 本次无新交接单需导出(全部已落库或无数据),结束。")
            return

Task 5韵达 yunda.py同构

Files: Modify inbound_verify/sites/yunda.py

(a) yunda_expected_download(page, force=False) + impl 加 forceactual 同理加形参);(b) 循环前加载 existing同 Task 4b站点"韵达"

(c) 循环内命中跳过L291 raw_no = ... 之后、L293 # 跳过已绑定的交接单 之前):

            raw_no = current_row.locator("td").nth(1).inner_text().strip()
            if raw_no in existing:
                print(f"      ⏭️ 交接单号 {raw_no} 已落库,跳过提交导出。")
                continue
            # 跳过已绑定的交接单
            bind_status = current_row.locator("td").nth(2).inner_text().strip()

(d) 空列表兜底 已存在L389-392无需新增。


Task 6安能 anneng.pyCDP

Files: Modify inbound_verify/sites/anneng.py

(a) anneng_expected_download(force=False) + impl 加 forceactual 同理);

(b) 主流程加载 existinganneng_expected_download_impl 内、L907 收集 target_ids 前,与 export_times = []L882并列同 Task 4b站点"安能"

(c) 主循环命中跳过L919 for 内、L920 print 之后、L921 activate_tab 之前):

    for i, ewbs_no in enumerate(target_ids, start=1):
        print(f"   ⏳ [{i}/{len(target_ids)}] 交接单号 {ewbs_no}")
        if ewbs_no in existing:
            print(f"      ⏭️ 交接单号 {ewbs_no} 已落库,跳过。")
            continue
        activate_tab(tab_cdp, "交接单信息")
        ...

(d) 空列表兜底(主循环后、进入 poll_and_download_tasks 之前):

    if not export_times:
        print(">> 本次无新交接单需导出(全部已落库或无数据),结束。")
        return

Task 7顺心 shunxin.py方式1 + 实勘)

顺心流程L315 点"运单列表" → L316 等"运单查询"label运单列表界面 → L319 点"导出"。方式1 在 L316 之后、L319 之前读交接单号。

Step 1实勘不改业务逻辑 确认运单列表界面里交接单号的 DOM 选择器(哪个元素/列)。实勘方式二选一(待用户同意):

  • 方式 A临时在 L316 后加调试打印dump 运单列表界面关键 DOM 文本),跑一次顺心应到,从日志定位选择器,再删调试代码。
  • 方式 Bdebug 模式(config.yaml debug.target_site=顺心CDP 9223单独挂载用 Playwright CLI 观察。

Step 2注入选择器 <HANDOVER_SELECTOR> 确认后替换):

(a) 入口与 impl 加 force

def shunxin_expected_download(pages, foreground=True, force=False):
    return with_retry(
        "顺心", "应到",
        lambda: shunxin_expected_download_impl(pages, foreground=foreground, force=force),
        lambda: shunxin_reset(pages),
    )

def shunxin_expected_download_impl(pages, foreground=True, force=False):
    ...
# shunxin_actual_download / impl 同样加 force=False 形参actual 不用)

(b) 循环前加载 existing同 Task 4b站点"顺心";两账号共享同一 existing)。

(c) 循环内:点运单列表 → 读交接单号 → 命中则退回跳过L315-316 之后、L319 点导出之前):

        for i in range(count):
            print(f"   ⏳ 正在处理第 {i+1}/{count} 个班次...")

            waybill_btns.nth(i).click()
            page.locator("label[title='运单查询']").wait_for(state="visible")

            # 【方式1】运单列表界面已加载读交接单号 → 已落库则退回列表跳过
            handover_no = page.locator("<HANDOVER_SELECTOR>").first.inner_text().strip()
            if handover_no in existing:
                print(f"      ⏭️ 交接单号 {handover_no} 已落库,跳过提交导出。")
                page.get_by_role("tab", name="车辆点到").click()   # 退回列表(复用 L336
                page.wait_for_timeout(500)
                continue

            # 4. 执行导出流程
            page.get_by_role("button", name="export 导出").click()
            ...

(d) 空列表兜底L339 循环后、进入导出任务管理页轮询之前):

        if not export_times:
            print(">> 本次无新班次需导出(全部已落库或无数据),结束。")
            return

Task 8前端 force 开关dashboard

Files: Modify dashboard/app/page.tsxBFF app/api/tasks/route.ts 是 generic 透传,不用改

(a) 加 stateL21 附近):

const [forceRedownload, setForceRedownload] = useState(false);

(b) trigger 加 force 形参并写入 bodyL27-50

  const trigger = useCallback(
    async (site: string, kind: string, label: string, force?: boolean) => {
      ...
        body: JSON.stringify({ site, kind, force: !!force }),
      ...
    },
    [refreshTasks],
  );

(c) triggerPrimary 透传L58-63

  const triggerPrimary = useCallback(
    async (cfg: SiteConfig) => {
      await trigger(cfg.name, cfg.primaryKind, `${cfg.name}·获取未到`, forceRedownload);
    },
    [trigger, forceRedownload],
  );

(d) UI在配置区/顶部加 checkbox默认不勾

<label className="inline-flex items-center gap-1 text-xs text-amber-700">
  <input
    type="checkbox"
    checked={forceRedownload}
    onChange={(e) => setForceRedownload(e.target.checked)}
  />
  强制重新下载(忽略已落库去重)
</label>

仅"获取未到"主按钮透传 force周期抓取不经前端、恒不 force。

验证: 前端勾选 → 触发 → 网络面板看到 POST /api/tasks body 含 force:true;后端日志 强制重下,跳过去重

6. 关键注意事项(陷阱)

  • 跳过的交接单号绝不 export_times.append:否则 len(export_times) > 实际提交数 → 下载数校验失败。所有 continue 都在 append 之前。
  • 全部跳过时必须 returnexport_times 为空时不进导出任务管理页轮询。
  • 顺心方式1跳过要退回:点进运单列表后命中已落库,需点"车辆点到"tab 退回再 continue(复用 L336
  • PG 降级只防漏不防重:查询失败 = 空集 = 当作未存在 = 继续提交UPSERT 兜底。
  • 不动 export_times 时间容差匹配
  • actual/百世 download 只加 force 形参兼容impl 不加去重。

7. 验证(无 pytest

  1. 单站联调config.yaml debug 单站):首次新单正常下载+入库;再触发同范围 → 已入库的全部 ⏭️ 跳过export_times 空,直接 return。
  2. force 开关:勾选"强制重下" → 已落库的也重新提交导出(日志 强制重下,跳过去重)。
  3. DB 核对SELECT site, handover_no, COUNT(*) FROM expected_record GROUP BY site, handover_no 无翻倍。
  4. cpolar 降级:断 PG → get_existing_handover_nos 返回空集 + 告警,流程仍正常下载(不漏)。
  5. Black + py_compile;前端 npm run build 或 dev 热更无类型错。

8. 实施顺序与依赖

Task 1 → 2 → 3骨架此时各站 download 的 force 形参在 4-7 补)→ 4/5/6/7各站连续做完让链路自洽→ 8前端。顺心 Task 7 的 Step1 实勘需在运行的服务上操作,实施时与用户协调时机。