2026-07-14 12:23:49 +08:00
|
|
|
"""
|
|
|
|
|
处理步骤监控服务 - 记录和监控产品处理流程
|
|
|
|
|
"""
|
|
|
|
|
import time
|
|
|
|
|
import uuid
|
|
|
|
|
import json
|
|
|
|
|
import threading
|
|
|
|
|
import logging
|
|
|
|
|
from datetime import datetime
|
|
|
|
|
from models.database import db
|
|
|
|
|
from services.search_service import search_service
|
|
|
|
|
from services.paramhub_client import paramhub_client
|
|
|
|
|
|
|
|
|
|
logger = logging.getLogger('process_monitor')
|
|
|
|
|
|
|
|
|
|
# 处理步骤定义
|
|
|
|
|
PROCESS_STEPS = [
|
|
|
|
|
{'num': 1, 'name': '搜索内容库', 'description': '从内容库搜索相关文章'},
|
|
|
|
|
{'num': 2, 'name': '搜索互联网', 'description': '从互联网搜索最新数据'},
|
|
|
|
|
{'num': 3, 'name': '抓取网页内容', 'description': '抓取搜索结果网页的详细内容'},
|
|
|
|
|
{'num': 4, 'name': '提取产品数据', 'description': '从抓取内容中提取产品相关数据'},
|
|
|
|
|
{'num': 5, 'name': '填充字段', 'description': '根据分类字段配置填充数据'},
|
|
|
|
|
{'num': 6, 'name': '提交审核', 'description': '提交到ParamHub待审核区'},
|
|
|
|
|
]
|
|
|
|
|
|
|
|
|
|
class ProcessMonitor:
|
|
|
|
|
"""处理步骤监控器"""
|
|
|
|
|
|
|
|
|
|
def __init__(self):
|
|
|
|
|
self.active_sessions = {}
|
|
|
|
|
self.step_timers = {}
|
|
|
|
|
|
|
|
|
|
def create_session_id(self):
|
|
|
|
|
"""生成会话ID"""
|
|
|
|
|
return f"proc_{datetime.now().strftime('%Y%m%d_%H%M%S')}_{uuid.uuid4().hex[:8]}"
|
|
|
|
|
|
|
|
|
|
def start_process(self, product_name, category=None, subcategory=None):
|
|
|
|
|
"""启动产品处理流程"""
|
|
|
|
|
session_id = self.create_session_id()
|
|
|
|
|
|
|
|
|
|
# 创建会话记录
|
|
|
|
|
db.create_process_session(session_id, product_name, category, subcategory)
|
|
|
|
|
|
|
|
|
|
# 初始化控制信息
|
|
|
|
|
self.active_sessions[session_id] = {
|
|
|
|
|
'paused': False,
|
|
|
|
|
'stop': False,
|
|
|
|
|
'current_step': 0
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
# 启动后台线程处理
|
|
|
|
|
thread = threading.Thread(
|
|
|
|
|
target=self._run_process,
|
|
|
|
|
args=(session_id, product_name, category, subcategory),
|
|
|
|
|
daemon=True
|
|
|
|
|
)
|
|
|
|
|
thread.start()
|
|
|
|
|
|
|
|
|
|
logger.info(f"启动处理会话: {session_id}, 产品: {product_name}")
|
|
|
|
|
return session_id
|
|
|
|
|
|
|
|
|
|
def _run_process(self, session_id, product_name, category, subcategory):
|
|
|
|
|
"""执行处理流程"""
|
|
|
|
|
try:
|
|
|
|
|
db.update_session_status(session_id, 'running')
|
|
|
|
|
|
|
|
|
|
result = {'success': False, 'message': '', 'review_id': None}
|
|
|
|
|
|
|
|
|
|
all_data = {
|
|
|
|
|
'library_results': [],
|
|
|
|
|
'internet_results': [],
|
|
|
|
|
'fetched_contents': [],
|
|
|
|
|
'extracted_data': None,
|
|
|
|
|
'filled_data': None
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
# 步骤1: 搜索内容库
|
|
|
|
|
if not self._check_pause(session_id):
|
|
|
|
|
self._start_step(session_id, product_name, 1, '搜索内容库')
|
|
|
|
|
try:
|
|
|
|
|
articles = db.search_articles(product_name, category)
|
|
|
|
|
all_data['library_results'] = articles
|
|
|
|
|
self._complete_step(session_id, 1, {'count': len(articles)})
|
|
|
|
|
logger.info(f"[{session_id}] 步骤1完成: 找到 {len(articles)} 篇文章")
|
|
|
|
|
except Exception as e:
|
|
|
|
|
self._fail_step(session_id, 1, str(e))
|
|
|
|
|
result['message'] = f'搜索内容库失败: {e}'
|
|
|
|
|
|
|
|
|
|
# 步骤2: 搜索互联网
|
|
|
|
|
if not self._check_pause(session_id) and not result.get('message'):
|
|
|
|
|
self._start_step(session_id, product_name, 2, '搜索互联网')
|
|
|
|
|
try:
|
|
|
|
|
internet_results = search_service.search_internet(product_name, max_results=10)
|
|
|
|
|
all_data['internet_results'] = internet_results
|
|
|
|
|
self._complete_step(session_id, 2, {'count': len(internet_results)})
|
|
|
|
|
logger.info(f"[{session_id}] 步骤2完成: 找到 {len(internet_results)} 条结果")
|
|
|
|
|
except Exception as e:
|
|
|
|
|
self._complete_step(session_id, 2, {'count': 0, 'error': str(e)})
|
|
|
|
|
|
|
|
|
|
# 步骤3: 抓取网页内容
|
|
|
|
|
if not self._check_pause(session_id) and all_data['internet_results']:
|
|
|
|
|
self._start_step(session_id, product_name, 3, '抓取网页内容')
|
|
|
|
|
try:
|
|
|
|
|
fetched = []
|
2026-07-14 15:36:18 +08:00
|
|
|
failed_count = 0
|
2026-07-14 12:23:49 +08:00
|
|
|
urls_to_fetch = [r['url'] for r in all_data['internet_results'][:5]]
|
|
|
|
|
|
|
|
|
|
for i, url in enumerate(urls_to_fetch):
|
|
|
|
|
if self._check_pause(session_id):
|
|
|
|
|
break
|
|
|
|
|
|
|
|
|
|
fetch_result = search_service.fetch_url_content(url)
|
|
|
|
|
if fetch_result.get('success'):
|
2026-07-14 15:13:31 +08:00
|
|
|
title = fetch_result.get('title', '')
|
|
|
|
|
content = fetch_result.get('content', '')
|
2026-07-14 12:23:49 +08:00
|
|
|
fetched.append({
|
|
|
|
|
'url': url,
|
2026-07-14 15:13:31 +08:00
|
|
|
'title': title,
|
|
|
|
|
'content': content[:500]
|
2026-07-14 12:23:49 +08:00
|
|
|
})
|
2026-07-14 15:13:31 +08:00
|
|
|
|
|
|
|
|
# 保存到内容库
|
|
|
|
|
try:
|
|
|
|
|
existing = db.search_articles(url)
|
|
|
|
|
if not any(a.get('url') == url for a in existing):
|
|
|
|
|
db.add_article(
|
|
|
|
|
product_names=[],
|
|
|
|
|
category=category or '',
|
|
|
|
|
keywords=[],
|
|
|
|
|
summary=content[:200] if content else '',
|
|
|
|
|
content=content,
|
|
|
|
|
source=url,
|
|
|
|
|
url=url,
|
|
|
|
|
search_title=title
|
|
|
|
|
)
|
|
|
|
|
logger.info(f"[{session_id}] 已保存到内容库: {title[:30]}")
|
|
|
|
|
except Exception as save_error:
|
|
|
|
|
logger.warning(f"[{session_id}] 保存内容库失败: {save_error}")
|
2026-07-14 15:36:18 +08:00
|
|
|
else:
|
|
|
|
|
# 记录失败URL
|
|
|
|
|
failed_count += 1
|
|
|
|
|
error_msg = fetch_result.get('error', '抓取失败')
|
|
|
|
|
try:
|
|
|
|
|
db.add_failed_url(url, product_name, error_msg, source='process_monitor')
|
|
|
|
|
logger.warning(f"[{session_id}] 抓取失败,已记录: {url}")
|
|
|
|
|
except Exception as e:
|
|
|
|
|
logger.error(f"[{session_id}] 记录失败URL出错: {e}")
|
2026-07-14 15:13:31 +08:00
|
|
|
|
2026-07-14 12:23:49 +08:00
|
|
|
time.sleep(0.3)
|
|
|
|
|
|
|
|
|
|
all_data['fetched_contents'] = fetched
|
2026-07-14 15:36:18 +08:00
|
|
|
self._complete_step(session_id, 3, {'count': len(fetched), 'failed': failed_count})
|
|
|
|
|
logger.info(f"[{session_id}] 步骤3完成: 抓取 {len(fetched)} 个网页, 失败 {failed_count} 个")
|
2026-07-14 12:23:49 +08:00
|
|
|
except Exception as e:
|
|
|
|
|
self._fail_step(session_id, 3, str(e))
|
|
|
|
|
|
|
|
|
|
# 步骤4: 提取产品数据
|
|
|
|
|
if not self._check_pause(session_id):
|
|
|
|
|
self._start_step(session_id, product_name, 4, '提取产品数据')
|
|
|
|
|
try:
|
|
|
|
|
extracted = self._extract_data(product_name, all_data)
|
|
|
|
|
all_data['extracted_data'] = extracted
|
|
|
|
|
|
|
|
|
|
if extracted:
|
|
|
|
|
self._complete_step(session_id, 4, {'has_data': True})
|
|
|
|
|
else:
|
|
|
|
|
self._complete_step(session_id, 4, {'has_data': False}, status='skipped')
|
|
|
|
|
result['message'] = '无法提取有效数据'
|
|
|
|
|
except Exception as e:
|
|
|
|
|
self._fail_step(session_id, 4, str(e))
|
|
|
|
|
|
|
|
|
|
# 步骤5: 填充字段
|
|
|
|
|
if not self._check_pause(session_id) and all_data['extracted_data']:
|
|
|
|
|
self._start_step(session_id, product_name, 5, '填充字段')
|
|
|
|
|
try:
|
|
|
|
|
filled = self._fill_fields(all_data['extracted_data'], category, subcategory)
|
|
|
|
|
all_data['filled_data'] = filled
|
|
|
|
|
|
|
|
|
|
if filled:
|
|
|
|
|
self._complete_step(session_id, 5, {'filled': True})
|
|
|
|
|
else:
|
|
|
|
|
self._fail_step(session_id, 5, '填充数据失败')
|
|
|
|
|
except Exception as e:
|
|
|
|
|
self._fail_step(session_id, 5, str(e))
|
|
|
|
|
|
|
|
|
|
# 步骤6: 提交审核
|
|
|
|
|
if not self._check_pause(session_id) and all_data['filled_data']:
|
|
|
|
|
self._start_step(session_id, product_name, 6, '提交审核')
|
|
|
|
|
try:
|
|
|
|
|
category_type = self._get_category_type(category)
|
|
|
|
|
success, review_id_or_error = paramhub_client.submit_for_review(
|
|
|
|
|
category_type,
|
|
|
|
|
all_data['filled_data'],
|
|
|
|
|
subcategory
|
|
|
|
|
)
|
|
|
|
|
|
|
|
|
|
if success:
|
|
|
|
|
self._complete_step(session_id, 6, {'review_id': review_id_or_error})
|
|
|
|
|
result['success'] = True
|
|
|
|
|
result['review_id'] = review_id_or_error
|
|
|
|
|
|
|
|
|
|
db.update_session_status(session_id, 'completed',
|
|
|
|
|
review_id=review_id_or_error,
|
|
|
|
|
result=json.dumps(result, ensure_ascii=False))
|
|
|
|
|
|
|
|
|
|
db.add_process_history(
|
|
|
|
|
product_name=product_name,
|
|
|
|
|
category=category,
|
|
|
|
|
subcategory=subcategory,
|
|
|
|
|
status='submitted',
|
|
|
|
|
review_id=review_id_or_error,
|
|
|
|
|
details=all_data
|
|
|
|
|
)
|
|
|
|
|
logger.info(f"[{session_id}] 步骤6完成: 提交成功")
|
|
|
|
|
else:
|
|
|
|
|
self._fail_step(session_id, 6, review_id_or_error)
|
|
|
|
|
db.update_session_status(session_id, 'failed')
|
|
|
|
|
except Exception as e:
|
|
|
|
|
self._fail_step(session_id, 6, str(e))
|
|
|
|
|
|
|
|
|
|
# 清理
|
|
|
|
|
if session_id in self.active_sessions:
|
|
|
|
|
del self.active_sessions[session_id]
|
|
|
|
|
|
|
|
|
|
return result
|
|
|
|
|
|
|
|
|
|
except Exception as e:
|
|
|
|
|
logger.error(f"处理会话异常: {session_id} - {e}")
|
|
|
|
|
db.update_session_status(session_id, 'failed')
|
|
|
|
|
return {'success': False, 'message': str(e)}
|
|
|
|
|
|
|
|
|
|
def _start_step(self, session_id, product_name, step_num, step_name):
|
|
|
|
|
"""开始步骤"""
|
|
|
|
|
db.update_session_status(session_id, 'running', current_step=step_num)
|
|
|
|
|
db.add_process_step(session_id, product_name, step_num, step_name)
|
|
|
|
|
if session_id not in self.step_timers:
|
|
|
|
|
self.step_timers[session_id] = {}
|
|
|
|
|
self.step_timers[session_id][step_num] = time.time()
|
|
|
|
|
|
|
|
|
|
def _complete_step(self, session_id, step_num, step_data=None, status='completed'):
|
|
|
|
|
"""完成步骤"""
|
|
|
|
|
duration_ms = None
|
|
|
|
|
if session_id in self.step_timers and step_num in self.step_timers[session_id]:
|
|
|
|
|
duration_ms = int((time.time() - self.step_timers[session_id][step_num]) * 1000)
|
|
|
|
|
db.update_step_status(session_id, step_num, status, step_data=step_data, duration_ms=duration_ms)
|
|
|
|
|
|
|
|
|
|
def _fail_step(self, session_id, step_num, error_message):
|
|
|
|
|
"""步骤失败"""
|
|
|
|
|
duration_ms = None
|
|
|
|
|
if session_id in self.step_timers and step_num in self.step_timers[session_id]:
|
|
|
|
|
duration_ms = int((time.time() - self.step_timers[session_id][step_num]) * 1000)
|
|
|
|
|
db.update_step_status(session_id, step_num, 'failed', error_message=error_message, duration_ms=duration_ms)
|
|
|
|
|
db.update_session_status(session_id, 'failed')
|
|
|
|
|
|
|
|
|
|
def _check_pause(self, session_id):
|
|
|
|
|
"""检查是否暂停"""
|
|
|
|
|
if session_id not in self.active_sessions:
|
|
|
|
|
return False
|
|
|
|
|
|
|
|
|
|
session = self.active_sessions[session_id]
|
|
|
|
|
|
|
|
|
|
if session.get('stop'):
|
|
|
|
|
return True
|
|
|
|
|
|
|
|
|
|
while session.get('paused'):
|
|
|
|
|
time.sleep(0.5)
|
|
|
|
|
if session.get('stop'):
|
|
|
|
|
return True
|
|
|
|
|
|
|
|
|
|
return False
|
|
|
|
|
|
|
|
|
|
def pause_session(self, session_id):
|
|
|
|
|
"""暂停会话"""
|
|
|
|
|
if session_id in self.active_sessions:
|
|
|
|
|
self.active_sessions[session_id]['paused'] = True
|
|
|
|
|
db.pause_session(session_id, '用户暂停')
|
|
|
|
|
return True
|
|
|
|
|
return False
|
|
|
|
|
|
|
|
|
|
def resume_session(self, session_id):
|
|
|
|
|
"""继续会话"""
|
|
|
|
|
if session_id in self.active_sessions:
|
|
|
|
|
self.active_sessions[session_id]['paused'] = False
|
|
|
|
|
db.resume_session(session_id)
|
|
|
|
|
return True
|
|
|
|
|
return False
|
|
|
|
|
|
|
|
|
|
def stop_session(self, session_id):
|
|
|
|
|
"""停止会话"""
|
|
|
|
|
if session_id in self.active_sessions:
|
|
|
|
|
self.active_sessions[session_id]['stop'] = True
|
|
|
|
|
self.active_sessions[session_id]['paused'] = False
|
|
|
|
|
db.update_session_status(session_id, 'stopped')
|
|
|
|
|
return True
|
|
|
|
|
return False
|
|
|
|
|
|
|
|
|
|
def get_session_status(self, session_id):
|
|
|
|
|
"""获取会话状态"""
|
|
|
|
|
session = db.get_process_session(session_id)
|
|
|
|
|
if session:
|
|
|
|
|
steps = db.get_process_steps(session_id)
|
|
|
|
|
return {'session': session, 'steps': steps}
|
|
|
|
|
return None
|
|
|
|
|
|
|
|
|
|
def _extract_data(self, product_name, all_data):
|
|
|
|
|
"""提取产品数据"""
|
|
|
|
|
all_content = []
|
|
|
|
|
|
|
|
|
|
for article in all_data.get('library_results', []):
|
|
|
|
|
content = article.get('content', '')
|
|
|
|
|
if content:
|
|
|
|
|
all_content.append(content)
|
|
|
|
|
|
|
|
|
|
for item in all_data.get('fetched_contents', []):
|
|
|
|
|
content = item.get('content', '')
|
|
|
|
|
if content:
|
|
|
|
|
all_content.append(content)
|
|
|
|
|
|
|
|
|
|
if not all_content:
|
|
|
|
|
return None
|
|
|
|
|
|
|
|
|
|
return {
|
|
|
|
|
'name': product_name,
|
|
|
|
|
'raw_content': '\n---\n'.join(all_content[:3])
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
def _fill_fields(self, extracted_data, category, subcategory):
|
|
|
|
|
"""填充字段"""
|
|
|
|
|
if not extracted_data:
|
|
|
|
|
return None
|
|
|
|
|
|
|
|
|
|
import re
|
|
|
|
|
|
|
|
|
|
filled = {
|
|
|
|
|
'name': extracted_data['name'],
|
|
|
|
|
'visible': True,
|
|
|
|
|
'is_pinned': False
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
content = extracted_data.get('raw_content', '')
|
|
|
|
|
|
|
|
|
|
params_match = re.search(r'(\d+(?:\.\d+)?)\s*[Bb]', content)
|
|
|
|
|
if params_match:
|
|
|
|
|
filled['parameters'] = f"{params_match.group(1)}B"
|
|
|
|
|
|
|
|
|
|
date_match = re.search(r'(\d{4}[-/]\d{1,2}[-/]\d{1,2})', content)
|
|
|
|
|
if date_match:
|
|
|
|
|
filled['publish_date'] = date_match.group(1).replace('/', '-')
|
|
|
|
|
|
|
|
|
|
filled['_source'] = 'auto_manager'
|
|
|
|
|
filled['_extracted_at'] = datetime.now().isoformat()
|
|
|
|
|
|
|
|
|
|
return filled
|
|
|
|
|
|
|
|
|
|
def _get_category_type(self, category):
|
|
|
|
|
"""获取分类类型"""
|
|
|
|
|
if not category:
|
|
|
|
|
return 'dynamic'
|
|
|
|
|
|
|
|
|
|
category_lower = category.lower()
|
|
|
|
|
if 'model' in category_lower or 'ai' in category_lower:
|
|
|
|
|
return 'model'
|
|
|
|
|
elif 'gpu' in category_lower:
|
|
|
|
|
return 'gpu'
|
|
|
|
|
elif 'cpu' in category_lower:
|
|
|
|
|
return 'cpu'
|
|
|
|
|
return 'dynamic'
|
|
|
|
|
|
|
|
|
|
# 全局处理监控实例
|
|
|
|
|
process_monitor = ProcessMonitor()
|