diff --git a/backend/crawler_copy/jd_pc_search/jd_competitor_report.py b/backend/crawler_copy/jd_pc_search/jd_competitor_report.py index 3d53292..323df0f 100644 --- a/backend/crawler_copy/jd_pc_search/jd_competitor_report.py +++ b/backend/crawler_copy/jd_pc_search/jd_competitor_report.py @@ -37,7 +37,12 @@ _ROOT = Path(__file__).resolve().parent if str(_ROOT) not in sys.path: sys.path.insert(0, str(_ROOT)) +_BACKEND_ROOT = Path(__file__).resolve().parents[2] +if str(_BACKEND_ROOT) not in sys.path: + sys.path.insert(0, str(_BACKEND_ROOT)) + import jd_keyword_pipeline as kpl # noqa: E402 +from pipeline.csv_schema import merged_csv_effective_total_sales # noqa: E402 # --------------------------------------------------------------------------- # 运行配置(按需改这里) @@ -1578,7 +1583,8 @@ def _competitor_matrix_md_line( ) cat = _md_cell(_detail_category_path_cell(row), 24) ing = _matrix_ingredients_cell(row) - cc = _md_cell(_cell(row, "销量口径(totalSales)", "评价量(commentFuzzy)"), 14) + ts_eff = merged_csv_effective_total_sales(row) + cc = _md_cell(ts_eff or _cell(row, "评价量(commentFuzzy)"), 14) prev = _md_cell(_cell(row, "comment_preview"), 72) return ( f"| {sku} | {title} | {brand} | {pj} | {df} | {shop} | {sell} | {rank} | " @@ -2665,7 +2671,7 @@ def build_competitor_brief( "category": _detail_category_path_cell(row), "selling_point": _cell(row, "卖点(sellingPoint)")[:240], "comment_fuzzy": _cell(row, "评价量(commentFuzzy)"), - "total_sales": _cell(row, "销量口径(totalSales)"), + "total_sales": merged_csv_effective_total_sales(row), } ) matrix_groups.append( diff --git a/backend/pipeline/backfill_merged_total_sales.py b/backend/pipeline/backfill_merged_total_sales.py deleted file mode 100644 index 9a89b5c..0000000 --- a/backend/pipeline/backfill_merged_total_sales.py +++ /dev/null @@ -1,288 +0,0 @@ -""" -不重新抓搜索:从已有 ``pc_search_raw/*.json``(及 ``.js``)解析 ``totalSales``, -并写回 ``keyword_pipeline_merged.csv`` /可选 ``pc_search_export.csv``。 -若原始 JSON 无该字段,则尝试从「销量楼层(commentSalesFloor)」单元格中抽取「已售…」片段。 - -用法(在 backend 目录下):: - - python pipeline/backfill_merged_total_sales.py --run-dir "../data/JD/pipeline_runs/某批次" - python pipeline/backfill_merged_total_sales.py --merged "D:/path/keyword_pipeline_merged.csv" --dry-run - -""" -from __future__ import annotations - -import argparse -import csv -import json -import re -import sys -from pathlib import Path -from typing import Any - -BACKEND_ROOT = Path(__file__).resolve().parent.parent -if str(BACKEND_ROOT) not in sys.path: - sys.path.insert(0, str(BACKEND_ROOT)) - -from pipeline.csv_schema import ( # noqa: E402 - JD_SEARCH_CSV_HEADERS, - MERGED_FIELD_TO_CSV_HEADER, -) - -COL_MERGED_TOTAL = MERGED_FIELD_TO_CSV_HEADER["total_sales"] -COL_MERGED_FLOOR = MERGED_FIELD_TO_CSV_HEADER["comment_sales_floor"] -COL_SKU_MERGED = MERGED_FIELD_TO_CSV_HEADER["sku_id"] - -COL_EXPORT_TOTAL = JD_SEARCH_CSV_HEADERS["total_sales"] -COL_EXPORT_FLOOR = JD_SEARCH_CSV_HEADERS["comment_sales_floor"] -COL_SKU_EXPORT = JD_SEARCH_CSV_HEADERS["sku_id"] - -# 仅列表响应文件,避免误扫 ``pc_request_*.json`` 请求元数据。 -_RAW_GLOBS = ("pc_search_*.json", "pc_search_*.js") - - -def infer_total_sales_from_sales_floor(cell: str) -> str: - """ - 从「销量楼层」合并列文案中截取可作销量口径的片段(供图表解析件数)。 - 例:``good:99%好评 | 已售50万+`` → ``已售50万+``。 - """ - t = (cell or "").strip() - if not t: - return "" - m = re.search(r"已售\s*[\d,,.+]*\s*[万亿]?\s*\+?", t) - if m: - return m.group(0).strip() - m2 = re.search(r"已售\s*[\d,,.+\s万千亿]+", t) - return m2.group(0).strip() if m2 else "" - - -def _load_json_payload(path: Path) -> Any | None: - try: - text = path.read_text(encoding="utf-8") - except OSError: - return None - try: - return json.loads(text) - except json.JSONDecodeError: - pass - jcr_root = BACKEND_ROOT / "crawler_copy" / "jd_pc_search" - if str(jcr_root) not in sys.path: - sys.path.insert(0, str(jcr_root)) - try: - from search.jd_h5_search_requests import _loads_json_or_jsonp # noqa: WPS433 - - return _loads_json_or_jsonp(text) - except Exception: - return None - - -def collect_total_sales_from_pc_search_raw(raw_dir: Path) -> dict[str, str]: - """ - 遍历 ``pc_search_raw``下保存的列表响应,按 SKU 汇总 ``total_sales``(后者覆盖前者)。 - """ - if str(BACKEND_ROOT / "crawler_copy" / "jd_pc_search") not in sys.path: - sys.path.insert(0, str(BACKEND_ROOT / "crawler_copy" / "jd_pc_search")) - from search.jd_h5_search_requests import parse_items_from_jd_json_payload # noqa: WPS433 - - out: dict[str, str] = {} - seen_paths: set[Path] = set() - for pattern in _RAW_GLOBS: - for p in sorted(raw_dir.glob(pattern)): - if p in seen_paths: - continue - seen_paths.add(p) - payload = _load_json_payload(p) - if payload is None: - continue - try: - rows = parse_items_from_jd_json_payload( - payload, keyword="", page=1 - ) - except Exception: - continue - for r in rows: - if not isinstance(r, dict): - continue - sku = str(r.get("sku_id") or "").strip() - ts = str(r.get("total_sales") or "").strip() - if sku and ts: - out[sku] = ts - return out - - -def _insert_column_after(fieldnames: list[str], col: str, after: str) -> list[str]: - fn = list(fieldnames) - if col in fn: - return fn - if after in fn: - i = fn.index(after) + 1 - fn.insert(i, col) - return fn - # 极旧表:无销量楼层时插在评价量后 - fallback_after = MERGED_FIELD_TO_CSV_HEADER["comment_fuzzy"] - if fallback_after in fn: - fn.insert(fn.index(fallback_after) + 1, col) - return fn - fn.append(col) - return fn - - -def _backfill_rows( - rows: list[dict[str, str]], - sku_col: str, - total_col: str, - floor_col: str, - sku_to_total: dict[str, str], - *, - use_floor_fallback: bool, -) -> int: - n = 0 - for row in rows: - cur = str(row.get(total_col) or "").strip() - sku = str(row.get(sku_col) or "").strip() - if not cur and sku: - ts = sku_to_total.get(sku, "") - if ts: - row[total_col] = ts - cur = ts - n += 1 - if not cur and use_floor_fallback: - floor = str(row.get(floor_col) or "").strip() - inf = infer_total_sales_from_sales_floor(floor) - if inf: - row[total_col] = inf - n += 1 - return n - - -def backfill_csv_file( - path: Path, - *, - sku_to_total: dict[str, str], - is_merged: bool, - use_floor_fallback: bool, - dry_run: bool, -) -> tuple[int, list[str]]: - raw = path.read_text(encoding="utf-8-sig") - lines = raw.splitlines() - if not lines: - return 0, [] - reader = csv.DictReader(lines) - old_fn = reader.fieldnames or [] - if is_merged: - sku_c, tot_c, fl_c = ( - COL_SKU_MERGED, - COL_MERGED_TOTAL, - COL_MERGED_FLOOR, - ) - fieldnames = _insert_column_after(list(old_fn), tot_c, COL_MERGED_FLOOR) - else: - sku_c, tot_c, fl_c = COL_SKU_EXPORT, COL_EXPORT_TOTAL, COL_EXPORT_FLOOR - fieldnames = _insert_column_after(list(old_fn), tot_c, COL_EXPORT_FLOOR) - rows = list(reader) - for r in rows: - for h in fieldnames: - r.setdefault(h, "") - filled = _backfill_rows( - rows, sku_c, tot_c, fl_c, sku_to_total, use_floor_fallback=use_floor_fallback - ) - if dry_run: - return filled, fieldnames - from io import StringIO - - sio = StringIO() - w2 = csv.DictWriter(sio, fieldnames=fieldnames, lineterminator="\n") - w2.writeheader() - w2.writerows(rows) - path.write_text("\ufeff" + sio.getvalue(), encoding="utf-8") - return filled, fieldnames - - -def _resolve_raw_dir( - run_dir: Path | None, merged_path: Path | None -) -> Path | None: - if run_dir is not None: - rd = (run_dir / "pc_search_raw").resolve() - if rd.is_dir(): - return rd - if merged_path is not None: - rd = (merged_path.parent / "pc_search_raw").resolve() - if rd.is_dir(): - return rd - return None - - -def main() -> None: - ap = argparse.ArgumentParser(description="从 pc_search_raw 补全销量口径列(不重新请求搜索)") - ap.add_argument("--run-dir", type=Path, default=None, help="批次目录(含 keyword_pipeline_merged.csv 与 pc_search_raw)") - ap.add_argument("--merged", type=Path, default=None, help="合并表路径(可单独指定)") - ap.add_argument( - "--raw-dir", - type=Path, - default=None, - help="原始搜索响应目录(默认:run-dir 或 merged 父目录下的 pc_search_raw)", - ) - ap.add_argument( - "--also-pc-search-export", - action="store_true", - help="同时处理同目录下的 pc_search_export.csv", - ) - ap.add_argument( - "--no-floor-fallback", - action="store_true", - help="禁用从销量楼层文案推断(仅用原始 JSON 中的 totalSales)", - ) - ap.add_argument("--dry-run", action="store_true", help="只统计将补全条数,不写文件") - args = ap.parse_args() - - run_dir = args.run_dir.resolve() if args.run_dir else None - merged_path = args.merged - if merged_path is None and run_dir is not None: - merged_path = run_dir / "keyword_pipeline_merged.csv" - if merged_path is None or not merged_path.is_file(): - ap.error("请指定有效的 --merged 或含 keyword_pipeline_merged.csv 的 --run-dir") - - merged_path = merged_path.resolve() - raw_dir = args.raw_dir.resolve() if args.raw_dir else _resolve_raw_dir(run_dir, merged_path) - sku_map: dict[str, str] = {} - if raw_dir is not None: - sku_map = collect_total_sales_from_pc_search_raw(raw_dir) - print(f"[backfill] 自 {raw_dir} 解析到带 totalSales 的 SKU:{len(sku_map)}", file=sys.stderr) - else: - print( - "[backfill] 未找到 pc_search_raw,将仅尝试销量楼层推断(若未加 --no-floor-fallback)", - file=sys.stderr, - ) - - use_floor = not args.no_floor_fallback - n_m, _ = backfill_csv_file( - merged_path, - sku_to_total=sku_map, - is_merged=True, - use_floor_fallback=use_floor, - dry_run=args.dry_run, - ) - print( - f"[backfill] merged:补全单元格数 {n_m}(空列→有值;dry_run={args.dry_run})", - file=sys.stderr, - ) - - if args.also_pc_search_export: - exp = merged_path.parent / "pc_search_export.csv" - if exp.is_file(): - n_e, _ = backfill_csv_file( - exp, - sku_to_total=sku_map, - is_merged=False, - use_floor_fallback=use_floor, - dry_run=args.dry_run, - ) - print( - f"[backfill] pc_search_export:补全单元格数 {n_e}", - file=sys.stderr, - ) - else: - print(f"[backfill] 跳过:无 {exp}", file=sys.stderr) - - -if __name__ == "__main__": - main() diff --git a/backend/pipeline/csv_schema.py b/backend/pipeline/csv_schema.py index ff1e6b4..0d90bfd 100644 --- a/backend/pipeline/csv_schema.py +++ b/backend/pipeline/csv_schema.py @@ -4,6 +4,8 @@ """ from __future__ import annotations +import re + # --- 搜索导出 pc_search_export.csv(列名为中文,与 jd_h5_search_requests.JD_EXPORT_COLUMN_HEADERS 一致)--- JD_SEARCH_INTERNAL_KEYS: tuple[str, ...] = ( "item_id", @@ -183,3 +185,27 @@ MERGED_CSV_TO_FIELD: dict[str, str] = dict(zip(MERGED_CSV_COLUMNS, MERGED_INTERN MERGED_FIELD_TO_CSV_HEADER: dict[str, str] = { internal: csv_h for csv_h, internal in MERGED_CSV_TO_FIELD.items() } + + +def infer_total_sales_from_sales_floor(cell: str) -> str: + """ + 从「销量楼层(commentSalesFloor)」列文案截取可作 ``销量口径(totalSales)`` 的片段(与列表接口未单独落 totalSales 列时的兜底一致)。 + """ + t = (cell or "").strip() + if not t: + return "" + m = re.search(r"已售\s*[\d,,.+]*\s*[万亿]?\s*\+?", t) + if m: + return m.group(0).strip() + m2 = re.search(r"已售\s*[\d,,.+\s万千亿]+", t) + return m2.group(0).strip() if m2 else "" + + +def merged_csv_effective_total_sales(row: dict[str, str]) -> str: + """合并表一行:优先已有 ``销量口径(totalSales)``,否则从销量楼层推断。""" + h_ts = MERGED_FIELD_TO_CSV_HEADER["total_sales"] + h_fl = MERGED_FIELD_TO_CSV_HEADER["comment_sales_floor"] + direct = str(row.get(h_ts) or "").strip() + if direct: + return direct + return infer_total_sales_from_sales_floor(str(row.get(h_fl) or "")) diff --git a/backend/pipeline/ingest.py b/backend/pipeline/ingest.py index 3de93cd..5549552 100644 --- a/backend/pipeline/ingest.py +++ b/backend/pipeline/ingest.py @@ -22,7 +22,9 @@ from .csv_schema import ( JD_SEARCH_INTERNAL_KEYS, MERGED_CSV_COLUMNS, MERGED_CSV_TO_FIELD, + MERGED_FIELD_TO_CSV_HEADER, SEARCH_CSV_HEADER_TO_FIELD, + merged_csv_effective_total_sales, ) from .models import ( JdJobCommentRow, @@ -84,6 +86,12 @@ def _comment_row_kwargs(row: dict[str, str]) -> dict[str, str]: } +def _normalize_merged_csv_total_sales(row: dict[str, str]) -> None: + """列表未写 totalSales 列时,用销量楼层推断,保证入库与快照与报告口径一致。""" + h = MERGED_FIELD_TO_CSV_HEADER["total_sales"] + row[h] = merged_csv_effective_total_sales(row) + + def _merged_row_kwargs(row: dict[str, str]) -> dict[str, str]: return { MERGED_CSV_TO_FIELD[col]: str(row.get(col) or "").strip() for col in MERGED_CSV_COLUMNS @@ -152,6 +160,7 @@ def ingest_job_dataset_rows(job: PipelineJob) -> dict[str, Any]: merged_rows = _read_csv_rows(merged_path) if merged_path.is_file() else [] m_objs: list[JdJobMergedRow] = [] for i, row in enumerate(merged_rows): + _normalize_merged_csv_total_sales(row) kw = _merged_row_kwargs(row) m_objs.append(JdJobMergedRow(job=job, row_index=i, **kw)) _bulk_create_in_chunks(JdJobMergedRow, m_objs) @@ -185,6 +194,7 @@ def ingest_job_merged_csv(job: PipelineJob) -> dict[str, Any]: sku = (row.get(SKU_FIELD_MERGED) or "").strip() if not sku: continue + _normalize_merged_csv_total_sales(row) payload = _payload_as_json(row) title = (row.get(TITLE_FIELD) or "")[:2000] ware = (row.get(WARE_FIELD) or "").strip()[:64] diff --git a/backend/pipeline/management/commands/refresh_jd_merged_total_sales.py b/backend/pipeline/management/commands/refresh_jd_merged_total_sales.py new file mode 100644 index 0000000..77b9858 --- /dev/null +++ b/backend/pipeline/management/commands/refresh_jd_merged_total_sales.py @@ -0,0 +1,54 @@ +# -*- coding: utf-8 -*- +""" +合并表入库后若 ``total_sales`` 为空,可按与入库相同的规则从 ``comment_sales_floor`` 补全。 + + python manage.py refresh_jd_merged_total_sales + python manage.py refresh_jd_merged_total_sales --job-id 42 +""" +from __future__ import annotations + +from django.core.management.base import BaseCommand + +from pipeline.csv_schema import MERGED_FIELD_TO_CSV_HEADER, merged_csv_effective_total_sales +from pipeline.models import JdJobMergedRow + + +class Command(BaseCommand): + help = "从销量楼层推断并回填 JdJobMergedRow.total_sales(与 ingest 口径一致)。" + + def add_arguments(self, parser) -> None: + parser.add_argument( + "--job-id", + type=int, + default=None, + help="仅处理该 PipelineJob;默认处理全部任务下的合并行", + ) + + def handle(self, *args, **options) -> None: + job_id = options.get("job_id") + qs = JdJobMergedRow.objects.all().order_by("id") + if job_id is not None: + qs = qs.filter(job_id=job_id) + + h_ts = MERGED_FIELD_TO_CSV_HEADER["total_sales"] + h_fl = MERGED_FIELD_TO_CSV_HEADER["comment_sales_floor"] + updates: list[JdJobMergedRow] = [] + n_changed = 0 + for r in qs.iterator(chunk_size=800): + row = {h_ts: r.total_sales or "", h_fl: r.comment_sales_floor or ""} + eff = merged_csv_effective_total_sales(row) + if eff and eff != (r.total_sales or "").strip(): + r.total_sales = eff + updates.append(r) + n_changed += 1 + if len(updates) >= 500: + JdJobMergedRow.objects.bulk_update(updates, ["total_sales"]) + updates.clear() + if updates: + JdJobMergedRow.objects.bulk_update(updates, ["total_sales"]) + + self.stdout.write( + self.style.SUCCESS( + f"refresh_jd_merged_total_sales 完成,更新行数约 {n_changed}" + ) + ) diff --git a/backend/pipeline/tests/test_competitor_brief.py b/backend/pipeline/tests/test_competitor_brief.py index 7f1eabd..400896e 100644 --- a/backend/pipeline/tests/test_competitor_brief.py +++ b/backend/pipeline/tests/test_competitor_brief.py @@ -8,7 +8,7 @@ from pathlib import Path from django.conf import settings from django.test import SimpleTestCase -from pipeline.backfill_merged_total_sales import infer_total_sales_from_sales_floor +from pipeline.csv_schema import infer_total_sales_from_sales_floor from pipeline.report_charts import _cn_volume_int