From 9939188a4bb6747be85aab8be4241ade02c0fc64 Mon Sep 17 00:00:00 2001 From: wsb1224 Date: Sat, 1 Aug 2026 03:15:43 +0800 Subject: [PATCH] =?UTF-8?q?=E4=B8=BB=E8=A6=81=E4=BF=AE=E5=A4=8D=EF=BC=9A?= =?UTF-8?q?=20=E6=B5=B7=E6=8A=A5=20PDF=20=E8=A7=A3=E6=9E=90=E4=BB=8E=20Gun?= =?UTF-8?q?icorn=20=E5=90=8E=E5=8F=B0=E7=BA=BF=E7=A8=8B=E8=BF=81=E7=A7=BB?= =?UTF-8?q?=E5=88=B0=20Celery=20Worker=EF=BC=8C=E6=B6=88=E9=99=A4=20gevent?= =?UTF-8?q?/asyncio.run=20=E5=86=B2=E7=AA=81=E5=92=8C=E4=BB=BB=E5=8A=A1?= =?UTF-8?q?=E5=8D=A1=E6=AD=BB=E3=80=82=20=E6=B5=B7=E6=8A=A5=E6=94=B9?= =?UTF-8?q?=E4=B8=BA=E7=B4=A7=E5=87=91=E8=A7=A3=E6=9E=90=EF=BC=8C=E5=8F=AA?= =?UTF-8?q?=E9=80=89=E6=8B=A9=E5=AE=A2=E6=88=B7=E8=B5=84=E6=96=99=E3=80=81?= =?UTF-8?q?=E4=BF=9D=E8=B4=B9=E5=92=8C=E6=A0=B8=E5=BF=83=E5=88=A9=E7=9B=8A?= =?UTF-8?q?=E9=A1=B5=EF=BC=9B=E6=8E=92=E9=99=A4=E6=8F=90=E9=A2=86=E6=96=B9?= =?UTF-8?q?=E6=A1=88=E3=80=81=E6=82=B2=E8=A7=82/=E4=B9=90=E8=A7=82?= =?UTF-8?q?=E6=83=85=E6=99=AF=E9=A1=B5=E3=80=82=20LLM=20=E8=B0=83=E7=94=A8?= =?UTF-8?q?=E7=94=B1=E5=8E=9F=E6=9D=A5=E7=9A=84=E7=BA=A6=2027=20=E6=AC=A1?= =?UTF-8?q?=E9=99=8D=E4=B8=BA=201=20=E6=AC=A1=E3=80=82=20=E8=A1=A5?= =?UTF-8?q?=E5=85=85=E5=B9=B4=E7=BC=B4=E4=BF=9D=E8=B4=B9=E3=80=81=E9=A6=96?= =?UTF-8?q?=E5=B9=B4=E5=AE=9E=E7=BC=B4=E3=80=81=E7=BC=B4=E8=B4=B9=E6=9C=9F?= =?UTF-8?q?=E3=80=81=E6=80=BB=E4=BF=9D=E8=B4=B9=E5=8F=8A=E7=AC=AC=201/5/10?= =?UTF-8?q?/15/20/25/30=20=E5=B9=B4=E9=80=80=E4=BF=9D=E4=BB=B7=E5=80=BC?= =?UTF-8?q?=E3=80=82=20=E5=A2=9E=E5=8A=A0=E7=9C=9F=E5=AE=9E=E8=A7=A3?= =?UTF-8?q?=E6=9E=90=E8=BF=9B=E5=BA=A6=E3=80=81=E9=94=99=E8=AF=AF=E4=BF=A1?= =?UTF-8?q?=E6=81=AF=E3=80=81=E4=BB=BB=E5=8A=A1=20ID=E3=80=81=E5=BF=83?= =?UTF-8?q?=E8=B7=B3=E5=92=8C=E5=AE=8C=E6=88=90=E6=97=B6=E9=97=B4=E3=80=82?= =?UTF-8?q?=20=E7=9B=B8=E5=90=8C=E7=94=A8=E6=88=B7=E9=87=8D=E5=A4=8D?= =?UTF-8?q?=E4=B8=8A=E4=BC=A0=E5=90=8C=E4=B8=80=E4=BB=BD=E8=AE=A1=E5=88=92?= =?UTF-8?q?=E4=B9=A6=E6=97=B6=E5=A4=8D=E7=94=A8=E7=8E=B0=E6=9C=89=E4=BB=BB?= =?UTF-8?q?=E5=8A=A1=E6=88=96=E7=BB=93=E6=9E=9C=E3=80=82=20=E5=89=8D?= =?UTF-8?q?=E7=AB=AF=E5=8F=96=E6=B6=88=20180=20=E7=A7=92=E6=9C=AC=E5=9C=B0?= =?UTF-8?q?=E5=81=87=E8=B6=85=E6=97=B6=EF=BC=8C=E6=94=B9=E4=B8=BA=E4=B8=B2?= =?UTF-8?q?=E8=A1=8C=E8=BD=AE=E8=AF=A2=E5=90=8E=E7=AB=AF=E7=9C=9F=E5=AE=9E?= =?UTF-8?q?=E7=8A=B6=E6=80=81=EF=BC=9B=E7=BD=91=E7=BB=9C=E6=B3=A2=E5=8A=A8?= =?UTF-8?q?=E4=B8=8D=E5=86=8D=E8=AF=AF=E5=88=A4=E8=A7=A3=E6=9E=90=E5=A4=B1?= =?UTF-8?q?=E8=B4=A5=E3=80=82=20=E5=A2=9E=E5=8A=A0=E6=9C=8D=E5=8A=A1?= =?UTF-8?q?=E9=87=8D=E5=90=AF=E5=90=8E=E7=9A=84=E8=BF=87=E6=9C=9F=E4=BB=BB?= =?UTF-8?q?=E5=8A=A1=E6=81=A2=E5=A4=8D=E6=9C=BA=E5=88=B6=E3=80=82=20?= =?UTF-8?q?=E4=BF=AE=E5=A4=8D=E8=A7=A3=E6=9E=90=E7=BB=93=E6=9E=9C=20JSON?= =?UTF-8?q?=20=E5=BA=8F=E5=88=97=E5=8C=96=E9=81=97=E6=BC=8F=E9=97=AE?= =?UTF-8?q?=E9=A2=98=E3=80=82=20=E5=85=B3=E9=94=AE=E6=96=87=E4=BB=B6?= =?UTF-8?q?=EF=BC=9A=20[extraction.py](D:/work/code/python/coding/baodanag?= =?UTF-8?q?ent/api/insurance/ppt/extraction.py)=20[tasks.py](D:/work/code/?= =?UTF-8?q?python/coding/baodanagent/api/insurance/poster/tasks.py)=20[cel?= =?UTF-8?q?ery=5Ftasks.py](D:/work/code/python/coding/baodanagent/api/insu?= =?UTF-8?q?rance/generation/celery=5Ftasks.py)=20[service.py](D:/work/code?= =?UTF-8?q?/python/coding/baodanagent/api/insurance/poster/service.py)=20[?= =?UTF-8?q?migrate=5F032.py](D:/work/code/python/coding/baodanagent/api/in?= =?UTF-8?q?surance/db/migrate=5F032.py)=20[PosterSourcePanel.vue](D:/work/?= =?UTF-8?q?code/python/coding/baodanagent/frontend/src/components/poster/w?= =?UTF-8?q?orkspace/PosterSourcePanel.vue)=20[=E5=9B=9E=E5=BD=92=E6=B5=8B?= =?UTF-8?q?=E8=AF=95](D:/work/code/python/coding/baodanagent/tests/ppt=5Fp?= =?UTF-8?q?oster=5Foptimization=5Ftest.py)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- api/insurance/db/migrate_032.py | 62 +++++ api/insurance/generation/celery_tasks.py | 46 +++ api/insurance/models/poster_case_upload.py | 36 ++- api/insurance/poster/service.py | 49 +++- api/insurance/poster/tasks.py | 192 ++++++------- api/insurance/ppt/extraction.py | 261 ++++++++++++++++-- api/insurance/routes.py | 2 + .../poster/workspace/PosterSourcePanel.vue | 142 ++++++---- tests/ppt_poster_optimization_test.py | 99 +++++++ 9 files changed, 692 insertions(+), 197 deletions(-) create mode 100644 api/insurance/db/migrate_032.py diff --git a/api/insurance/db/migrate_032.py b/api/insurance/db/migrate_032.py new file mode 100644 index 0000000..8cd4032 --- /dev/null +++ b/api/insurance/db/migrate_032.py @@ -0,0 +1,62 @@ +"""迁移 032:海报计划书解析任务进度与文件去重字段。""" +import logging + +from sqlalchemy import inspect, text + +logger = logging.getLogger(__name__) + + +def migrate(): + from insurance.db.compat import db + + table_name = "poster_case_uploads" + columns = { + item["name"] + for item in inspect(db.engine).get_columns(table_name) + } + additions = { + "parse_progress": "INTEGER NOT NULL DEFAULT 0", + "parse_message": "VARCHAR(500) NOT NULL DEFAULT ''", + "parse_error": "TEXT", + "parse_task_id": "VARCHAR(200)", + "parse_started_at": "TIMESTAMP", + "parse_heartbeat_at": "TIMESTAMP", + "parse_finished_at": "TIMESTAMP", + "file_hash": "VARCHAR(64)", + } + for name, column_type in additions.items(): + if name not in columns: + db.session.execute(text( + f"ALTER TABLE {table_name} ADD COLUMN {name} {column_type}" + )) + + indexes = { + item["name"] + for item in inspect(db.engine).get_indexes(table_name) + } + if "idx_poster_case_owner_hash" not in indexes: + db.session.execute(text( + "CREATE INDEX idx_poster_case_owner_hash " + "ON poster_case_uploads (user_id, file_hash)" + )) + + db.session.execute(text( + "UPDATE poster_case_uploads SET parse_status = 'failed', parse_progress = 100, " + "parse_message = '解析任务已中断', " + "parse_error = '服务升级后请重新解析', parse_finished_at = CURRENT_TIMESTAMP " + "WHERE parse_status IN ('queued', 'parsing')" + )) + db.session.execute(text( + "UPDATE poster_case_uploads SET parse_progress = 100 " + "WHERE parse_status IN ('parsed', 'partial', 'failed')" + )) + db.session.execute(text( + "UPDATE poster_case_uploads SET parse_message = CASE " + "WHEN parse_status = 'parsed' THEN '解析完成' " + "WHEN parse_status = 'partial' THEN '解析完成,部分字段需核对' " + "WHEN parse_status = 'failed' THEN '解析失败' " + "ELSE COALESCE(parse_message, '') END " + "WHERE parse_message IS NULL OR parse_message = ''" + )) + db.session.commit() + logger.info("[migrate_032] 海报计划书解析任务字段迁移完成") diff --git a/api/insurance/generation/celery_tasks.py b/api/insurance/generation/celery_tasks.py index 7b535e4..d0a3d8d 100644 --- a/api/insurance/generation/celery_tasks.py +++ b/api/insurance/generation/celery_tasks.py @@ -1332,6 +1332,52 @@ def _execute_poster_generate(task_id: str): ) +# ─── 海报计划书解析任务 ────────────────────────────────────── + +POSTER_CASE_PARSE_SOFT_TIMEOUT = 300 +POSTER_CASE_PARSE_HARD_TIMEOUT = 360 + + +@shared_task( + bind=True, + name="insurance.parse_poster_case", + queue="insurance", + soft_time_limit=POSTER_CASE_PARSE_SOFT_TIMEOUT, + time_limit=POSTER_CASE_PARSE_HARD_TIMEOUT, +) +def parse_poster_case_task(self, case_upload_id: int): + """在 Worker 中解析海报计划书,数据库状态转换保证幂等。""" + from insurance.db.compat import db + + claimed = db.session.execute( + db.text( + "UPDATE poster_case_uploads SET parse_status = 'parsing', " + "parse_progress = 10, parse_message = '正在读取 PDF 文本', " + "parse_error = NULL, parse_started_at = NOW(), " + "parse_heartbeat_at = NOW(), parse_finished_at = NULL " + "WHERE id = :id AND parse_status = 'queued'" + ), + {"id": case_upload_id}, + ) + db.session.commit() + if claimed.rowcount == 0: + logger.info("海报计划书 %s 已被领取或不在 queued 状态,跳过", case_upload_id) + return + + try: + from insurance.poster.tasks import _execute_case_parse + + _execute_case_parse(case_upload_id) + except Exception as exc: + logger.error("海报计划书解析失败 [%s]: %s", case_upload_id, exc, exc_info=True) + from insurance.poster.tasks import _mark_case_failed + + _mark_case_failed(case_upload_id, str(exc)) + raise + finally: + db.session.remove() + + # ─── 产品小册子解析任务 ────────────────────────────────────── MANUAL_PARSE_SOFT_TIMEOUT = 300 # 5 分钟软超时 diff --git a/api/insurance/models/poster_case_upload.py b/api/insurance/models/poster_case_upload.py index b0bcf90..91be4b5 100644 --- a/api/insurance/models/poster_case_upload.py +++ b/api/insurance/models/poster_case_upload.py @@ -1,5 +1,5 @@ """海报计划书上传模型。""" -from sqlalchemy import Column, String, Text, BigInteger, TIMESTAMP, func +from sqlalchemy import Column, String, Text, BigInteger, Integer, TIMESTAMP, func from insurance.db.compat import db @@ -14,7 +14,19 @@ class PosterCaseUpload(db.Model): product_source_id = Column(String(64), nullable=True, comment="产品来源 ID") product_snapshot_json = Column(Text, nullable=True, comment="产品信息快照 JSON") source_file_url = Column(String(500), nullable=False, comment="源文件地址") - parse_status = Column(String(20), default="pending", comment="解析状态: pending/parsed/failed") + parse_status = Column( + String(20), + default="pending", + comment="解析状态: pending/queued/parsing/parsed/partial/failed", + ) + parse_progress = Column(Integer, nullable=False, default=0, comment="解析进度 0-100") + parse_message = Column(String(500), nullable=False, default="", comment="解析进度说明") + parse_error = Column(Text, nullable=True, comment="解析失败原因") + parse_task_id = Column(String(200), nullable=True, comment="Celery 任务 ID") + parse_started_at = Column(TIMESTAMP, nullable=True) + parse_heartbeat_at = Column(TIMESTAMP, nullable=True) + parse_finished_at = Column(TIMESTAMP, nullable=True) + file_hash = Column(String(64), nullable=True, index=True, comment="源文件 SHA-256") parsed_data = Column(Text, nullable=True, comment="系统解析结果 JSON") confirmed_data = Column(Text, nullable=True, comment="人工核对后最终结果 JSON") confirmed_by = Column(String(50), nullable=True, comment="核对人") @@ -23,6 +35,7 @@ class PosterCaseUpload(db.Model): def to_dict(self): import json + def _safe_json(text): if not text: return None @@ -31,6 +44,18 @@ class PosterCaseUpload(db.Model): except (json.JSONDecodeError, TypeError): return None + def _safe_error(value): + if not value: + return "" + message = str(value).strip() + unsafe_markers = ( + "traceback", "file \"", "\\", "/app/", "/api/", + "http://", "https://", "api_key", "token=", + ) + if "\n" in message or any(marker in message.lower() for marker in unsafe_markers): + return "解析失败,请重试;如多次失败请联系管理员" + return message[:300] + return { "id": self.id, "userId": self.user_id, @@ -42,6 +67,13 @@ class PosterCaseUpload(db.Model): "productSnapshot": _safe_json(self.product_snapshot_json), "sourceFileUrl": self.source_file_url, "parseStatus": self.parse_status, + "parseProgress": self.parse_progress or 0, + "parseMessage": self.parse_message or "", + "parseError": _safe_error(self.parse_error), + "parseTaskId": self.parse_task_id, + "parseStartedAt": self.parse_started_at.isoformat() if self.parse_started_at else None, + "parseHeartbeatAt": self.parse_heartbeat_at.isoformat() if self.parse_heartbeat_at else None, + "parseFinishedAt": self.parse_finished_at.isoformat() if self.parse_finished_at else None, "parsedData": _safe_json(self.parsed_data), "confirmedData": _safe_json(self.confirmed_data), "confirmedBy": self.confirmed_by, diff --git a/api/insurance/poster/service.py b/api/insurance/poster/service.py index 7a182a7..00e5d44 100644 --- a/api/insurance/poster/service.py +++ b/api/insurance/poster/service.py @@ -1,13 +1,13 @@ """海报业务逻辑服务。""" import asyncio import copy +import hashlib import json import os import re import uuid import logging from datetime import datetime, timedelta -from flask import current_app from insurance.db.compat import db from insurance.models.ppt_config import PptProduct, PptCompany from insurance.models.poster_case_upload import PosterCaseUpload @@ -131,6 +131,25 @@ class PosterService: code = 4003 if "密码" in err_msg else 4002 return {"code": code, "message": err_msg, "data": None} + digest = hashlib.sha256(pdf_bytes).hexdigest() + existing = PosterCaseUpload.query.filter_by( + user_id=user_id, + file_hash=digest, + product_source_type=context["sourceType"], + product_source_id=context["sourceId"], + ).filter(PosterCaseUpload.parse_status.in_([ + "queued", "parsing", "parsed", "partial", + ])).order_by(PosterCaseUpload.id.desc()).first() + if existing: + data = existing.to_dict() + data["deduplicated"] = True + message = ( + "相同计划书正在解析" + if existing.parse_status in ("queued", "parsing") + else "已使用相同计划书的解析结果" + ) + return {"code": 0, "message": message, "data": data} + # 保存文件(使用持久化存储) safe_uid = _safe_user_id(user_id) from insurance.config import get_storage_root @@ -149,17 +168,22 @@ class PosterService: product_source_id=context["sourceId"], product_snapshot_json=json.dumps(product_snapshot(context), ensure_ascii=False), source_file_url=filepath, - parse_status="pending", + file_hash=digest, + parse_status="queued", + parse_progress=5, + parse_message="任务已提交,等待解析...", ) db.session.add(record) db.session.commit() # 启动后台解析任务 from insurance.poster.tasks import start_case_parse_task - record.parse_status = "queued" - db.session.commit() - if not start_case_parse_task(current_app._get_current_object(), record.id): + if not start_case_parse_task(record.id): record.parse_status = "failed" + record.parse_progress = 100 + record.parse_message = "任务提交失败" + record.parse_error = "解析任务未能提交到后台队列" + record.parse_finished_at = datetime.now() db.session.commit() return {"code": 5001, "message": "任务排队失败,请重试", "data": None} @@ -192,18 +216,29 @@ class PosterService: if not record.source_file_url or not os.path.exists(record.source_file_url): return {"code": 1002, "message": "源计划书不存在,请重新上传", "data": None} if record.parse_status in ("queued", "parsing"): - return {"code": 1002, "message": "计划书正在解析中", "data": record.to_dict()} + return {"code": 0, "message": "计划书正在解析中", "data": record.to_dict()} from insurance.poster.tasks import start_case_parse_task record.parse_status = "queued" + record.parse_progress = 5 + record.parse_message = "任务已提交,等待解析..." + record.parse_error = None + record.parse_task_id = None + record.parse_started_at = None + record.parse_heartbeat_at = None + record.parse_finished_at = None record.parsed_data = None record.confirmed_data = None record.confirmed_by = None record.confirmed_at = None db.session.commit() - if not start_case_parse_task(current_app._get_current_object(), record.id): + if not start_case_parse_task(record.id): record.parse_status = "failed" + record.parse_progress = 100 + record.parse_message = "任务提交失败" + record.parse_error = "解析任务未能提交到后台队列" + record.parse_finished_at = datetime.now() db.session.commit() return {"code": 5001, "message": "重新解析任务启动失败,请重试", "data": None} return {"code": 0, "data": record.to_dict()} diff --git a/api/insurance/poster/tasks.py b/api/insurance/poster/tasks.py index 560af71..4c579e5 100644 --- a/api/insurance/poster/tasks.py +++ b/api/insurance/poster/tasks.py @@ -1,11 +1,6 @@ -"""海报后台任务管理。 - -使用与 PPT parse_worker 相同的后台线程 + Redis 锁模式。 -任务状态通过数据库字段追踪,前端轮询获取进度。 -""" +"""海报后台任务管理。""" import json import logging -import threading import os from datetime import datetime @@ -13,10 +8,6 @@ from insurance.db.compat import db logger = logging.getLogger(__name__) -_local_locks: set = set() -_local_locks_guard = threading.Lock() -_redis_locks: set = set() - # 超过此时间仍为 queued/generating 的任务视为过期(秒) STALE_TASK_TIMEOUT = 600 # 10 分钟 @@ -28,6 +19,7 @@ def recover_stale_tasks(): """ from insurance.models.poster_record import PosterRecord from insurance.models.poster_case_upload import PosterCaseUpload + from sqlalchemy import inspect cutoff = datetime.now().timestamp() - STALE_TASK_TIMEOUT @@ -43,12 +35,21 @@ def recover_stale_tasks(): logger.warning(f"恢复过期海报任务: record_id={record.id}") # 恢复过期的计划书解析任务 - stale_cases = PosterCaseUpload.query.filter( - PosterCaseUpload.parse_status.in_(["queued", "parsing"]), - PosterCaseUpload.created_at < datetime.fromtimestamp(cutoff), - ).all() + case_columns = { + item["name"] for item in inspect(db.engine).get_columns("poster_case_uploads") + } + stale_cases = [] + if {"parse_progress", "parse_message", "parse_error"}.issubset(case_columns): + stale_cases = PosterCaseUpload.query.filter( + PosterCaseUpload.parse_status.in_(["queued", "parsing"]), + PosterCaseUpload.created_at < datetime.fromtimestamp(cutoff), + ).all() for case in stale_cases: case.parse_status = "failed" + case.parse_progress = 100 + case.parse_message = "解析任务已中断" + case.parse_error = "任务因服务重启或 Worker 中断,请重新解析" + case.parse_finished_at = datetime.now() logger.warning(f"恢复过期解析任务: case_id={case.id}") if stale_records or stale_cases: @@ -58,32 +59,25 @@ def recover_stale_tasks(): # ─── 计划书解析任务 ─────────────────────────────────────── -def start_case_parse_task(app, case_upload_id: int) -> bool: - """启动后台计划书解析任务,返回是否新启动。""" - lock_key = f"poster_case_parse:{case_upload_id}" - if not _acquire_lock(lock_key): +def start_case_parse_task(case_upload_id: int) -> bool: + """将计划书解析提交到 Celery Worker。""" + try: + from insurance.generation.celery_tasks import parse_poster_case_task + from insurance.models.poster_case_upload import PosterCaseUpload + + result = parse_poster_case_task.apply_async( + args=[case_upload_id], + queue="insurance", + ) + record = PosterCaseUpload.query.get(case_upload_id) + if record: + record.parse_task_id = result.id + db.session.commit() + return True + except Exception as exc: + logger.error("海报计划书解析任务提交失败 [%s]: %s", case_upload_id, exc, exc_info=True) return False - thread = threading.Thread( - target=_run_case_parse_task, - args=(app, case_upload_id, lock_key), - daemon=True, - ) - thread.start() - return True - - -def _run_case_parse_task(app, case_upload_id: int, lock_key: str): - with app.app_context(): - try: - _execute_case_parse(case_upload_id) - except Exception as exc: - logger.error(f"计划书解析任务失败 [{case_upload_id}]: {exc}", exc_info=True) - _mark_case_failed(case_upload_id, str(exc)) - finally: - _release_lock(lock_key) - db.session.remove() - def _execute_case_parse(case_upload_id: int): from insurance.models.poster_case_upload import PosterCaseUpload @@ -94,12 +88,7 @@ def _execute_case_parse(case_upload_id: int): filepath = record.source_file_url if not filepath or not os.path.exists(filepath): - record.parse_status = "failed" - db.session.commit() - return - - record.parse_status = "parsing" - db.session.commit() + raise FileNotFoundError("计划书源文件不存在") from insurance.ppt.extraction import ExtractionOrchestrator orchestrator = ExtractionOrchestrator(use_cache=False) @@ -127,40 +116,52 @@ def _execute_case_parse(case_upload_id: int): or "" ) - # 优先使用完整解析(获取利益演示表) + def report_progress(progress: int, message: str): + _update_case_progress(case_upload_id, progress, message) + + # 优先使用海报紧凑解析,只提取基础字段和代表性利益年度。 try: - result = asyncio.run(orchestrator.extract_plan( + compact_data = asyncio.run(orchestrator.extract_for_poster( filepath, plan_type=plan_type, - company_id=company_id, - product_id=record.product_id or "", - product_name_hint=product_name, - product_aliases=product_aliases, + progress_callback=report_progress, )) - if result.status != "error" and result.data: - data = result.data - parsed = _map_extract_plan_fields( - data, result.plan_type, result.status, product_context=product_context, - ) - logger.info(f"使用 extract_plan 解析成功: case_id={case_upload_id}, type={result.plan_type}") + compact_data = compact_data or {} + if product_name and not compact_data.get("product_name"): + compact_data["product_name"] = product_name + parsed = _map_extract_plan_fields( + compact_data, plan_type, "poster_compact", product_context=product_context, + ) + if (parsed.get("meta") or {}).get("status") == "failed": + parsed = None + else: + logger.info("使用海报紧凑解析成功: case_id=%s, type=%s", case_upload_id, plan_type) except Exception as exc: - logger.warning(f"extract_plan 失败,降级到 extract_for_poster: {exc}") + logger.warning("海报紧凑解析失败,降级到完整解析: %s", exc) - # 降级:使用轻量解析 + # 降级:完整解析仍跳过海报不使用的提领方案和销售分析。 if not parsed: try: - fallback_data = asyncio.run(orchestrator.extract_for_poster(filepath)) - parsed = _map_extract_plan_fields( - fallback_data, plan_type, "partial", product_context=product_context, - ) - logger.info(f"使用 extract_for_poster 降级解析: case_id={case_upload_id}") + result = asyncio.run(orchestrator.extract_plan( + filepath, + plan_type=plan_type, + progress_callback=report_progress, + company_id=company_id, + product_id=record.product_id or "", + product_name_hint=product_name, + product_aliases=product_aliases, + include_withdrawal=False, + )) + if result.status != "error" and result.data: + parsed = _map_extract_plan_fields( + result.data, result.plan_type, result.status, product_context=product_context, + ) + logger.info("使用完整解析降级成功: case_id=%s, type=%s", case_upload_id, result.plan_type) except Exception as exc: - logger.error(f"extract_for_poster 也失败: {exc}") - record = PosterCaseUpload.query.get(case_upload_id) - if record: - record.parse_status = "failed" - db.session.commit() - return + logger.error("完整解析降级也失败: %s", exc) + + if not parsed: + raise ValueError("计划书中未识别到可用于海报的字段") record = PosterCaseUpload.query.get(case_upload_id) if not record: @@ -168,6 +169,11 @@ def _execute_case_parse(case_upload_id: int): parse_status = (parsed.get("meta") or {}).get("status", "failed") record.parsed_data = json.dumps(parsed, ensure_ascii=False) record.parse_status = parse_status + record.parse_progress = 100 + record.parse_message = "解析完成" if parse_status == "parsed" else "解析完成,部分字段需核对" + record.parse_error = None + record.parse_heartbeat_at = datetime.now() + record.parse_finished_at = datetime.now() db.session.commit() @@ -395,41 +401,21 @@ def _mark_case_failed(case_upload_id: int, error: str): if not record: return record.parse_status = "failed" + record.parse_progress = 100 + record.parse_message = "解析失败" + record.parse_error = error[:1000] + record.parse_heartbeat_at = datetime.now() + record.parse_finished_at = datetime.now() db.session.commit() -# ─── 锁管理 ──────────────────────────────────────────────── +def _update_case_progress(case_upload_id: int, progress: int, message: str): + from insurance.models.poster_case_upload import PosterCaseUpload -def _acquire_lock(key: str) -> bool: - """获取任务锁(优先 Redis,降级本地内存)。""" - redis_key = f"poster_task_lock:{key}" - try: - from insurance.db.compat import redis_client - if redis_client and redis_client.set(redis_key, "1", nx=True, ex=1800): - _redis_locks.add(key) - return True - if redis_client: - return False - except Exception: - pass - - with _local_locks_guard: - if key in _local_locks: - return False - _local_locks.add(key) - return True - - -def _release_lock(key: str): - """释放任务锁。""" - if key in _redis_locks: - try: - from insurance.db.compat import redis_client - if redis_client: - redis_client.delete(f"poster_task_lock:{key}") - except Exception: - pass - _redis_locks.discard(key) - - with _local_locks_guard: - _local_locks.discard(key) + record = PosterCaseUpload.query.get(case_upload_id) + if not record: + return + record.parse_progress = max(0, min(int(progress), 99)) + record.parse_message = message[:500] + record.parse_heartbeat_at = datetime.now() + db.session.commit() diff --git a/api/insurance/ppt/extraction.py b/api/insurance/ppt/extraction.py index 937c2d9..4a89491 100644 --- a/api/insurance/ppt/extraction.py +++ b/api/insurance/ppt/extraction.py @@ -160,6 +160,118 @@ def _select_page_chunks( return chunks +def _select_poster_pages( + pdf_text: str, + plan_type: str = "savings", + max_pages: int = 3, + max_chars: int = 12000, +) -> str: + """选择海报字段所需的摘要页,排除提领及压力情景表。""" + pages = _split_marked_pages(pdf_text) + if not pages: + return pdf_text[:max_chars] + + target_years = () if plan_type == "ci" else (1, 5, 10, 15, 20, 25, 30) + withdrawal_terms = ( + "款项提取", "款項提取", "提领", "提領", "提款", "withdrawal", + ) + scenario_terms = ( + "悲观情景", "悲觀情景", "乐观情景", "樂觀情景", + "pessimistic scenario", "optimistic scenario", + ) + table_terms = ( + "保单年度", "保單年度", "退保价值", "退保價值", + "policy year", "cash value", "surrender value", "account value", + ) + identity_terms = ( + "拟受保人", "擬受保人", "受保人", "被保人", "年龄", "年齡", + "性别", "性別", "保单货币", "保單貨幣", "annual premium", + ) + + scored = [] + for page_num, content in pages: + lowered = content.lower() + years = { + year for year in target_years + if re.search(rf"(? 0: + selected.append(identity) + selected_pages.add(identity["page"]) + + benefit_candidates = [ + item for item in scored + if item["table_hits"] >= 2 and not item["excluded"] + ] + if not benefit_candidates: + benefit_candidates = [item for item in scored if item["table_hits"] >= 2] + benefit_candidates.sort(key=lambda item: (-item["benefit_score"], item["page"])) + + if benefit_candidates and len(selected) < max_pages: + best = benefit_candidates[0] + if best["page"] not in selected_pages: + selected.append(best) + selected_pages.add(best["page"]) + covered_years = set(best["years"]) + for candidate in benefit_candidates[1:]: + if len(selected) >= max_pages or covered_years.issuperset(target_years): + break + if candidate["page"] in selected_pages or not (candidate["years"] - covered_years): + continue + selected.append(candidate) + selected_pages.add(candidate["page"]) + covered_years.update(candidate["years"]) + + if not selected: + selected.append({"page": pages[0][0], "content": pages[0][1]}) + + output = [] + used_chars = 0 + for item in sorted(selected, key=lambda value: value["page"]): + page_text = f"[PAGE {item['page']}]\n{item['content']}" + remaining = max_chars - used_chars + if remaining <= 50: + break + output.append(page_text[:remaining]) + used_chars += min(len(page_text), remaining) + return "\n\n".join(output) + + def _merge_rows_by_policy_year(rows: list[dict]) -> list[dict]: """按保单年度合并分块结果,重复年度保留字段更完整的一行。""" by_year: dict[int, dict] = {} @@ -946,6 +1058,7 @@ async def _llm_extract_split( plan_type: str, llm_client, progress_callback=None, + include_withdrawal: bool = True, ) -> tuple[dict, any]: """分块 LLM 提取:身份字段 + 利益表 + 提领表分开调用。 @@ -1172,7 +1285,7 @@ async def _llm_extract_split( split_diagnostics["benefit"]["rowCount"] = len(merged_data["benefit_illustration"]) # ── 第三次调用:提领表(可选,仅储蓄险/IUL)── - if normalized_type in ("savings", "iul"): + if include_withdrawal and normalized_type in ("savings", "iul"): withdrawal_chunks = _select_page_chunks( pdf_text, ("提取", "提款", "提领", "提領", "领取", "領取", "withdrawal"), @@ -1286,6 +1399,7 @@ class ExtractionOrchestrator: product_id: str = "", product_name_hint: str = "", product_aliases: Optional[list[str]] = None, + include_withdrawal: bool = True, ) -> ExtractionResult: """从 PDF 提取结构化数据。 @@ -1429,6 +1543,7 @@ class ExtractionOrchestrator: data, response = await _llm_extract_split( pdf_text, plan_type, llm_client, progress_callback, + include_withdrawal=include_withdrawal, ) split_meta = ((data.get("_meta") or {}).get("split_extraction") or {}) extraction_stats["llm_blocks"] = { @@ -1494,7 +1609,7 @@ class ExtractionOrchestrator: product_name = _normalized_product_name(data) status, extraction_error = assess_extraction_payload(data, detected_type) # 只缓存完整结果,避免后续复用错误或缺字段的解析。 - if self.use_cache and status == "success": + if self.use_cache and include_withdrawal and status == "success": self._save_to_cache(abs_path, data) total_ms = (time.time() - start) * 1000 @@ -1570,46 +1685,136 @@ class ExtractionOrchestrator: except Exception as e: logger.warning(f"缓存写入失败: {e}") - async def extract_for_poster(self, filepath: str) -> dict: - """提取海报所需的关键字段(精简 prompt,降低 LLM 成本)。 - - 返回: - {age, gender, currency, sum_assured, premium_term, - annual_premium, coverage_period, key_benefits} - """ + async def extract_for_poster( + self, + filepath: str, + plan_type: str = "savings", + progress_callback: Optional[Callable[[int, str], None]] = None, + ) -> dict: + """一次 LLM 调用提取海报字段和代表性利益年度。""" from insurance.ppt.llm_client import llm_client abs_path = os.path.abspath(filepath) - text, _page_qualities = _extract_pdf_text(abs_path) + if progress_callback: + progress_callback(15, "正在读取 PDF 文本") + text, page_qualities = _extract_pdf_text(abs_path) if not text: raise ValueError("无法提取 PDF 文本") - text = text[:6000] + if not text.startswith("[OCR]") and page_qualities: + low_count = sum( + 1 for quality in page_qualities + if quality.get("quality") in ("corrupted", "low") + ) + if low_count: + if progress_callback: + progress_callback(30, f"检测到 {low_count} 个低质量页,正在 OCR 补充") + text = _merge_page_ocr(text, page_qualities, abs_path) - system_prompt = ( - "你是一位保险计划书解析专家。请从以下计划书内容中提取海报所需的关键字段。\n" - "输出 JSON 格式(不要包含 markdown 代码块标记):\n" - '{"age": 35, "gender": "male", "currency": "USD", "sum_assured": 500000, ' - '"premium_term": 5, "annual_premium": 100000, "coverage_period": "终身", ' - '"key_benefits": ["身故赔偿", "全残保障"]}\n' - "注意:数值用数字,不要带货币符号。gender 用 male/female。" + selected_text = _select_poster_pages(text, plan_type=plan_type) + if progress_callback: + progress_callback(55, "正在识别客户信息和关键利益年度") + + prompt = ( + "从以下保险计划书页面中提取生成海报需要的数据。\n" + f"产品类型:{plan_type}\n" + "只提取正文明确出现的值,无法确认时填 null,禁止编造。\n" + "利益演示只保留页面中存在的保单年度 1、5、10、15、20、25、30;" + "普通退保价值不要与提领、提款或压力情景混淆。\n" + "annual_premium 必须取基本计划的标准每年保费;折扣后的首年应缴金额写入 " + "first_year_amount_due,不能覆盖 annual_premium。\n" + "每一行利益数据必须直接输出 guaranteed_cash_value、non_guaranteed_cash_value、" + "total_surrender_value、death_benefit,不要嵌套 surrender_value 或 death_benefit 对象。\n" + "金额去除货币符号和千位逗号并输出数字,gender 使用 male/female。\n\n" + f"PDF 关键页面:\n{selected_text}" ) - result, _response = await llm_client.structured_output( - text, system_prompt, + result, response = await llm_client.structured_output( + prompt=prompt, + system_prompt="你是保险计划书数据提取专家。只输出符合 Schema 的 JSON。", schema={ "type": "object", "properties": { - "age": {"type": "number"}, - "gender": {"type": "string"}, - "currency": {"type": "string"}, - "sum_assured": {"type": "number"}, - "premium_term": {"type": "number"}, - "annual_premium": {"type": "number"}, - "coverage_period": {"type": "string"}, + "product_name": {"type": ["string", "null"]}, + "insured": { + "type": "object", + "properties": { + "age": {"type": ["number", "null"]}, + "gender": {"type": ["string", "null"]}, + "smoking_status": {"type": ["string", "null"]}, + }, + }, + "policy": { + "type": "object", + "properties": { + "currency": {"type": ["string", "null"]}, + "sum_insured": {"type": ["number", "null"]}, + "annual_premium": {"type": ["number", "null"]}, + "basic_plan_annual_premium": {"type": ["number", "null"]}, + "first_year_amount_due": {"type": ["number", "null"]}, + "premium_payment_period": {"type": ["number", "null"]}, + "coverage_period": {"type": ["string", "null"]}, + "total_premium": {"type": ["number", "null"]}, + "initial_death_benefit": {"type": ["number", "null"]}, + }, + }, + "benefit_illustration": { + "type": "array", + "items": { + "type": "object", + "properties": { + "policy_year": {"type": "number"}, + "total_premium_paid": {"type": ["number", "null"]}, + "guaranteed_cash_value": {"type": ["number", "null"]}, + "non_guaranteed_cash_value": {"type": ["number", "null"]}, + "total_surrender_value": {"type": ["number", "null"]}, + "death_benefit": {"type": ["number", "null"]}, + "source_page": {"type": ["number", "null"]}, + }, + "required": ["policy_year"], + }, + }, "key_benefits": {"type": "array", "items": {"type": "string"}}, }, - "required": ["age", "gender", "currency", "sum_assured", "annual_premium"], + "required": ["insured", "policy", "benefit_illustration"], + "additionalProperties": False, }, + temperature=0, ) + target_years = {1, 5, 10, 15, 20, 25, 30} + normalized_rows = [] + for row in result.get("benefit_illustration") or []: + if not isinstance(row, dict): + continue + normalized = dict(row) + surrender = normalized.pop("surrender_value", None) + if isinstance(surrender, dict): + normalized.setdefault("guaranteed_cash_value", surrender.get("guaranteed")) + normalized.setdefault("non_guaranteed_cash_value", surrender.get("non_guaranteed")) + normalized.setdefault("total_surrender_value", surrender.get("total")) + death = normalized.get("death_benefit") + if isinstance(death, dict): + normalized["death_benefit"] = death.get("total") + normalized_rows.append(normalized) + rows = _merge_rows_by_policy_year(normalized_rows) + filtered_rows = [] + for row in rows: + try: + year = int(row.get("policy_year")) + except (TypeError, ValueError): + continue + if year in target_years: + filtered_rows.append(row) + result["benefit_illustration"] = filtered_rows + result = _apply_filename_hints(result, abs_path, plan_type) + result.setdefault("_meta", {}).update({ + "method": "poster_compact", + "low_quality_pages": [ + quality["page"] for quality in page_qualities + if quality.get("quality") in ("corrupted", "low") + ], + "latency_ms": getattr(response, "latency_ms", 0) or 0, + }) + if progress_callback: + progress_callback(90, "关键字段识别完成,正在校验") return result diff --git a/api/insurance/routes.py b/api/insurance/routes.py index 85eb510..41c59ea 100644 --- a/api/insurance/routes.py +++ b/api/insurance/routes.py @@ -86,6 +86,8 @@ def _recover_stale_tasks(app: Flask): recover_stale_tasks() from insurance.generation.celery_tasks import recover_stale_manual_tasks recover_stale_manual_tasks() + from insurance.poster.tasks import recover_stale_tasks as recover_stale_poster_tasks + recover_stale_poster_tasks() except Exception as e: app.logger.warning(f"恢复过期任务失败: {e}") diff --git a/frontend/src/components/poster/workspace/PosterSourcePanel.vue b/frontend/src/components/poster/workspace/PosterSourcePanel.vue index 8d94510..4bf468e 100644 --- a/frontend/src/components/poster/workspace/PosterSourcePanel.vue +++ b/frontend/src/components/poster/workspace/PosterSourcePanel.vue @@ -43,10 +43,16 @@ -
+
{{ parseMessage || '解析中...' }} + {{ parseProgress }}%
@@ -179,7 +185,7 @@