主要修复:

海报 PDF 解析从 Gunicorn 后台线程迁移到 Celery Worker,消除 gevent/asyncio.run 冲突和任务卡死。
海报改为紧凑解析,只选择客户资料、保费和核心利益页;排除提领方案、悲观/乐观情景页。
LLM 调用由原来的约 27 次降为 1 次。
补充年缴保费、首年实缴、缴费期、总保费及第 1/5/10/15/20/25/30 年退保价值。
增加真实解析进度、错误信息、任务 ID、心跳和完成时间。
相同用户重复上传同一份计划书时复用现有任务或结果。
前端取消 180 秒本地假超时,改为串行轮询后端真实状态;网络波动不再误判解析失败。
增加服务重启后的过期任务恢复机制。
修复解析结果 JSON 序列化遗漏问题。
关键文件:
[extraction.py](D:/work/code/python/coding/baodanagent/api/insurance/ppt/extraction.py)
[tasks.py](D:/work/code/python/coding/baodanagent/api/insurance/poster/tasks.py)
[celery_tasks.py](D:/work/code/python/coding/baodanagent/api/insurance/generation/celery_tasks.py)
[service.py](D:/work/code/python/coding/baodanagent/api/insurance/poster/service.py)
[migrate_032.py](D:/work/code/python/coding/baodanagent/api/insurance/db/migrate_032.py)
[PosterSourcePanel.vue](D:/work/code/python/coding/baodanagent/frontend/src/components/poster/workspace/PosterSourcePanel.vue)
[回归测试](D:/work/code/python/coding/baodanagent/tests/ppt_poster_optimization_test.py)
This commit is contained in:
wsb1224 2026-08-01 03:15:43 +08:00
parent 9dc494d968
commit 9939188a4b
9 changed files with 692 additions and 197 deletions

View File

@ -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] 海报计划书解析任务字段迁移完成")

View File

@ -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 分钟软超时

View File

@ -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,

View File

@ -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()}

View File

@ -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()

View File

@ -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"(?<!\d){year}(?!\d)", content)
}
table_hits = sum(term.lower() in lowered for term in table_terms)
is_withdrawal = any(term.lower() in lowered for term in withdrawal_terms)
is_scenario = any(term.lower() in lowered for term in scenario_terms)
is_summary = any(term in content for term in (
"说明摘要", "說明摘要", "建议书摘要", "建議書摘要", "利益摘要",
)) or "illustration summary" in lowered
is_main_benefit = any(term in content for term in (
"说明 退保价值及身故赔偿", "說明 退保價值及身故賠償",
"退保价值及身故赔偿", "退保價值及身故賠償",
))
identity_hits = sum(term.lower() in lowered for term in identity_terms)
identity_score = identity_hits * 2 + (8 if is_summary else 0) + (2 if page_num <= 3 else 0)
benefit_score = (
table_hits * 3
+ len(years) * 2
+ (12 if is_summary else 0)
+ (8 if is_main_benefit else 0)
- (30 if is_withdrawal else 0)
- (20 if is_scenario else 0)
)
scored.append({
"page": page_num,
"content": content,
"years": years,
"identity_score": identity_score,
"benefit_score": benefit_score,
"table_hits": table_hits,
"excluded": is_withdrawal or is_scenario,
})
selected = []
selected_pages = set()
identity = max(scored, key=lambda item: item["identity_score"])
if identity["identity_score"] > 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

View File

@ -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}")

View File

