diff --git a/backend/pipeline/dataset_api.py b/backend/pipeline/dataset_api.py index b4fa63d..511d36b 100644 --- a/backend/pipeline/dataset_api.py +++ b/backend/pipeline/dataset_api.py @@ -1,18 +1,35 @@ -"""库内数据浏览 API:排序、价格与类目筛选(查询参数解析与 QuerySet 变换)。""" +"""库内数据浏览 API:排序、价格与报告细类筛选(查询参数解析与 QuerySet 变换)。""" from __future__ import annotations from typing import Any -from django.db.models import F, QuerySet +from django.db.models import F, Q, QuerySet from django.db.models.expressions import OrderBy from rest_framework.request import Request -SEARCH_SORT_FIELDS = frozenset({"row_index", "price", "sku_id", "title", "leaf_category"}) +SEARCH_SORT_FIELDS = frozenset( + {"row_index", "price", "sku_id", "title", "leaf_category", "matrix_group_label"} +) DETAIL_SORT_FIELDS = frozenset( - {"row_index", "price", "sku_id", "detail_category_path", "detail_brand"} + { + "row_index", + "price", + "sku_id", + "detail_category_path", + "detail_brand", + "matrix_group_label", + } ) MERGED_SORT_FIELDS = frozenset( - {"row_index", "price", "sku_id", "title", "leaf_category", "detail_category_path"} + { + "row_index", + "price", + "sku_id", + "title", + "leaf_category", + "detail_category_path", + "matrix_group_label", + } ) @@ -39,11 +56,9 @@ def price_bounds_from_request(request: Request) -> tuple[float | None, float | N ) -def category_norm_id_from_request(request: Request) -> int | None: - raw = (request.query_params.get("category_norm_id") or "").strip() - if raw.isdigit(): - return int(raw) - return None +def report_group_from_request(request: Request) -> str: + """与 §5 矩阵一致的细类名(如「饼干」「米」);对应查询参数 ``report_group``。""" + return (request.query_params.get("report_group") or "").strip() def detail_category_q_from_request(request: Request) -> str: @@ -52,7 +67,7 @@ def detail_category_q_from_request(request: Request) -> str: def filter_echo( *, - category_norm_id: int | None, + report_group: str, price_min: float | None, price_max: float | None, detail_category_q: str, @@ -60,7 +75,7 @@ def filter_echo( desc: bool, ) -> dict[str, Any]: return { - "category_norm_id": category_norm_id, + "report_group": report_group or None, "price_min": price_min, "price_max": price_max, "detail_category_q": detail_category_q or None, @@ -70,9 +85,9 @@ def filter_echo( def apply_search_filters(qs: QuerySet, request: Request) -> QuerySet: - cid = category_norm_id_from_request(request) - if cid is not None: - qs = qs.filter(leaf_category_norm_id=cid) + rg = report_group_from_request(request) + if rg: + qs = qs.filter(Q(matrix_group_label=rg) | Q(leaf_category=rg)) pmin, pmax = price_bounds_from_request(request) if pmin is not None: qs = qs.filter(price_value__gte=pmin) @@ -94,6 +109,7 @@ def apply_search_order(qs: QuerySet, sort: str, desc: bool) -> QuerySet: "sku_id": "sku_id", "title": "title", "leaf_category": "leaf_category", + "matrix_group_label": "matrix_group_label", }[sort] return qs.order_by( OrderBy(F(field), descending=desc, nulls_last=True), @@ -102,6 +118,9 @@ def apply_search_order(qs: QuerySet, sort: str, desc: bool) -> QuerySet: def apply_detail_filters(qs: QuerySet, request: Request) -> QuerySet: + rg = report_group_from_request(request) + if rg: + qs = qs.filter(matrix_group_label=rg) q = detail_category_q_from_request(request) if q: qs = qs.filter(detail_category_path__icontains=q) @@ -126,6 +145,7 @@ def apply_detail_order(qs: QuerySet, sort: str, desc: bool) -> QuerySet: "sku_id": "sku_id", "detail_category_path": "detail_category_path", "detail_brand": "detail_brand", + "matrix_group_label": "matrix_group_label", }[sort] return qs.order_by( OrderBy(F(field), descending=desc, nulls_last=True), @@ -134,9 +154,9 @@ def apply_detail_order(qs: QuerySet, sort: str, desc: bool) -> QuerySet: def apply_merged_filters(qs: QuerySet, request: Request) -> QuerySet: - cid = category_norm_id_from_request(request) - if cid is not None: - qs = qs.filter(leaf_category_norm_id=cid) + rg = report_group_from_request(request) + if rg: + qs = qs.filter(matrix_group_label=rg) q = detail_category_q_from_request(request) if q: qs = qs.filter(detail_category_path__icontains=q) @@ -162,6 +182,7 @@ def apply_merged_order(qs: QuerySet, sort: str, desc: bool) -> QuerySet: "title": "title", "leaf_category": "leaf_category", "detail_category_path": "detail_category_path", + "matrix_group_label": "matrix_group_label", }[sort] return qs.order_by( OrderBy(F(field), descending=desc, nulls_last=True), diff --git a/backend/pipeline/dataset_nonempty.py b/backend/pipeline/dataset_nonempty.py index 844c8dc..f4405a0 100644 --- a/backend/pipeline/dataset_nonempty.py +++ b/backend/pipeline/dataset_nonempty.py @@ -14,6 +14,8 @@ from .csv_schema import ( from .models import JdJobCommentRow, JdJobDetailRow, JdJobMergedRow, JdJobSearchRow, PipelineJob from .row_serialize import COMMENT_FIELDS_ORDER, DETAIL_FIELDS_ORDER +MATRIX_GROUP_COLUMN = {"key": "matrix_group_label", "label": "报告细类(§5矩阵)"} + def _is_nonempty(val) -> bool: if val is None: @@ -62,7 +64,13 @@ def nonempty_merged_fields_for_job(job: PipelineJob) -> list[str]: def search_columns_for_api(job: PipelineJob) -> list[dict[str, str]]: - return [{"key": k, "label": JD_SEARCH_CSV_HEADERS[k]} for k in nonempty_search_keys_for_job(job)] + cols = [ + {"key": k, "label": JD_SEARCH_CSV_HEADERS[k]} + for k in nonempty_search_keys_for_job(job) + ] + if JdJobSearchRow.objects.filter(job=job).exclude(matrix_group_label="").exists(): + cols.append(dict(MATRIX_GROUP_COLUMN)) + return cols def _detail_field_to_csv_col(field: str) -> str: @@ -73,10 +81,13 @@ def _detail_field_to_csv_col(field: str) -> str: def detail_columns_for_api(job: PipelineJob) -> list[dict[str, str]]: - return [ + cols = [ {"key": f, "label": _detail_field_to_csv_col(f)} for f in nonempty_detail_fields_for_job(job) ] + if JdJobDetailRow.objects.filter(job=job).exclude(matrix_group_label="").exists(): + cols.append(dict(MATRIX_GROUP_COLUMN)) + return cols def _comment_field_to_csv_col(field: str) -> str: @@ -94,21 +105,30 @@ def comment_columns_for_api(job: PipelineJob) -> list[dict[str, str]]: def merged_columns_for_api(job: PipelineJob) -> list[dict[str, str]]: - return [ + cols = [ {"key": k, "label": MERGED_FIELD_TO_CSV_HEADER[k]} for k in nonempty_merged_fields_for_job(job) ] + if JdJobMergedRow.objects.filter(job=job).exclude(matrix_group_label="").exists(): + cols.append(dict(MATRIX_GROUP_COLUMN)) + return cols def search_export_headers(job: PipelineJob) -> list[str]: keys = nonempty_search_keys_for_job(job) - return ["id", "row_index"] + [JD_SEARCH_CSV_HEADERS[k] for k in keys] + h = ["id", "row_index"] + [JD_SEARCH_CSV_HEADERS[k] for k in keys] + if JdJobSearchRow.objects.filter(job=job).exclude(matrix_group_label="").exists(): + h.append(MATRIX_GROUP_COLUMN["label"]) + return h def detail_export_headers(job: PipelineJob) -> list[str]: fields = set(nonempty_detail_fields_for_job(job)) cols = [c for c in DETAIL_CSV_COLUMNS if DETAIL_CSV_TO_FIELD[c] in fields] - return ["id", "row_index"] + cols + h = ["id", "row_index"] + cols + if JdJobDetailRow.objects.filter(job=job).exclude(matrix_group_label="").exists(): + h.append(MATRIX_GROUP_COLUMN["label"]) + return h def comment_export_headers(job: PipelineJob) -> list[str]: @@ -119,4 +139,7 @@ def comment_export_headers(job: PipelineJob) -> list[str]: def merged_export_headers(job: PipelineJob) -> list[str]: keys = nonempty_merged_fields_for_job(job) - return ["id", "row_index"] + [MERGED_FIELD_TO_CSV_HEADER[k] for k in keys] + h = ["id", "row_index"] + [MERGED_FIELD_TO_CSV_HEADER[k] for k in keys] + if JdJobMergedRow.objects.filter(job=job).exclude(matrix_group_label="").exists(): + h.append(MATRIX_GROUP_COLUMN["label"]) + return h diff --git a/backend/pipeline/export_job.py b/backend/pipeline/export_job.py index 5f0be75..a357cc7 100644 --- a/backend/pipeline/export_job.py +++ b/backend/pipeline/export_job.py @@ -16,6 +16,7 @@ from .csv_schema import ( MERGED_FIELD_TO_CSV_HEADER, ) from .dataset_nonempty import ( + MATRIX_GROUP_COLUMN, comment_export_headers, detail_export_headers, merged_export_headers, @@ -34,20 +35,30 @@ from .row_serialize import ( ) -def _search_row_csv_dict(r: JdJobSearchRow, internal_keys: list[str]) -> dict[str, Any]: +def _search_row_csv_dict( + r: JdJobSearchRow, internal_keys: list[str], headers: list[str] +) -> dict[str, Any]: d = search_row_to_dict(r) out: dict[str, Any] = {"id": d["id"], "row_index": d["row_index"]} for k in internal_keys: out[JD_SEARCH_CSV_HEADERS[k]] = d.get(k, "") + zh = MATRIX_GROUP_COLUMN["label"] + if zh in headers: + out[zh] = d.get("matrix_group_label", "") return out -def _detail_row_csv_dict(r: JdJobDetailRow, csv_cols: list[str]) -> dict[str, Any]: +def _detail_row_csv_dict( + r: JdJobDetailRow, csv_cols: list[str], headers: list[str] +) -> dict[str, Any]: d = detail_row_to_dict(r) out: dict[str, Any] = {"id": d["id"], "row_index": d["row_index"]} for col in csv_cols: fn = DETAIL_CSV_TO_FIELD[col] out[col] = d.get(fn, "") + zh = MATRIX_GROUP_COLUMN["label"] + if zh in headers: + out[zh] = d.get("matrix_group_label", "") return out @@ -60,11 +71,16 @@ def _comment_row_csv_dict(r: JdJobCommentRow, csv_cols: list[str]) -> dict[str, return out -def _merged_row_csv_dict(r: JdJobMergedRow, internal_keys: list[str]) -> dict[str, Any]: +def _merged_row_csv_dict( + r: JdJobMergedRow, internal_keys: list[str], headers: list[str] +) -> dict[str, Any]: d = merged_row_to_dict(r) out: dict[str, Any] = {"id": d["id"], "row_index": d["row_index"]} for k in internal_keys: out[MERGED_FIELD_TO_CSV_HEADER[k]] = d.get(k, "") + zh = MATRIX_GROUP_COLUMN["label"] + if zh in headers: + out[zh] = d.get("matrix_group_label", "") return out @@ -98,20 +114,30 @@ def _prune_merged_dict(d: dict[str, Any], fields: list[str]) -> dict[str, Any]: def _rows_as_list_search(job: PipelineJob) -> list[dict[str, Any]]: keys = nonempty_search_keys_for_job(job) + extra_mg = JdJobSearchRow.objects.filter(job=job).exclude(matrix_group_label="").exists() qs = JdJobSearchRow.objects.filter(job=job) - return [ - _prune_search_dict(search_row_to_dict(obj), keys) - for obj in qs.order_by("row_index").iterator(chunk_size=400) - ] + out: list[dict[str, Any]] = [] + for obj in qs.order_by("row_index").iterator(chunk_size=400): + d = search_row_to_dict(obj) + row = _prune_search_dict(d, keys) + if extra_mg: + row["matrix_group_label"] = d.get("matrix_group_label", "") + out.append(row) + return out def _rows_as_list_detail(job: PipelineJob) -> list[dict[str, Any]]: fields = nonempty_detail_fields_for_job(job) + extra_mg = JdJobDetailRow.objects.filter(job=job).exclude(matrix_group_label="").exists() qs = JdJobDetailRow.objects.filter(job=job) - return [ - _prune_detail_dict(detail_row_to_dict(obj), fields) - for obj in qs.order_by("row_index").iterator(chunk_size=400) - ] + out: list[dict[str, Any]] = [] + for obj in qs.order_by("row_index").iterator(chunk_size=400): + d = detail_row_to_dict(obj) + row = _prune_detail_dict(d, fields) + if extra_mg: + row["matrix_group_label"] = d.get("matrix_group_label", "") + out.append(row) + return out def _rows_as_list_comment(job: PipelineJob) -> list[dict[str, Any]]: @@ -125,11 +151,16 @@ def _rows_as_list_comment(job: PipelineJob) -> list[dict[str, Any]]: def _rows_as_list_merged(job: PipelineJob) -> list[dict[str, Any]]: fields = nonempty_merged_fields_for_job(job) + extra_mg = JdJobMergedRow.objects.filter(job=job).exclude(matrix_group_label="").exists() qs = JdJobMergedRow.objects.filter(job=job) - return [ - _prune_merged_dict(merged_row_to_dict(obj), fields) - for obj in qs.order_by("row_index").iterator(chunk_size=400) - ] + out: list[dict[str, Any]] = [] + for obj in qs.order_by("row_index").iterator(chunk_size=400): + d = merged_row_to_dict(obj) + row = _prune_merged_dict(d, fields) + if extra_mg: + row["matrix_group_label"] = d.get("matrix_group_label", "") + out.append(row) + return out def build_json_bytes(*, job: PipelineJob, kind: str) -> tuple[bytes, str]: @@ -178,18 +209,21 @@ def _write_csv_from_qs( def build_csv_bytes(*, job: PipelineJob, kind: str) -> tuple[bytes, str]: if kind == "search": sk = nonempty_search_keys_for_job(job) + headers = search_export_headers(job) text = _write_csv_from_qs( qs=JdJobSearchRow.objects.filter(job=job), - headers=search_export_headers(job), - row_fn=lambda o, _sk=sk: _search_row_csv_dict(o, _sk), + headers=headers, + row_fn=lambda o, _sk=sk, _h=headers: _search_row_csv_dict(o, _sk, _h), ) name = f"job_{job.id}_search.csv" elif kind == "detail": - dcols = [c for c in detail_export_headers(job) if c not in ("id", "row_index")] + headers = detail_export_headers(job) + zh = MATRIX_GROUP_COLUMN["label"] + dcols = [c for c in headers if c not in ("id", "row_index", zh)] text = _write_csv_from_qs( qs=JdJobDetailRow.objects.filter(job=job), - headers=detail_export_headers(job), - row_fn=lambda o, _dc=dcols: _detail_row_csv_dict(o, _dc), + headers=headers, + row_fn=lambda o, _dc=dcols, _h=headers: _detail_row_csv_dict(o, _dc, _h), ) name = f"job_{job.id}_detail.csv" elif kind == "comments": @@ -202,22 +236,28 @@ def build_csv_bytes(*, job: PipelineJob, kind: str) -> tuple[bytes, str]: name = f"job_{job.id}_comments.csv" elif kind == "all": sk = nonempty_search_keys_for_job(job) - dcols = [c for c in detail_export_headers(job) if c not in ("id", "row_index")] + sheaders = search_export_headers(job) + dheaders = detail_export_headers(job) + zh = MATRIX_GROUP_COLUMN["label"] + dcols = [c for c in dheaders if c not in ("id", "row_index", zh)] ccols = [c for c in comment_export_headers(job) if c not in ("id", "row_index")] mk = nonempty_merged_fields_for_job(job) + mheaders = merged_export_headers(job) parts = [ "# search", _write_csv_from_qs( qs=JdJobSearchRow.objects.filter(job=job), - headers=search_export_headers(job), - row_fn=lambda o, _sk=sk: _search_row_csv_dict(o, _sk), + headers=sheaders, + row_fn=lambda o, _sk=sk, _h=sheaders: _search_row_csv_dict(o, _sk, _h), ), "", "# detail", _write_csv_from_qs( qs=JdJobDetailRow.objects.filter(job=job), - headers=detail_export_headers(job), - row_fn=lambda o, _dc=dcols: _detail_row_csv_dict(o, _dc), + headers=dheaders, + row_fn=lambda o, _dc=dcols, _h=dheaders: _detail_row_csv_dict( + o, _dc, _h + ), ), "", "# comments", @@ -230,18 +270,19 @@ def build_csv_bytes(*, job: PipelineJob, kind: str) -> tuple[bytes, str]: "# merged", _write_csv_from_qs( qs=JdJobMergedRow.objects.filter(job=job), - headers=merged_export_headers(job), - row_fn=lambda o, _mk=mk: _merged_row_csv_dict(o, _mk), + headers=mheaders, + row_fn=lambda o, _mk=mk, _h=mheaders: _merged_row_csv_dict(o, _mk, _h), ), ] text = "\n".join(parts) name = f"job_{job.id}_all.csv" elif kind == "merged": mk = nonempty_merged_fields_for_job(job) + headers = merged_export_headers(job) text = _write_csv_from_qs( qs=JdJobMergedRow.objects.filter(job=job), - headers=merged_export_headers(job), - row_fn=lambda o, _mk=mk: _merged_row_csv_dict(o, _mk), + headers=headers, + row_fn=lambda o, _mk=mk, _h=headers: _merged_row_csv_dict(o, _mk, _h), ) name = f"job_{job.id}_merged.csv" else: @@ -262,22 +303,25 @@ def build_xlsx_bytes(*, job: PipelineJob, kind: str) -> tuple[bytes, str]: ws = wb.active ws.title = "search"[:31] sk = nonempty_search_keys_for_job(job) + sheaders = search_export_headers(job) _append_sheet( ws, - search_export_headers(job), + sheaders, JdJobSearchRow.objects.filter(job=job), - lambda o, _sk=sk: _search_row_csv_dict(o, _sk), + lambda o, _sk=sk, _h=sheaders: _search_row_csv_dict(o, _sk, _h), ) name = f"job_{job.id}_search.xlsx" elif kind == "detail": ws = wb.active ws.title = "detail"[:31] - dcols = [c for c in detail_export_headers(job) if c not in ("id", "row_index")] + dheaders = detail_export_headers(job) + zh = MATRIX_GROUP_COLUMN["label"] + dcols = [c for c in dheaders if c not in ("id", "row_index", zh)] _append_sheet( ws, - detail_export_headers(job), + dheaders, JdJobDetailRow.objects.filter(job=job), - lambda o, _dc=dcols: _detail_row_csv_dict(o, _dc), + lambda o, _dc=dcols, _h=dheaders: _detail_row_csv_dict(o, _dc, _h), ) name = f"job_{job.id}_detail.xlsx" elif kind == "comments": @@ -293,23 +337,27 @@ def build_xlsx_bytes(*, job: PipelineJob, kind: str) -> tuple[bytes, str]: name = f"job_{job.id}_comments.xlsx" elif kind == "all": sk = nonempty_search_keys_for_job(job) - dcols = [c for c in detail_export_headers(job) if c not in ("id", "row_index")] + sheaders = search_export_headers(job) + dheaders = detail_export_headers(job) + zh = MATRIX_GROUP_COLUMN["label"] + dcols = [c for c in dheaders if c not in ("id", "row_index", zh)] ccols = [c for c in comment_export_headers(job) if c not in ("id", "row_index")] mk = nonempty_merged_fields_for_job(job) + mheaders = merged_export_headers(job) ws1 = wb.active ws1.title = "search"[:31] _append_sheet( ws1, - search_export_headers(job), + sheaders, JdJobSearchRow.objects.filter(job=job), - lambda o, _sk=sk: _search_row_csv_dict(o, _sk), + lambda o, _sk=sk, _h=sheaders: _search_row_csv_dict(o, _sk, _h), ) ws2 = wb.create_sheet("detail"[:31]) _append_sheet( ws2, - detail_export_headers(job), + dheaders, JdJobDetailRow.objects.filter(job=job), - lambda o, _dc=dcols: _detail_row_csv_dict(o, _dc), + lambda o, _dc=dcols, _h=dheaders: _detail_row_csv_dict(o, _dc, _h), ) ws3 = wb.create_sheet("comments"[:31]) _append_sheet( @@ -321,20 +369,21 @@ def build_xlsx_bytes(*, job: PipelineJob, kind: str) -> tuple[bytes, str]: ws4 = wb.create_sheet("merged"[:31]) _append_sheet( ws4, - merged_export_headers(job), + mheaders, JdJobMergedRow.objects.filter(job=job), - lambda o, _mk=mk: _merged_row_csv_dict(o, _mk), + lambda o, _mk=mk, _h=mheaders: _merged_row_csv_dict(o, _mk, _h), ) name = f"job_{job.id}_all.xlsx" elif kind == "merged": mk = nonempty_merged_fields_for_job(job) ws = wb.active ws.title = "merged"[:31] + mheaders = merged_export_headers(job) _append_sheet( ws, - merged_export_headers(job), + mheaders, JdJobMergedRow.objects.filter(job=job), - lambda o, _mk=mk: _merged_row_csv_dict(o, _mk), + lambda o, _mk=mk, _h=mheaders: _merged_row_csv_dict(o, _mk, _h), ) name = f"job_{job.id}_merged.xlsx" else: diff --git a/backend/pipeline/ingest.py b/backend/pipeline/ingest.py index 19b36c1..f3a8e46 100644 --- a/backend/pipeline/ingest.py +++ b/backend/pipeline/ingest.py @@ -29,12 +29,12 @@ from .csv_schema import ( search_csv_effective_total_sales, strip_buyer_ranking_line_prefix, ) +from .matrix_group_label import matrix_group_label_from_detail_path from .models import ( JdJobCommentRow, JdJobDetailRow, JdJobMergedRow, JdJobSearchRow, - JdLeafCategoryNorm, JdProduct, JdProductSnapshot, PipelineJob, @@ -119,6 +119,33 @@ def _bulk_create_in_chunks(model, objects: list[Any]) -> None: model.objects.bulk_create(objects[i : i + BULK_CHUNK]) +def _sync_search_rows_matrix_labels( + job: PipelineJob, merged_kw_list: list[tuple[int, dict[str, str]]] +) -> None: + """按 SKU 将合并表解析出的报告细类回填到搜索行(与 §5 矩阵口径一致)。""" + sku_to_mg: dict[str, str] = {} + for _, kw in merged_kw_list: + sk = (kw.get("sku_id") or "").strip() + if not sk: + continue + mg = matrix_group_label_from_detail_path(kw.get("detail_category_path") or "") + if mg: + sku_to_mg[sk] = mg + if not sku_to_mg: + return + chunk: list[JdJobSearchRow] = [] + for r in JdJobSearchRow.objects.filter(job=job).iterator(chunk_size=400): + sk = (r.sku_id or "").strip() + if sk and sk in sku_to_mg: + r.matrix_group_label = sku_to_mg[sk] + chunk.append(r) + if len(chunk) >= 400: + JdJobSearchRow.objects.bulk_update(chunk, ["matrix_group_label"]) + chunk.clear() + if chunk: + JdJobSearchRow.objects.bulk_update(chunk, ["matrix_group_label"]) + + def _run_dir(job: PipelineJob) -> Path: return Path(job.run_dir or "").expanduser().resolve() @@ -174,32 +201,12 @@ def ingest_job_dataset_rows(job: PipelineJob) -> dict[str, Any]: if not search_rows and search_path.is_file() is False: pass search_kw_list: list[tuple[int, dict[str, str]]] = [] - leaf_labels: set[str] = set() for i, row in enumerate(search_rows): _normalize_search_csv_total_sales(row) kw = _search_row_kwargs(row) search_kw_list.append((i, kw)) - lc = (kw.get("leaf_category") or "").strip()[:512] - if lc: - leaf_labels.add(lc) - norm_map: dict[str, JdLeafCategoryNorm] = {} - if leaf_labels: - have = set( - JdLeafCategoryNorm.objects.filter(label__in=leaf_labels).values_list( - "label", flat=True - ) - ) - missing = [JdLeafCategoryNorm(label=l) for l in leaf_labels if l not in have] - if missing: - JdLeafCategoryNorm.objects.bulk_create(missing, ignore_conflicts=True) - norm_map = { - n.label: n - for n in JdLeafCategoryNorm.objects.filter(label__in=leaf_labels) - } s_objs: list[JdJobSearchRow] = [] for i, kw in search_kw_list: - lc = (kw.get("leaf_category") or "").strip()[:512] - norm = norm_map.get(lc) if lc else None pv = effective_list_price_value( kw.get("coupon_price"), kw.get("price"), kw.get("original_price") ) @@ -207,7 +214,7 @@ def ingest_job_dataset_rows(job: PipelineJob) -> dict[str, Any]: JdJobSearchRow( job=job, row_index=i, - leaf_category_norm=norm, + matrix_group_label="", price_value=pv, **kw, ) @@ -221,9 +228,14 @@ def ingest_job_dataset_rows(job: PipelineJob) -> dict[str, Any]: for i, row in enumerate(detail_rows): kw = _detail_row_kwargs(row) dpv = float_price_from_cell(kw.get("detail_price_final")) + mg = matrix_group_label_from_detail_path(kw.get("detail_category_path") or "") d_objs.append( JdJobDetailRow( - job=job, row_index=i, detail_price_value=dpv, **kw + job=job, + row_index=i, + matrix_group_label=mg, + detail_price_value=dpv, + **kw, ) ) _bulk_create_in_chunks(JdJobDetailRow, d_objs) @@ -241,34 +253,13 @@ def ingest_job_dataset_rows(job: PipelineJob) -> dict[str, Any]: merged_path = run_dir / FILE_MERGED_CSV merged_rows = _read_csv_rows(merged_path) if merged_path.is_file() else [] merged_kw_list: list[tuple[int, dict[str, str]]] = [] - leaf_labels_m: set[str] = set() for i, row in enumerate(merged_rows): _normalize_merged_csv_total_sales(row) kw = _merged_row_kwargs(row) merged_kw_list.append((i, kw)) - lc = (kw.get("leaf_category") or "").strip()[:512] - if lc: - leaf_labels_m.add(lc) - norm_map_m: dict[str, JdLeafCategoryNorm] = {} - if leaf_labels_m: - have_m = set( - JdLeafCategoryNorm.objects.filter(label__in=leaf_labels_m).values_list( - "label", flat=True - ) - ) - missing_m = [ - JdLeafCategoryNorm(label=l) for l in leaf_labels_m if l not in have_m - ] - if missing_m: - JdLeafCategoryNorm.objects.bulk_create(missing_m, ignore_conflicts=True) - norm_map_m = { - n.label: n - for n in JdLeafCategoryNorm.objects.filter(label__in=leaf_labels_m) - } m_objs: list[JdJobMergedRow] = [] for i, kw in merged_kw_list: - lc = (kw.get("leaf_category") or "").strip()[:512] - norm = norm_map_m.get(lc) if lc else None + mg = matrix_group_label_from_detail_path(kw.get("detail_category_path") or "") pv = effective_list_price_value( kw.get("coupon_price"), kw.get("price"), kw.get("original_price") ) @@ -276,13 +267,14 @@ def ingest_job_dataset_rows(job: PipelineJob) -> dict[str, Any]: JdJobMergedRow( job=job, row_index=i, - leaf_category_norm=norm, + matrix_group_label=mg, price_value=pv, **kw, ) ) _bulk_create_in_chunks(JdJobMergedRow, m_objs) stats["merged_table_rows"] = len(m_objs) + _sync_search_rows_matrix_labels(job, merged_kw_list) return stats diff --git a/backend/pipeline/matrix_group_label.py b/backend/pipeline/matrix_group_label.py new file mode 100644 index 0000000..6034715 --- /dev/null +++ b/backend/pipeline/matrix_group_label.py @@ -0,0 +1,63 @@ +""" +与 ``jd_competitor_report._matrix_group_label_from_path`` 同源: +从商详 ``detail_category_path`` 解析 §5 竞品矩阵用的细类展示名(如「饼干」「米」)。 +""" +from __future__ import annotations + +import re + + +def category_token_meaningless(seg: str) -> bool: + """纯数字类目 ID、空串或疑似内部编码的段,不宜直接作为矩阵分组展示名。""" + t = (seg or "").strip() + if not t: + return True + if t.isdigit(): + return True + if len(t) >= 14 and re.fullmatch(r"[A-Za-z0-9_\-]+", t): + return True + return False + + +def matrix_display_segment_from_parts(parts: list[str]) -> str | None: + """ + 与历史逻辑一致的主选段;若该段无意义则自右向左找第一段可读文本 + (避免「仅类目码」或中间段为数字 ID 时整组成品名式乱桶)。 + """ + if not parts: + return None + if len(parts) >= 4: + preferred = parts[-2] + elif len(parts) >= 3: + preferred = parts[1] + elif len(parts) >= 2: + preferred = parts[1] + else: + preferred = parts[0] + order: list[str] = [] + if preferred: + order.append(preferred) + if len(parts) >= 2: + order.append(parts[-2]) + order.append(parts[-1]) + order.extend(reversed(parts)) + seen: set[str] = set() + for cand in order: + if not cand or cand in seen: + continue + seen.add(cand) + if not category_token_meaningless(cand): + return cand.strip() + return None + + +def matrix_group_label_from_detail_path(path: str) -> str: + """由 ``detail_category_path`` 文本解析细类展示名;空或无可读段则返回空串。""" + t = (path or "").strip() + if not t: + return "" + parts = [p.strip() for p in t.replace(">", ">").split(">") if p.strip()] + if not parts: + return "" + key = matrix_display_segment_from_parts(parts) + return (key[:80] if key else "") diff --git a/backend/pipeline/migrations/0018_report_group_matrix_label.py b/backend/pipeline/migrations/0018_report_group_matrix_label.py new file mode 100644 index 0000000..3bdf135 --- /dev/null +++ b/backend/pipeline/migrations/0018_report_group_matrix_label.py @@ -0,0 +1,92 @@ +# Generated by Django 5.2.1 on 2026-04-16 01:55 + +from django.db import migrations, models + + +def backfill_matrix_group_labels(apps, schema_editor): + from pipeline.matrix_group_label import matrix_group_label_from_detail_path + + JdJobDetailRow = apps.get_model("pipeline", "JdJobDetailRow") + JdJobMergedRow = apps.get_model("pipeline", "JdJobMergedRow") + JdJobSearchRow = apps.get_model("pipeline", "JdJobSearchRow") + + chunk: list = [] + for r in JdJobDetailRow.objects.all().iterator(chunk_size=400): + r.matrix_group_label = matrix_group_label_from_detail_path( + r.detail_category_path or "" + ) + chunk.append(r) + if len(chunk) >= 400: + JdJobDetailRow.objects.bulk_update(chunk, ["matrix_group_label"]) + chunk.clear() + if chunk: + JdJobDetailRow.objects.bulk_update(chunk, ["matrix_group_label"]) + + chunk = [] + for r in JdJobMergedRow.objects.all().iterator(chunk_size=400): + r.matrix_group_label = matrix_group_label_from_detail_path( + r.detail_category_path or "" + ) + chunk.append(r) + if len(chunk) >= 400: + JdJobMergedRow.objects.bulk_update(chunk, ["matrix_group_label"]) + chunk.clear() + if chunk: + JdJobMergedRow.objects.bulk_update(chunk, ["matrix_group_label"]) + + sku_to_mg: dict[str, str] = {} + for r in JdJobMergedRow.objects.exclude(matrix_group_label="").iterator( + chunk_size=400 + ): + sk = str(r.sku_id or "").strip() + if sk: + sku_to_mg[sk] = r.matrix_group_label + + chunk = [] + for r in JdJobSearchRow.objects.all().iterator(chunk_size=400): + sk = str(r.sku_id or "").strip() + r.matrix_group_label = sku_to_mg.get(sk, "") + chunk.append(r) + if len(chunk) >= 400: + JdJobSearchRow.objects.bulk_update(chunk, ["matrix_group_label"]) + chunk.clear() + if chunk: + JdJobSearchRow.objects.bulk_update(chunk, ["matrix_group_label"]) + + +class Migration(migrations.Migration): + + dependencies = [ + ('pipeline', '0017_dataset_browse_filters'), + ] + + operations = [ + migrations.AddField( + model_name='jdjobdetailrow', + name='matrix_group_label', + field=models.CharField(blank=True, db_index=True, default='', help_text='与 §5 矩阵同源:由 detail_category_path 解析', max_length=80, verbose_name='报告细类'), + ), + migrations.AddField( + model_name='jdjobmergedrow', + name='matrix_group_label', + field=models.CharField(blank=True, db_index=True, default='', help_text='与 §5 矩阵同源:由 detail_category_path 解析', max_length=80, verbose_name='报告细类'), + ), + migrations.AddField( + model_name='jdjobsearchrow', + name='matrix_group_label', + field=models.CharField(blank=True, db_index=True, default='', help_text='与 §5 矩阵同源:由合并表商详路径解析;可按 SKU 从合并表回填', max_length=80, verbose_name='报告细类'), + ), + migrations.AddIndex( + model_name='jdjobdetailrow', + index=models.Index(fields=['job', 'matrix_group_label'], name='pipeline_jd_job_id_5595d3_idx'), + ), + migrations.AddIndex( + model_name='jdjobmergedrow', + index=models.Index(fields=['job', 'matrix_group_label'], name='pipeline_jd_job_id_163e3f_idx'), + ), + migrations.AddIndex( + model_name='jdjobsearchrow', + index=models.Index(fields=['job', 'matrix_group_label'], name='pipeline_jd_job_id_38fae5_idx'), + ), + migrations.RunPython(backfill_matrix_group_labels, migrations.RunPython.noop), + ] diff --git a/backend/pipeline/models.py b/backend/pipeline/models.py index 4cf3b25..793f3b8 100644 --- a/backend/pipeline/models.py +++ b/backend/pipeline/models.py @@ -190,6 +190,14 @@ class JdJobSearchRow(models.Model): on_delete=models.SET_NULL, related_name="search_rows", ) + matrix_group_label = models.CharField( + max_length=80, + blank=True, + default="", + db_index=True, + verbose_name="报告细类", + help_text="与 §5 矩阵同源:由合并表商详路径解析;可按 SKU 从合并表回填", + ) price_value = models.FloatField(null=True, blank=True, db_index=True) platform = models.TextField(blank=True, default="") keyword = models.TextField(blank=True, default="") @@ -206,6 +214,7 @@ class JdJobSearchRow(models.Model): indexes = [ models.Index(fields=["job", "sku_id"]), models.Index(fields=["job", "leaf_category_norm"]), + models.Index(fields=["job", "matrix_group_label"]), models.Index(fields=["job", "price_value"]), ] @@ -232,6 +241,14 @@ class JdJobDetailRow(models.Model): buyer_ranking_line = models.TextField(blank=True, default="") buyer_promo_text = models.TextField(blank=True, default="") detail_price_value = models.FloatField(null=True, blank=True, db_index=True) + matrix_group_label = models.CharField( + max_length=80, + blank=True, + default="", + db_index=True, + verbose_name="报告细类", + help_text="与 §5 矩阵同源:由 detail_category_path 解析", + ) class Meta: ordering = ["row_index"] @@ -244,6 +261,7 @@ class JdJobDetailRow(models.Model): indexes = [ models.Index(fields=["job", "sku_id"]), models.Index(fields=["job", "detail_price_value"]), + models.Index(fields=["job", "matrix_group_label"]), ] def __str__(self) -> str: @@ -317,6 +335,14 @@ class JdJobMergedRow(models.Model): on_delete=models.SET_NULL, related_name="merged_rows", ) + matrix_group_label = models.CharField( + max_length=80, + blank=True, + default="", + db_index=True, + verbose_name="报告细类", + help_text="与 §5 矩阵同源:由 detail_category_path 解析", + ) price_value = models.FloatField(null=True, blank=True, db_index=True) keyword = models.TextField(blank=True, default="") page = models.TextField(blank=True, default="") @@ -342,6 +368,7 @@ class JdJobMergedRow(models.Model): indexes = [ models.Index(fields=["job", "sku_id"]), models.Index(fields=["job", "leaf_category_norm"]), + models.Index(fields=["job", "matrix_group_label"]), models.Index(fields=["job", "price_value"]), ] diff --git a/backend/pipeline/row_serialize.py b/backend/pipeline/row_serialize.py index d6ec855..8d42f19 100644 --- a/backend/pipeline/row_serialize.py +++ b/backend/pipeline/row_serialize.py @@ -22,6 +22,7 @@ def search_row_to_dict(r: JdJobSearchRow) -> dict[str, Any]: out: dict[str, Any] = {"id": r.id, "row_index": r.row_index} for k in JD_SEARCH_INTERNAL_KEYS: out[k] = getattr(r, k) or "" + out["matrix_group_label"] = r.matrix_group_label or "" return out @@ -29,6 +30,7 @@ def detail_row_to_dict(r: JdJobDetailRow) -> dict[str, Any]: out: dict[str, Any] = {"id": r.id, "row_index": r.row_index} for k in DETAIL_FIELDS_ORDER: out[k] = getattr(r, k) or "" + out["matrix_group_label"] = r.matrix_group_label or "" return out @@ -43,4 +45,5 @@ def merged_row_to_dict(r: JdJobMergedRow) -> dict[str, Any]: out: dict[str, Any] = {"id": r.id, "row_index": r.row_index} for k in MERGED_FIELDS_ORDER: out[k] = getattr(r, k) or "" + out["matrix_group_label"] = r.matrix_group_label or "" return out diff --git a/backend/pipeline/views.py b/backend/pipeline/views.py index 489d401..21c48f1 100644 --- a/backend/pipeline/views.py +++ b/backend/pipeline/views.py @@ -27,11 +27,11 @@ from .dataset_api import ( apply_merged_order, apply_search_filters, apply_search_order, - category_norm_id_from_request, detail_category_q_from_request, filter_echo, parse_sort_meta, price_bounds_from_request, + report_group_from_request, ) from .dataset_nonempty import ( comment_columns_for_api, @@ -66,7 +66,6 @@ from .models import ( JdJobDetailRow, JdJobMergedRow, JdJobSearchRow, - JdLeafCategoryNorm, JdProduct, JdProductSnapshot, JobStatus, @@ -756,17 +755,15 @@ def _read_page_params(request) -> tuple[int, int]: return page, page_size -def _category_norm_options_for_job(job: PipelineJob, RowModel: type) -> list[dict[str, Any]]: - ids = ( - RowModel.objects.filter(job=job, leaf_category_norm_id__isnull=False) - .values_list("leaf_category_norm_id", flat=True) +def _report_group_options_for_job(job: PipelineJob) -> list[str]: + """与 §5 矩阵一致的细类名列表(来自合并表 ``detail_category_path`` 解析)。""" + qs = ( + JdJobMergedRow.objects.filter(job=job) + .exclude(matrix_group_label="") + .values_list("matrix_group_label", flat=True) .distinct() ) - return list( - JdLeafCategoryNorm.objects.filter(id__in=ids) - .order_by("label") - .values("id", "label") - ) + return sorted({str(x) for x in qs if x}) def _detail_category_path_options(job: PipelineJob) -> list[str]: @@ -797,12 +794,7 @@ class JobDatasetSummaryView(APIView): "detail_columns": detail_columns_for_api(job), "comment_columns": comment_columns_for_api(job), "merged_columns": merged_columns_for_api(job), - "search_category_options": _category_norm_options_for_job( - job, JdJobSearchRow - ), - "merged_category_options": _category_norm_options_for_job( - job, JdJobMergedRow - ), + "report_group_options": _report_group_options_for_job(job), "detail_category_path_options": _detail_category_path_options(job), "dataset_sort_help": { "search": sorted(SEARCH_SORT_FIELDS), @@ -820,7 +812,7 @@ class JobDatasetSearchView(APIView): page, page_size = _read_page_params(request) sort, desc = parse_sort_meta(request) sort_eff = sort if sort in SEARCH_SORT_FIELDS else "row_index" - cid = category_norm_id_from_request(request) + rg = report_group_from_request(request) pmin, pmax = price_bounds_from_request(request) dcq = detail_category_q_from_request(request) qs = JdJobSearchRow.objects.filter(job=job) @@ -835,7 +827,7 @@ class JobDatasetSearchView(APIView): "page": page, "page_size": page_size, "filters": filter_echo( - category_norm_id=cid, + report_group=rg, price_min=pmin, price_max=pmax, detail_category_q=dcq, @@ -853,7 +845,7 @@ class JobDatasetDetailView(APIView): page, page_size = _read_page_params(request) sort, desc = parse_sort_meta(request) sort_eff = sort if sort in DETAIL_SORT_FIELDS else "row_index" - cid = category_norm_id_from_request(request) + rg = report_group_from_request(request) pmin, pmax = price_bounds_from_request(request) dcq = detail_category_q_from_request(request) qs = JdJobDetailRow.objects.filter(job=job) @@ -868,7 +860,7 @@ class JobDatasetDetailView(APIView): "page": page, "page_size": page_size, "filters": filter_echo( - category_norm_id=cid, + report_group=rg, price_min=pmin, price_max=pmax, detail_category_q=dcq, @@ -908,7 +900,7 @@ class JobDatasetMergedView(APIView): page, page_size = _read_page_params(request) sort, desc = parse_sort_meta(request) sort_eff = sort if sort in MERGED_SORT_FIELDS else "row_index" - cid = category_norm_id_from_request(request) + rg = report_group_from_request(request) pmin, pmax = price_bounds_from_request(request) dcq = detail_category_q_from_request(request) qs = JdJobMergedRow.objects.filter(job=job) @@ -923,7 +915,7 @@ class JobDatasetMergedView(APIView): "page": page, "page_size": page_size, "filters": filter_echo( - category_norm_id=cid, + report_group=rg, price_min=pmin, price_max=pmax, detail_category_q=dcq, diff --git a/frontend/src/components/JobDatasetModal.vue b/frontend/src/components/JobDatasetModal.vue index e3b9f39..d20edc0 100644 --- a/frontend/src/components/JobDatasetModal.vue +++ b/frontend/src/components/JobDatasetModal.vue @@ -24,6 +24,7 @@ const SORT_LABELS = { sku_id: 'SKU', title: '标题', leaf_category: '叶类目', + matrix_group_label: '报告细类', detail_category_path: '类目路径', detail_brand: '品牌', } @@ -39,7 +40,8 @@ const err = ref('') const commentSkuFilter = ref('') const sortField = ref('row_index') const sortOrder = ref('asc') -const categoryNormId = ref('') +/** 与 §5 矩阵一致的细类名(如饼干、米),对应接口参数 report_group */ +const reportGroup = ref('') const priceMin = ref('') const priceMax = ref('') const detailCategoryQ = ref('') @@ -68,12 +70,7 @@ const sortOptions = computed(() => { return keys.map((k) => ({ value: k, label: SORT_LABELS[k] || k })) }) -const categoryNormOptions = computed(() => { - if (!summary.value) return [] - if (tab.value === 'search') return summary.value.search_category_options || [] - if (tab.value === 'merged') return summary.value.merged_category_options || [] - return [] -}) +const reportGroupOptions = computed(() => summary.value?.report_group_options || []) const displayColumns = computed(() => { const s = summary.value @@ -164,7 +161,7 @@ async function refreshList() { : { sort: sortField.value, order: sortOrder.value, - categoryNormId: categoryNormId.value.trim(), + reportGroup: reportGroup.value.trim(), priceMin: priceMin.value, priceMax: priceMax.value, detailCategoryQ: detailCategoryQ.value.trim(), @@ -198,7 +195,7 @@ watch( page.value = 1 sortField.value = 'row_index' sortOrder.value = 'asc' - categoryNormId.value = '' + reportGroup.value = '' priceMin.value = '' priceMax.value = '' detailCategoryQ.value = '' @@ -216,7 +213,7 @@ watch(tab, () => { exportPanelOpen.value = false sortField.value = 'row_index' sortOrder.value = 'asc' - categoryNormId.value = '' + reportGroup.value = '' priceMin.value = '' priceMax.value = '' detailCategoryQ.value = '' @@ -226,7 +223,7 @@ watch( [ sortField, sortOrder, - categoryNormId, + reportGroup, priceMin, priceMax, detailCategoryQ, @@ -245,7 +242,7 @@ watch( commentSkuFilter, sortField, sortOrder, - categoryNormId, + reportGroup, priceMin, priceMax, detailCategoryQ, @@ -430,17 +427,13 @@ async function runExport(format) { - +