Files
tyapi-server/scripts/batch_flxg7e8f_analyze.py
2026-06-27 23:56:18 +08:00

574 lines
20 KiB
Python
Raw Permalink Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

#!/usr/bin/env python3
"""
批量调用远程 FLXG7E8F分析执行/失信/限高三类案件中的金额字段。
金额字段含义(见 个人司法涉诉查询_返回字段说明:
执行案件 implement.cases[]:
n_sqzxbdje 申请执行标的金额
n_sjdwje 实际到位金额(已还款/已执行到位repaidAmount 应对应此字段)
n_wzxje 未执行金额(未还款/未执行到位,正确的「未还款金额」应对应此字段)
n_jabdje 结案标的金额(民事结案金额,不能当作还款金额)
失信 breachCaseList[]:
estimatedJudgementAmount 判决金额估计(非未还款字段)
无 n_wzxje可尝试按案号关联 implement 案件取 n_wzxje
限高 consumptionRestrictionList[]:
无金额字段;可尝试按案号关联 implement 案件取 n_wzxje
依赖: pip install pycryptodome requests
示例:
python batch_flxg7e8f_analyze.py \\
--input ..\\query_1-2026-06-27_32042_decrypted.json \\
--access-id <你的AccessId> \\
--secret-key f507c58537aac8227b5f4a99cbd317ff \\
--output ..\\flxg7e8f_analysis.json \\
--limit 3
获取 access_id:
SELECT access_id, secret_key FROM api_users
WHERE user_id = '2acfef92-0700-4638-b386-bd03a8ad09c3';
"""
from __future__ import annotations
import argparse
import base64
import json
import os
import re
import sys
import time
from pathlib import Path
from typing import Any
try:
import requests
from Crypto.Cipher import AES
from Crypto.Random import get_random_bytes
except ImportError:
print("请先安装: pip install pycryptodome requests", file=sys.stderr)
sys.exit(1)
DEFAULT_BASE_URL = "https://api.tianyuanapi.com"
API_NAME = "FLXG7E8F"
def aes_encrypt(plain_text: bytes, key_hex: str) -> str:
key = bytes.fromhex(key_hex.strip())
block = 16
pad = block - len(plain_text) % block
plain_text = plain_text + bytes([pad]) * pad
iv = get_random_bytes(16)
ct = AES.new(key, AES.MODE_CBC, iv).encrypt(plain_text)
return base64.b64encode(iv + ct).decode()
def aes_decrypt(cipher_b64: str, key_hex: str) -> Any:
key = bytes.fromhex(key_hex.strip())
raw = base64.b64decode(cipher_b64.strip())
iv, ciphertext = raw[:16], raw[16:]
plain = AES.new(key, AES.MODE_CBC, iv).decrypt(ciphertext)
pad = plain[-1]
plain = plain[:-pad]
return json.loads(plain.decode("utf-8"))
def to_float(val: Any) -> float | None:
if val is None or val == "":
return None
if isinstance(val, (int, float)):
return float(val)
if isinstance(val, str):
s = val.strip().replace(",", "")
if not s:
return None
try:
return float(s)
except ValueError:
return None
return None
def normalize_case_no(case_no: str) -> str:
if not case_no:
return ""
s = case_no.strip().upper()
s = re.sub(r"\s+", "", s)
return s
def pick_amount_fields(case_map: dict) -> dict[str, Any]:
sqzx = to_float(case_map.get("n_sqzxbdje"))
sjdw = to_float(case_map.get("n_sjdwje"))
wzx = to_float(case_map.get("n_wzxje"))
jabd = to_float(case_map.get("n_jabdje"))
computed_unpaid = None
if sqzx is not None and sjdw is not None:
computed_unpaid = max(sqzx - sjdw, 0.0)
has_correct_unpaid_field = wzx is not None
has_correct_repaid_field = sjdw is not None
wrong_if_use_jabdje_as_repaid = (
not has_correct_repaid_field and jabd is not None and jabd > 0
)
return {
"n_sqzxbdje": sqzx,
"n_sjdwje": sjdw,
"n_wzxje": wzx,
"n_jabdje": jabd,
"computed_unpaid_sqzx_minus_sjdw": computed_unpaid,
"has_correct_unpaid_field": has_correct_unpaid_field,
"has_correct_repaid_field": has_correct_repaid_field,
"wrong_if_use_jabdje_as_repaid": wrong_if_use_jabdje_as_repaid,
"unpaid_amount": wzx if wzx is not None else computed_unpaid,
}
def build_implement_index(implement_cases: list[dict]) -> dict[str, dict]:
index: dict[str, dict] = {}
for case in implement_cases:
case_no = normalize_case_no(str(case.get("c_ah") or ""))
if not case_no:
continue
amounts = pick_amount_fields(case)
index[case_no] = {
"case_number": case.get("c_ah"),
"case_status": case.get("n_ajjzjd"),
"filing_date": case.get("d_larq"),
"close_date": case.get("d_jarq"),
**amounts,
}
return index
def is_case_fully_paid(amounts: dict) -> bool:
"""执行案件是否已全部到位。"""
wzx = amounts.get("n_wzxje")
if wzx is not None:
return wzx <= 0
sqzx = amounts.get("n_sqzxbdje")
sjdw = amounts.get("n_sjdwje")
if sqzx is not None and sjdw is not None:
return sjdw >= sqzx
unpaid = amounts.get("unpaid_amount")
if unpaid is not None:
return unpaid <= 0
return True
def has_outstanding_unpaid(amounts: dict) -> bool:
wzx = amounts.get("n_wzxje")
if wzx is not None:
return wzx > 0
unpaid = amounts.get("unpaid_amount")
if unpaid is not None:
return unpaid > 0
sqzx = amounts.get("n_sqzxbdje")
sjdw = amounts.get("n_sjdwje")
if sqzx is not None and sjdw is not None and sqzx > sjdw:
return True
return False
def classify_analysis(analysis: dict) -> tuple[str, str]:
"""
判定是否有「金额映射问题」需要关注:
NO_ISSUE - 无执行/失信/限高,或金额已全部到位
HAS_OUTSTANDING_UNPAID - 存在未结/未执行金额,需用 n_wzxje / n_sjdwje 正确映射
"""
if not analysis.get("hit"):
return "NO_ISSUE", "未命中司法数据"
implement_cases = analysis.get("implement_cases") or []
breach_cases = analysis.get("breach_cases") or []
limit_cases = analysis.get("limit_cases") or []
if not implement_cases and not breach_cases and not limit_cases:
return "NO_ISSUE", "无执行/失信/限高案件"
outstanding_details: list[str] = []
for case in implement_cases:
if has_outstanding_unpaid(case):
outstanding_details.append(f"执行 {case.get('case_number')} 未执行金额={case.get('n_wzxje') or case.get('unpaid_amount')}")
for case in breach_cases:
fulfill = (case.get("fulfill_status") or "").strip()
linked = case.get("linked_implement") or {}
if fulfill == "全部未履行":
if linked and has_outstanding_unpaid(linked):
outstanding_details.append(f"失信 {case.get('case_number')} 关联未执行={linked.get('n_wzxje')}")
elif not linked:
outstanding_details.append(f"失信 {case.get('case_number')} 全部未履行(无关联执行明细)")
elif linked and has_outstanding_unpaid(linked):
outstanding_details.append(f"失信 {case.get('case_number')} 关联未执行={linked.get('n_wzxje')}")
for case in limit_cases:
linked = case.get("linked_implement") or {}
if linked and has_outstanding_unpaid(linked):
outstanding_details.append(f"限高 {case.get('case_number')} 关联未执行={linked.get('n_wzxje')}")
if outstanding_details:
return "HAS_OUTSTANDING_UNPAID", "; ".join(outstanding_details)
return "NO_ISSUE", "有执行/失信/限高案件,但金额已全部到位或无未结金额"
def analyze_judicial_data(judicial_data: dict) -> dict[str, Any]:
if not judicial_data or judicial_data == -1:
return {"hit": False, "message": "未命中司法数据"}
lawsuit_stat = judicial_data.get("lawsuitStat") or {}
implement_section = lawsuit_stat.get("implement") or {}
implement_raw = implement_section.get("cases") or []
implement_cases = []
for item in implement_raw:
if isinstance(item, dict):
case_no = item.get("c_ah")
amounts = pick_amount_fields(item)
implement_cases.append(
{
"case_number": case_no,
"case_status": item.get("n_ajjzjd"),
"filing_date": item.get("d_larq"),
"close_date": item.get("d_jarq"),
**amounts,
}
)
implement_index = build_implement_index(
[c for c in implement_raw if isinstance(c, dict)]
)
breach_cases = []
for item in judicial_data.get("breachCaseList") or []:
if not isinstance(item, dict):
continue
case_no = normalize_case_no(str(item.get("caseNumber") or ""))
linked = implement_index.get(case_no)
est = to_float(item.get("estimatedJudgementAmount"))
breach_cases.append(
{
"case_number": item.get("caseNumber"),
"fulfill_status": item.get("fulfillStatus"),
"estimated_judgement_amount": est,
"issue_date": item.get("issueDate"),
"linked_implement": linked,
"has_correct_unpaid_via_link": bool(
linked and linked.get("has_correct_unpaid_field")
),
"unpaid_amount": linked.get("unpaid_amount") if linked else None,
}
)
limit_cases = []
for item in judicial_data.get("consumptionRestrictionList") or []:
if not isinstance(item, dict):
continue
case_no = normalize_case_no(str(item.get("caseNumber") or ""))
linked = implement_index.get(case_no)
limit_cases.append(
{
"case_number": item.get("caseNumber"),
"issue_date": item.get("issueDate"),
"file_date": item.get("fileDate"),
"court": item.get("executiveCourt"),
"linked_implement": linked,
"has_correct_unpaid_via_link": bool(
linked and linked.get("has_correct_unpaid_field")
),
"unpaid_amount": linked.get("unpaid_amount") if linked else None,
}
)
implement_with_outstanding = [c for c in implement_cases if has_outstanding_unpaid(c)]
breach_with_outstanding = [
c
for c in breach_cases
if (c.get("fulfill_status") == "全部未履行")
or (c.get("linked_implement") and has_outstanding_unpaid(c["linked_implement"]))
]
limit_with_outstanding = [
c
for c in limit_cases
if c.get("linked_implement") and has_outstanding_unpaid(c["linked_implement"])
]
result = {
"hit": True,
"summary": {
"implement_count": len(implement_cases),
"breach_count": len(breach_cases),
"limit_count": len(limit_cases),
"implement_with_outstanding_unpaid": len(implement_with_outstanding),
"breach_with_outstanding": len(breach_with_outstanding),
"limit_with_outstanding": len(limit_with_outstanding),
},
"implement_cases": implement_cases,
"breach_cases": breach_cases,
"limit_cases": limit_cases,
}
flag, reason = classify_analysis(result)
result["flag"] = flag
result["flag_reason"] = reason
return result
def call_flxg7e8f(
base_url: str,
access_id: str,
secret_key: str,
name: str,
id_card: str,
timeout: int = 60,
) -> tuple[bool, dict[str, Any]]:
url = f"{base_url.rstrip('/')}/api/v1/{API_NAME}"
params = {"name": name, "id_card": id_card}
encrypted = aes_encrypt(json.dumps(params, ensure_ascii=False).encode("utf-8"), secret_key)
payload = {"data": encrypted, "options": {"json": True}}
resp = requests.post(
url,
headers={"Access-Id": access_id, "Content-Type": "application/json"},
json=payload,
timeout=timeout,
)
resp.raise_for_status()
body = resp.json()
if body.get("code") != 0:
return False, {
"error_code": body.get("code"),
"error_message": body.get("message"),
"transaction_id": body.get("transaction_id"),
}
decrypted = aes_decrypt(body["data"], secret_key)
return True, {
"transaction_id": body.get("transaction_id"),
"response": decrypted,
}
def load_input_records(path: Path) -> list[dict]:
data = json.loads(path.read_text(encoding="utf-8-sig"))
if not isinstance(data, list):
raise ValueError("输入 JSON 应为数组")
return data
def dedupe_persons(records: list[dict]) -> list[dict]:
seen: set[str] = set()
persons = []
for rec in records:
rp = rec.get("request_params") or {}
if isinstance(rp, str):
continue
id_card = (rp.get("id_card") or "").strip().upper()
name = (rp.get("name") or "").strip()
if not id_card or not name:
continue
if id_card in seen:
continue
seen.add(id_card)
persons.append(
{
"name": name,
"id_card": id_card,
"source_ids": [rec.get("id")],
"source_transaction_ids": [rec.get("transaction_id")],
}
)
# attach all source ids for duplicates
id_to_person = {p["id_card"]: p for p in persons}
for rec in records:
rp = rec.get("request_params") or {}
if isinstance(rp, str):
continue
id_card = (rp.get("id_card") or "").strip().upper()
if id_card in id_to_person:
p = id_to_person[id_card]
rid, rtx = rec.get("id"), rec.get("transaction_id")
if rid and rid not in p["source_ids"]:
p["source_ids"].append(rid)
if rtx and rtx not in p["source_transaction_ids"]:
p["source_transaction_ids"].append(rtx)
return persons
def main() -> None:
parser = argparse.ArgumentParser(description="批量 FLXG7E8F + 未还款金额分析")
parser.add_argument("--input", required=True, help="解密后的 COMBMY01 JSON")
parser.add_argument("--output", default="flxg7e8f_analysis.json")
parser.add_argument("--raw-output", default="flxg7e8f_raw.json", help="FLXG7E8F 解密后源数据汇总")
parser.add_argument("--access-id", default=os.getenv("TYAPI_ACCESS_ID"))
parser.add_argument("--secret-key", default=os.getenv("TYAPI_SECRET_KEY"))
parser.add_argument("--base-url", default=os.getenv("TYAPI_BASE_URL", DEFAULT_BASE_URL))
parser.add_argument("--delay", type=float, default=0.5, help="每次请求间隔秒数")
parser.add_argument("--limit", type=int, default=0, help="仅处理前 N 个唯一身份证0=全部)")
parser.add_argument("--resume", help="已有结果 JSON跳过已成功的人员")
parser.add_argument("--dry-run", action="store_true", help="只列出待调用人员,不请求 API")
args = parser.parse_args()
if not args.secret_key:
print("请提供 --secret-key 或环境变量 TYAPI_SECRET_KEY", file=sys.stderr)
sys.exit(1)
input_path = Path(args.input)
records = load_input_records(input_path)
persons = dedupe_persons(records)
if args.limit > 0:
persons = persons[: args.limit]
print(f"输入记录 {len(records)} 条, 去重后 {len(persons)}")
existing: dict[str, dict] = {}
if args.resume and Path(args.resume).is_file():
prev = json.loads(Path(args.resume).read_text(encoding="utf-8"))
for item in prev.get("persons", []):
if item.get("api_success"):
existing[item["id_card"]] = item
if args.dry_run:
for p in persons:
print(f"{p['name']}\t{p['id_card']}\tsources={len(p['source_ids'])}")
return
if not args.access_id:
print("请提供 --access-id 或环境变量 TYAPI_ACCESS_ID", file=sys.stderr)
sys.exit(1)
results = []
raw_records = []
ok_count = fail_count = no_issue_count = outstanding_count = 0
for idx, person in enumerate(persons, start=1):
id_card = person["id_card"]
if id_card in existing:
cached = existing[id_card]
results.append(cached)
if cached.get("raw_response"):
raw_records.append(
{
"name": cached.get("name"),
"id_card": id_card,
"transaction_id": cached.get("transaction_id"),
"raw_response": cached.get("raw_response"),
}
)
print(f"[{idx}/{len(persons)}] 跳过(已存在) {person['name']} {id_card}")
continue
print(f"[{idx}/{len(persons)}] 调用 {person['name']} {id_card} ...", flush=True)
item = {
"name": person["name"],
"id_card": id_card,
"source_ids": person["source_ids"],
"source_transaction_ids": person["source_transaction_ids"],
"api_success": False,
}
try:
success, api_result = call_flxg7e8f(
args.base_url, args.access_id, args.secret_key, person["name"], id_card
)
item["api_success"] = success
if not success:
item.update(api_result)
fail_count += 1
else:
item["transaction_id"] = api_result.get("transaction_id")
item["raw_response"] = api_result.get("response")
raw_records.append(
{
"name": person["name"],
"id_card": id_card,
"transaction_id": item["transaction_id"],
"source_ids": person["source_ids"],
"raw_response": item["raw_response"],
}
)
judicial = (api_result.get("response") or {}).get("judicial_data")
analysis = analyze_judicial_data(judicial)
item["analysis"] = analysis
item["flag"] = analysis.get("flag", "NO_ISSUE")
item["flag_reason"] = analysis.get("flag_reason", "")
ok_count += 1
if item["flag"] == "HAS_OUTSTANDING_UNPAID":
outstanding_count += 1
else:
no_issue_count += 1
except Exception as exc:
item["api_success"] = False
item["error_message"] = str(exc)
fail_count += 1
results.append(item)
# 增量保存,防止中断丢数据
out_path = Path(args.output)
raw_path = Path(args.raw_output)
payload = {
"meta": {
"input": str(input_path),
"total_records": len(records),
"unique_persons": len(persons),
"processed": len(results),
},
"persons": results,
}
out_path.write_text(json.dumps(payload, ensure_ascii=False, indent=2), encoding="utf-8")
raw_path.write_text(
json.dumps(
{
"meta": payload["meta"],
"records": raw_records,
},
ensure_ascii=False,
indent=2,
),
encoding="utf-8",
)
if idx < len(persons) and args.delay > 0:
time.sleep(args.delay)
summary = {
"api_ok": ok_count,
"api_fail": fail_count,
"no_issue": no_issue_count,
"has_outstanding_unpaid": outstanding_count,
}
final = {
"meta": {
"input": str(input_path),
"total_records": len(records),
"unique_persons": len(persons),
"summary": summary,
"raw_output": str(Path(args.raw_output).resolve()),
},
"persons": results,
}
Path(args.output).write_text(json.dumps(final, ensure_ascii=False, indent=2), encoding="utf-8")
Path(args.raw_output).write_text(
json.dumps({"meta": final["meta"], "records": raw_records}, ensure_ascii=False, indent=2),
encoding="utf-8",
)
print("\n=== 完成 ===")
print(json.dumps(summary, ensure_ascii=False, indent=2))
print(f"分析结果: {Path(args.output).resolve()}")
print(f"源数据: {Path(args.raw_output).resolve()}")
if __name__ == "__main__":
main()