@ -43,10 +43,16 @@
</div>
<!-- 上传中/解析中紧凑状态行 -->
<div v-else-if="parseStatus === 'uploading' || parseStatus === 'queued' || parseStatus === 'parsing'" class="parsing-row">
<div
v-else-if="parseStatus === 'uploading' || parseStatus === 'queued' || parseStatus === 'parsing'"
class="parsing-row"
role="status"
aria-live="polite"
>
<el-icon class="spin"><Loading /></el-icon>
<span class="parse-msg">{{ parseMessage || '解析中...' }}</span>
<el-progress :percentage="parseProgress" :stroke-width="4" :show-text="false" style="flex: 1" />
<span class="parse-percent">{{ parseProgress }}%</span>
</div>
<!-- 解析失败 -->
@ -179,7 +185,7 @@
</template>
<script setup lang="ts">
import { ref, computed, watch } from 'vue'
import { ref, computed, watch, onBeforeUnmount } from 'vue'
import { Document, ArrowDown, Upload, Loading, CircleCheckFilled } from '@element-plus/icons-vue'
import { ElMessage } from 'element-plus'
import { posterApi } from '@/utils/poster-api'
@ -303,17 +309,49 @@ function formatSize(bytes: number) {
return `${(bytes / 1048576).toFixed(1)}MB`
}
let pollTimer: ReturnType<typeof setInterval> | null = null
let pollRetries = 0
const MAX_POLL_RETRIES = 90
let pollTimer: ReturnType<typeof setTimeout> | null = null
let pollGeneration = 0
let consecutivePollErrors = 0
function stopPolling() {
pollGeneration++
if (pollTimer) {
clearInterval(pollTimer)
clearTimeout(pollTimer)
pollTimer = null
}
}
function applyCaseData(data: any) {
const status = data?.parseStatus || 'queued'
if (status === 'parsed' || status === 'partial') {
emit('update:parse', {
parseStatus: status,
parseProgress: 100,
parseMessage: data.parseMessage || '',
parseFailMessage: '',
parsedFields: data.parsedData || {},
})
if (data.parsedData) Object.assign(localFields.value, data.parsedData)
return true
}
if (status === 'failed') {
emit('update:parse', {
parseStatus: 'failed',
parseProgress: 100,
parseMessage: '',
parseFailMessage: data.parseError || '计划书解析失败,请重新解析或更换文件。',
})
return true
}
emit('update:parse', {
parseStatus: status,
parseProgress: Number(data?.parseProgress ?? 5),
parseMessage: data?.parseMessage || (status === 'queued' ? '任务排队中...' : '正在解析计划书...'),
parseFailMessage: '',
})
return false
}
async function handleFileChange(file: any) {
const raw = file.raw || file
emit('update:parse', {
@ -337,13 +375,11 @@ async function handleFileChange(file: any) {
caseUploadId: data?.id,
caseFileName: raw.name || '计划书.pdf',
caseFileSize: raw.size || 0,
parseStatus: data?.parseStatus || 'parsing',
parseProgress: 20,
parseMessage: '计划书已上传,正在解析...',
})
localPassword.value = ''
if (data?.id) {
startParsePolling(data.id)
const finished = applyCaseData(data)
if (!finished) startParsePolling(data.id)
}
} catch (e: any) {
emit('update:parse', {
@ -357,52 +393,36 @@ async function handleFileChange(file: any) {
function startParsePolling(id: number) {
stopPolling()
pollRetries = 0
pollTimer = setInterval(async () => {
pollRetries++
if (pollRetries > MAX_POLL_RETRIES) {
stopPolling()
emit('update:parse', {
parseStatus: 'failed',
parseMessage: '',
parseFailMessage: '解析超时,请重试',
})
return
}
const generation = pollGeneration
consecutivePollErrors = 0
const poll = async () => {
if (generation !== pollGeneration) return
try {
const res: any = await posterApi.getCaseUpload(id)
const data = res?.data
if (!data) return
const status = data.parseStatus
if (status === 'parsed' || status === 'partial') {
stopPolling()
emit('update:parse', {
parseStatus: status,
parseProgress: 100,
parseMessage: '',
parsedFields: data.parsedData || {},
})
if (data.parsedData) {
Object.assign(localFields.value, data.parsedData)
if (data) {
consecutivePollErrors = 0
if (applyCaseData(data)) {
stopPolling()
return
}
} else if (status === 'failed') {
stopPolling()
emit('update:parse', {
parseStatus: 'failed',
parseMessage: '',
parseFailMessage: '计划书解析失败,请重新上传。',
})
} else {
emit('update:parse', {
parseProgress: Math.min(90, props.parseProgress + 8),
parseMessage: status === 'queued' ? '任务排队中...' : '正在解析计划书...',
})
}
} catch {
//
consecutivePollErrors++
if (consecutivePollErrors >= 3) {
emit('update:parse', {
parseMessage: '网络连接不稳定,后台仍在解析,正在重新连接...',
})
}
} finally {
if (generation === pollGeneration) {
pollTimer = setTimeout(poll, 2000)
}
}
}, 2000)
}
pollTimer = setTimeout(poll, 600)
}
async function onConfirm() {
@ -438,14 +458,9 @@ async function onReparse() {
try {
const res: any = await posterApi.reparseCaseUpload(props.caseUploadId)
const data = res?.data
emit('update:parse', {
parseStatus: data?.parseStatus || 'queued',
parseProgress: 15,
parseMessage: '正在使用新版解析器重新识别...',
parsedFields: {},
dataConfirmed: false,
})
startParsePolling(props.caseUploadId)
emit('update:parse', { parsedFields: {}, dataConfirmed: false })
const finished = applyCaseData(data)
if (!finished) startParsePolling(props.caseUploadId)
} catch (e: any) {
ElMessage.error(e?.message || '重新解析失败')
}
@ -467,6 +482,8 @@ function onReset() {
localFields.value = emptyLocalFields()
editingFields.value = false
}
onBeforeUnmount(stopPolling)
</script>
<style scoped>
@ -567,11 +584,22 @@ function onReset() {
}
.parse-msg {
min-width: 0;
overflow: hidden;
text-overflow: ellipsis;
font-size: 13px;
color: var(--poster-muted);
white-space: nowrap;
}
.parse-percent {
flex-shrink: 0;
min-width: 32px;
color: var(--poster-muted);
font-size: 12px;
text-align: end;
}
.spin {
animation: spin 1s linear infinite;
flex-shrink: 0;

View File

@ -1,6 +1,8 @@
"""PPT/海报优化规则回归测试。"""
import asyncio
import sys
from pathlib import Path
from types import SimpleNamespace
sys.path.insert(0, str(Path(__file__).resolve().parents[1] / "api"))
@ -723,3 +725,100 @@ def test_poster_case_mapping_preserves_parse_diagnostics():
assert mapped["meta"]["method"] == "regex+ocr"
assert mapped["meta"]["lowQualityPages"] == [3, 8]
assert mapped["meta"]["provenance"]["insured.age"]["confidence"] == 0.8
def test_poster_page_selection_prefers_summary_and_excludes_withdrawal_scenarios():
from insurance.ppt.extraction import _select_poster_pages
pdf_text = "\n".join([
"[PAGE 1]\n受保人 年龄 性别 保单货币 年缴保费",
"[PAGE 2]\n款项提取说明 保单年度 退保价值 10 20 30",
"[PAGE 3]\n基本计划说明摘要 保单年度 退保价值 保证现金 1 5 10 15 20 25 30",
"[PAGE 4]\n悲观情景 保单年度 退保价值 10 20 30",
])
selected = _select_poster_pages(pdf_text, "savings")
assert "[PAGE 1]" in selected
assert "[PAGE 3]" in selected
assert "[PAGE 2]" not in selected
assert "[PAGE 4]" not in selected
def test_compact_poster_extraction_uses_one_call_and_keeps_target_years(monkeypatch):
from insurance.ppt import extraction
from insurance.ppt.llm_client import llm_client
pdf_text = (
"[PAGE 1]\n受保人 年龄 35 性别 男 保单货币 USD 年缴保费 100000\n"
"[PAGE 2]\n说明摘要 保单年度 退保价值 保证现金 10 20 30"
)
monkeypatch.setattr(extraction, "_extract_pdf_text", lambda _path: (pdf_text, []))
calls = []
async def fake_structured_output(**kwargs):
calls.append(kwargs["prompt"])
return {
"product_name": "测试储蓄计划",
"insured": {"age": 35, "gender": "male", "smoking_status": "non-smoker"},
"policy": {
"currency": "USD", "sum_insured": 500000,
"annual_premium": 100000, "first_year_amount_due": 92000,
"premium_payment_period": 5,
},
"benefit_illustration": [
{"policy_year": 10, "surrender_value": {"total": 300000}},
{"policy_year": 20, "total_surrender_value": 700000},
{"policy_year": 30, "total_surrender_value": 1200000},
{"policy_year": 40, "total_surrender_value": 1800000},
],
}, SimpleNamespace(latency_ms=120)
monkeypatch.setattr(llm_client, "structured_output", fake_structured_output)
progress = []
result = asyncio.run(extraction.ExtractionOrchestrator(use_cache=False).extract_for_poster(
"poster.pdf",
plan_type="savings",
progress_callback=lambda value, message: progress.append((value, message)),
))
assert len(calls) == 1
assert "提领、提款或压力情景" in calls[0]
assert [row["policy_year"] for row in result["benefit_illustration"]] == [10, 20, 30]
assert result["benefit_illustration"][0]["total_surrender_value"] == 300000
assert result["policy"]["annual_premium"] == 100000
assert result["policy"]["first_year_amount_due"] == 92000
assert result["_meta"]["method"] == "poster_compact"
assert progress[-1][0] == 90
def test_poster_source_polling_uses_backend_progress_without_local_timeout():
root = Path(__file__).resolve().parents[1]
source = (root / "frontend/src/components/poster/workspace/PosterSourcePanel.vue").read_text(encoding="utf-8")
assert "MAX_POLL_RETRIES" not in source
assert "setInterval" not in source
assert "data?.parseProgress" in source
assert "data.parseError" in source
assert "网络连接不稳定,后台仍在解析" in source
def test_poster_case_serializes_parsed_data_and_hides_unsafe_errors():
from insurance.models.poster_case_upload import PosterCaseUpload
record = PosterCaseUpload(
user_id="user-1",
product_id="product-1",
source_file_url="poster.pdf",
parse_status="failed",
parse_progress=100,
parse_error='Traceback: File "/app/insurance/tasks.py"',
parsed_data='{"annual_premium": 100000}',
product_snapshot_json='{"name": "测试产品"}',
)
data = record.to_dict()
assert data["parsedData"] == {"annual_premium": 100000}
assert data["productSnapshot"] == {"name": "测试产品"}
assert data["parseError"] == "解析失败,请重试;如多次失败请联系管理员"