diff --git a/README.md b/README.md index 4f77c55..2d9a9d3 100644 --- a/README.md +++ b/README.md @@ -19,6 +19,10 @@ - 检查已发布产品(模型/GPU/CPU) - 检查待审核产品 - 已存在则自动跳过,避免重复处理 +- ⚠️ **异常产品管理**:自动识别无法处理的产品 + - 内容库和互联网均无搜索结果时存入异常库 + - 支持人工排查和重新处理 + - 提供异常产品列表、详情、解决、删除 API ## 系统架构 @@ -322,7 +326,13 @@ BATCH_SIZE = 5 # 批量处理数量 ## 版本历史 -- v1.2.0 (2026-07-16): 产品状态检查 +- v1.17.0 (2026-07-17): 异常产品管理 + - 新增异常产品库,自动存储无法处理的产品 + - 内容库和互联网均无搜索结果时存入异常库 + - 新增异常产品 API(查询/详情/解决/删除/重试) + - 优化错误处理流程 + +- v1.16.0 (2026-07-16): 产品状态检查 - 新增产品状态检查功能 - 处理前检查产品是否已存在(已发布/待审核) - 避免重复处理已有产品 diff --git a/models/database.py b/models/database.py index be6f28b..69be7e1 100644 --- a/models/database.py +++ b/models/database.py @@ -204,9 +204,30 @@ class Database: ) ''') + # 异常产品表 + cursor.execute(''' + CREATE TABLE IF NOT EXISTS abnormal_products ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + product_name TEXT NOT NULL UNIQUE, + category TEXT, + subcategory TEXT, + abnormal_type TEXT DEFAULT 'no_search_results', + abnormal_reason TEXT, + search_results TEXT, + retry_count INTEGER DEFAULT 0, + status TEXT DEFAULT 'pending', + resolution TEXT, + resolved_at DATETIME, + resolved_by TEXT, + created_at DATETIME DEFAULT CURRENT_TIMESTAMP, + last_retry_at DATETIME + ) + ''') + # 创建索引 cursor.execute('CREATE INDEX IF NOT EXISTS idx_process_steps_session ON process_steps(process_id)') cursor.execute('CREATE INDEX IF NOT EXISTS idx_process_sessions_status ON process_sessions(status)') + cursor.execute('CREATE INDEX IF NOT EXISTS idx_abnormal_products_status ON abnormal_products(status)') conn.commit() @@ -830,5 +851,99 @@ class Database: intervention_data=intervention_data ) +# ========== 异常产品操作 ========== + def add_abnormal_product(self, product_name, category=None, subcategory=None, + abnormal_type='no_search_results', abnormal_reason=None, + search_results=None): + """添加异常产品""" + with self.get_connection() as conn: + cursor = conn.cursor() + try: + cursor.execute(''' + INSERT INTO abnormal_products + (product_name, category, subcategory, abnormal_type, abnormal_reason, search_results) + VALUES (?, ?, ?, ?, ?, ?) + ''', (product_name, category, subcategory, abnormal_type, abnormal_reason, + json.dumps(search_results, ensure_ascii=False) if search_results else None)) + conn.commit() + return cursor.lastrowid + except sqlite3.IntegrityError: + # 产品已存在,更新重试次数 + cursor.execute(''' + UPDATE abnormal_products + SET retry_count = retry_count + 1, + last_retry_at = CURRENT_TIMESTAMP, + abnormal_reason = ?, + search_results = ? + WHERE product_name = ? + ''', (abnormal_reason, + json.dumps(search_results, ensure_ascii=False) if search_results else None, + product_name)) + conn.commit() + return None + + def get_abnormal_products(self, limit=100, status='pending'): + """获取异常产品列表""" + with self.get_connection() as conn: + cursor = conn.cursor() + if status == 'all': + cursor.execute(''' + SELECT * FROM abnormal_products + ORDER BY created_at DESC + LIMIT ? + ''', (limit,)) + else: + cursor.execute(''' + SELECT * FROM abnormal_products + WHERE status = ? + ORDER BY created_at DESC + LIMIT ? + ''', (status, limit)) + return [dict(row) for row in cursor.fetchall()] + + def get_abnormal_count(self, status='pending'): + """获取异常产品数量""" + with self.get_connection() as conn: + cursor = conn.cursor() + if status == 'all': + cursor.execute('SELECT COUNT(*) FROM abnormal_products') + else: + cursor.execute('SELECT COUNT(*) FROM abnormal_products WHERE status = ?', (status,)) + return cursor.fetchone()[0] + + def resolve_abnormal_product(self, product_name, resolution, resolved_by='manual'): + """标记异常产品为已解决""" + with self.get_connection() as conn: + cursor = conn.cursor() + cursor.execute(''' + UPDATE abnormal_products + SET status = 'resolved', + resolution = ?, + resolved_by = ?, + resolved_at = CURRENT_TIMESTAMP + WHERE product_name = ? + ''', (resolution, resolved_by, product_name)) + conn.commit() + return cursor.rowcount > 0 + + def delete_abnormal_product(self, product_name): + """删除异常产品记录""" + with self.get_connection() as conn: + cursor = conn.cursor() + cursor.execute('DELETE FROM abnormal_products WHERE product_name = ?', (product_name,)) + conn.commit() + return cursor.rowcount > 0 + + def get_abnormal_product(self, product_name): + """获取异常产品详情""" + with self.get_connection() as conn: + cursor = conn.cursor() + cursor.execute('SELECT * FROM abnormal_products WHERE product_name = ?', (product_name,)) + row = cursor.fetchone() + result = dict(row) if row else None + if result and result.get('search_results'): + result['search_results'] = json.loads(result['search_results']) + return result + # 全局数据库实例 db = Database() \ No newline at end of file diff --git a/routes/products.py b/routes/products.py index 293b133..f7687bd 100644 --- a/routes/products.py +++ b/routes/products.py @@ -213,4 +213,122 @@ def process_batch(): 'success': True, 'processed': len(results), 'results': results - }) \ No newline at end of file + }) + +# ========== 异常产品 API ========== + +@bp.route('/abnormal', methods=['GET']) +def list_abnormal(): + """获取异常产品列表""" + limit = request.args.get('limit', 100, type=int) + status = request.args.get('status', 'pending') + + products = db.get_abnormal_products(limit=limit, status=status) + count = db.get_abnormal_count(status=status) + + # 解析 JSON 字段 + for item in products: + if item.get('search_results'): + try: + item['search_results'] = __import__('json').loads(item['search_results']) + except: + pass + + return jsonify({ + 'success': True, + 'products': products, + 'count': count + }) + + +@bp.route('/abnormal/', methods=['GET']) +def get_abnormal(product_name): + """获取异常产品详情""" + product = db.get_abnormal_product(product_name) + + if not product: + return jsonify({'error': '异常产品不存在'}), 404 + + return jsonify({ + 'success': True, + 'product': product + }) + + +@bp.route('/abnormal//resolve', methods=['POST']) +def resolve_abnormal(product_name): + """解决异常产品""" + data = request.get_json() + resolution = data.get('resolution', '人工处理完成') + resolved_by = data.get('resolved_by', 'manual') + + success = db.resolve_abnormal_product(product_name, resolution, resolved_by) + + if success: + return jsonify({ + 'success': True, + 'message': '异常产品已标记为已解决' + }) + else: + return jsonify({'error': '异常产品不存在'}), 404 + + +@bp.route('/abnormal/', methods=['DELETE']) +def delete_abnormal(product_name): + """删除异常产品记录""" + success = db.delete_abnormal_product(product_name) + + if success: + return jsonify({ + 'success': True, + 'message': '异常产品记录已删除' + }) + else: + return jsonify({'error': '异常产品不存在'}), 404 + + +@bp.route('/abnormal//retry', methods=['POST']) +def retry_abnormal(product_name): + """重试处理异常产品""" + # 获取异常产品详情 + abnormal = db.get_abnormal_product(product_name) + + if not abnormal: + return jsonify({'error': '异常产品不存在'}), 404 + + # 检查是否正在处理 + processing = db.get_processing_products() + if any(p['product_name'] == product_name for p in processing): + return jsonify({'error': '该产品正在处理中'}), 400 + + # 添加到处理中列表 + db.start_processing( + product_name=product_name, + category=abnormal.get('category'), + subcategory=abnormal.get('subcategory') + ) + + try: + # 执行处理 + result = process_service.process_product({ + 'product_name': product_name, + 'category': abnormal.get('category'), + 'subcategory': abnormal.get('subcategory') + }) + + # 如果处理成功,从异常库移除 + if result['success']: + db.resolve_abnormal_product( + product_name, + f"重试处理成功: {result['message']}", + 'auto_retry' + ) + + return jsonify({ + 'success': result['success'], + 'message': result['message'], + 'review_id': result.get('review_id') + }) + finally: + # 完成处理,从处理中列表移除 + db.finish_processing(product_name) diff --git a/services/process_service.py b/services/process_service.py index ac79c36..c8c0813 100644 --- a/services/process_service.py +++ b/services/process_service.py @@ -84,9 +84,38 @@ class DataProcessService: ) if search_results['total'] == 0: - # 没有找到数据,发送通知 - paramhub_client.send_notification(f"未找到产品 '{product_name}' 的相关数据") - result['message'] = '未找到相关数据' + # 没有找到数据,存入异常产品库 + print(f"[异常] 未找到产品 '{product_name}' 的相关数据,存入异常库") + + # 添加到异常产品库 + db.add_abnormal_product( + product_name=product_name, + category=category, + subcategory=subcategory, + abnormal_type='no_search_results', + abnormal_reason='内容库和互联网均未搜索到相关数据', + search_results=search_results + ) + + # 发送通知 + paramhub_client.send_notification( + f"⚠️ 产品 '{product_name}' 未找到相关数据\n" + f"已存入异常产品库,请人工排查" + ) + + # 记录处理历史 + db.add_process_history( + product_name=product_name, + category=category, + subcategory=subcategory, + status='abnormal', + details={ + 'reason': 'no_search_results', + 'message': '内容库和互联网均未搜索到相关数据' + } + ) + + result['message'] = f"产品 '{product_name}' 未找到相关数据,已存入异常库" return result # 2. 提取对应产品的具体内容(排除无关产品)