代码优化

This commit is contained in:
liujianjiang 2025-11-28 18:32:41 +08:00
parent ef48cec1d5
commit 7de0c8c265
4 changed files with 13 additions and 31 deletions

View File

@ -1,16 +0,0 @@
# -*- coding: utf-8 -*-
from typing import Dict, Any
from public_function.redis_task_manager import RedisTaskManager
class CrawlerManagement:
def __init__(self, config_data: Dict[str, Any]):
self.config_data = config_data
self.redis_conn = RedisTaskManager(self.config_data)
def get_task_item(self, task_data: Dict[str, Any]):
redis_key = f"{task_data["region"].lower()}_{task_data['app_name'].lower()}"
result = self.redis_conn.get_random_task_and_delete(redis_key)
if result:
return result
return None

17
main.py
View File

@ -38,17 +38,6 @@ def get_account_manager():
return None
def get_crawler_manager():
"""任务管理器实例"""
config = get_config()
try:
from crawler_management.crawler_management import CrawlerManagement
return CrawlerManagement(config)
except ImportError:
print(f"未找到AllTask类返回模拟实例")
return None
def get_task_manager():
"""任务获取管理器实例"""
config = get_config()
@ -136,7 +125,7 @@ async def update_account(account_data: AccountUpdate, account_manager: Any = Dep
@app.post("/get_crawler_task", summary="获取任务")
async def get_crawler_task(task_data: CrawlerItem, task_manager: Any = Depends(get_crawler_manager)):
async def get_crawler_task(task_data: CrawlerItem, task_manager: Any = Depends(get_task_manager)):
"""
获取指定应用的可用账号
- **app_name**: 应用名称
@ -146,6 +135,8 @@ async def get_crawler_task(task_data: CrawlerItem, task_manager: Any = Depends(g
params = task_data.model_dump()
result = task_manager.get_task_item(params)
if result:
data = {"task_id": result["task_id"], "status": 4}
await task_manager.update_task_record(data)
return {"code": 200, "message": "任务获取成功", "data": result}
raise HTTPException(status_code=404, detail="队列暂时没有任务,请等待一段时间后重新尝试")
except Exception as e:
@ -190,8 +181,6 @@ async def get_goods_info(task_data: GoodsInfo, task_manager: Any = Depends(get_t
raise HTTPException(status_code=503, detail="抓取商品数据失败,请重新尝试")
except Exception as e:
print(f"任务异常 - ID: {task_id}, 错误: {str(e)}")
data = {"task_id": params["task_id"], "status": 3}
await task_manager.update_task_record(data)
raise HTTPException(status_code=503, detail="抓取商品数据失败,请重新尝试")

View File

@ -7,7 +7,7 @@ create table crawler_task_record_info
item_id VARCHAR(50) NOT NULL COMMENT '商品ID',
region VARCHAR(50) NOT NULL COMMENT '客户所在地区',
token VARCHAR(128) NOT NULL COMMENT '用户标识',
status int NOT NULL DEFAULT 1 COMMENT '任务状态1、开始执行2、执行成功3、执行失败',
status int NOT NULL DEFAULT 1 COMMENT '任务状态1、开始执行2、执行成功3、提交到redis队列4、从redis消耗',
create_time DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP COMMENT '创建时间',
update_time DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP COMMENT '更新时间',
PRIMARY KEY (id),

View File

@ -16,6 +16,13 @@ class AllTask:
self.redis_conn = RedisTaskManager(self.config_data)
self.db_pool: Optional[AsyncMySQL] = AsyncMySQL(self.config_data["advert_policy"])
def get_task_item(self, task_data: Dict[str, Any]):
redis_key = f"{task_data["region"].lower()}_{task_data['app_name'].lower()}"
result = self.redis_conn.get_random_task_and_delete(redis_key)
if result:
return result
return None
async def deal_shopee_task(self, param: Dict[str, Any]):
"""
处理Shopee任务的异步方法
@ -32,6 +39,8 @@ class AllTask:
print(f"{key}:{param['shop_id']}:{param['item_id']} 在redis中获取数据失败将任务提交到队列")
# 确保add_task_to_set是异步方法或正确处理
self.redis_conn.add_task_to_set(task_data=param)
update_parms = {"task_id": param["task_id"], "status": 3}
await self.update_task_record(update_parms)
# 任务结束后开始等待
endtime = time.time() + 55
counter = 0