93 lines
3.2 KiB
Python
93 lines
3.2 KiB
Python
from __future__ import annotations
|
|
|
|
from online_mining_common import (
|
|
case_matches,
|
|
compact_text,
|
|
emit_error,
|
|
emit_success,
|
|
extract_case,
|
|
load_json_payload,
|
|
search_docs,
|
|
source_config,
|
|
string_values,
|
|
write_optional_jsonl,
|
|
)
|
|
|
|
|
|
def build_index_filters(source: str, filters: dict) -> list[dict]:
|
|
"""尽量把可索引条件下推到 ES,其余条件在客户端二次过滤。"""
|
|
config = source_config(source)
|
|
result: list[dict] = []
|
|
if source == "main":
|
|
for key in ("domain", "func"):
|
|
values = string_values(filters.get(key))
|
|
if values:
|
|
result.append({"terms": {key: values}})
|
|
contains = string_values(filters.get("query_contains"))
|
|
if contains:
|
|
# query 是 keyword 字段,wildcard 可用但不要放太宽,调用方要限制日期和 scan_size。
|
|
for item in contains:
|
|
result.append({"wildcard": {"query": {"value": f"*{item}*"}}})
|
|
if filters.get("request_id"):
|
|
result.append({"terms": {str(config["id_field"]): string_values(filters.get("request_id"))}})
|
|
return result
|
|
|
|
|
|
def main() -> int:
|
|
try:
|
|
payload = load_json_payload()
|
|
source = str(payload.get("source") or "pre_processing")
|
|
date = payload.get("date")
|
|
date_from = payload.get("date_from")
|
|
date_to = payload.get("date_to")
|
|
lookback_days = payload.get("lookback_days")
|
|
lookback_days = int(lookback_days) if lookback_days not in (None, "") else None
|
|
size = int(payload.get("size") or 50)
|
|
scan_size = int(payload.get("scan_size") or max(size * 5, size))
|
|
filters = payload.get("filters") or {}
|
|
if not isinstance(filters, dict):
|
|
raise ValueError("filters must be an object")
|
|
docs = search_docs(
|
|
source=source,
|
|
date=date,
|
|
size=scan_size,
|
|
query_filters=build_index_filters(source, filters),
|
|
date_from=date_from,
|
|
date_to=date_to,
|
|
lookback_days=lookback_days,
|
|
)
|
|
cases = []
|
|
for doc in docs:
|
|
case = extract_case(source, doc)
|
|
if not case_matches(case, filters):
|
|
continue
|
|
cases.append(case)
|
|
if len(cases) >= size:
|
|
break
|
|
output_path = write_optional_jsonl(str(payload.get("output_path") or ""), cases)
|
|
emit_success(
|
|
{
|
|
"source": source,
|
|
"date": date or (f"{date_from}..{date_to}" if date_from and date_to else f"past-{lookback_days}d" if lookback_days else "past-48h"),
|
|
"scanned": len(docs),
|
|
"matched": len(cases),
|
|
"output_path": output_path,
|
|
"cases": [
|
|
{
|
|
key: compact_text(value, 1200)
|
|
for key, value in case.items()
|
|
if key != "raw" and value not in (None, "", [], {})
|
|
}
|
|
for case in cases
|
|
],
|
|
}
|
|
)
|
|
return 0
|
|
except Exception as exc: # noqa: BLE001
|
|
emit_error(exc)
|
|
return 1
|
|
|
|
|
|
if __name__ == "__main__":
|
|
raise SystemExit(main())
|