This commit is contained in:
HuangHai
2026-01-13 19:19:55 +08:00
parent 11ae4b5abf
commit 9fdbdaa338
3 changed files with 65 additions and 10 deletions

View File

@@ -1,8 +1,9 @@
# 采集配置
SCROLL_DISTANCE_RATIO = 0.6
MAX_SCROLLS = 100
MAX_CRAWL_DISTANCE = 50
# 最大采集场站数量,达到此次数后停止采集
MAX_STATIONS_COUNT = 100
# 场站去重过期时间(秒),在此时间内重复出现的场站不会再次点击进入详情页
REDIS_STATION_EXPIRE = 120
DATA_RETENTION_DAYS = 365

View File

@@ -21,9 +21,8 @@ from Apps.AiTeJiYiChong.Service import AiTeJiYiChongService
from Config.Config import TEMP_IMAGE_DIR
from Apps.AiTeJiYiChong.Config.Setting import (
SCROLL_DISTANCE_RATIO,
MAX_SCROLLS, REDIS_STATION_EXPIRE,
MAX_STATIONS_COUNT, REDIS_STATION_EXPIRE,
WAIT_AFTER_SCROLL,
MAX_CRAWL_DISTANCE,
SAFE_EXCLUDE_RATIO,
BOTTOM_SAFE_EXCLUDE_RATIO
)
@@ -43,7 +42,7 @@ async def is_already_crawled(redis_kit, station_name):
return True
return False
async def get_station_list(d, service, max_scrolls=MAX_SCROLLS):
async def get_station_list(d, service, max_stations_count=MAX_STATIONS_COUNT):
"""
获取场站列表并处理翻页
"""
@@ -55,15 +54,18 @@ async def get_station_list(d, service, max_scrolls=MAX_SCROLLS):
device_info['width'] = w
device_info['height'] = h
logger.info(f"开始爬取列表,设备: {device_info.get('productName')} | 分辨率: {w}x{h}")
logger.info(f"开始爬取列表,设备: {device_info.get('productName')} | 分辨率: {w}x{h} | 目标数量: {max_stations_count}")
# 用于追踪后台分析任务
background_tasks = []
last_list_md5 = None
no_new_data_count = 0
total_processed_count = 0
scroll_count = 0
for i in range(max_scrolls + 1):
logger.info(f"正在处理第 {i + 1} 页...")
while total_processed_count < max_stations_count:
scroll_count += 1
logger.info(f"正在处理第 {scroll_count} 次滚动 (已采集: {total_processed_count}/{max_stations_count})...")
# 1. 拍摄截图
image_uuid = str(uuid.uuid4())
@@ -127,8 +129,9 @@ async def get_station_list(d, service, max_scrolls=MAX_SCROLLS):
# 正常处理新场站
click_x, click_y = card["click_point"]
logger.info(f"准备处理第 {card_idx + 1} 个场站: {station_name}, 点击坐标: ({click_x}, {click_y})")
logger.info(f">>> 发现新场站 '{station_name}',开始处理... ({total_processed_count + 1}/{max_stations_count})")
new_stations_processed += 1
total_processed_count += 1
d.click(click_x, click_y)
# 等待二级页面加载
@@ -230,6 +233,13 @@ async def get_station_list(d, service, max_scrolls=MAX_SCROLLS):
full_name_key = f"crawled:aite:{Kit.clean_station_name(station_name)}"
await redis_kit.set_data(full_name_key, "1", expire=REDIS_STATION_EXPIRE)
# 检查是否已达到最大采集数量
if total_processed_count >= max_stations_count:
logger.info(f"已达到目标采集数量 {max_stations_count},准备结束采集。")
break
if total_processed_count >= max_stations_count:
break
# 在每一页结束时,清理已完成的任务
done_tasks = [t for t in background_tasks if t.done()]
for t in done_tasks:
@@ -258,7 +268,7 @@ async def get_station_list(d, service, max_scrolls=MAX_SCROLLS):
logger.info(f"正在等待剩余 {len(background_tasks)} 个后台分析任务完成...")
await asyncio.gather(*background_tasks, return_exceptions=True)
logger.info("达到最大翻页次数,爬取结束")
logger.info(f"采集任务完成,共采集 {total_processed_count} 个场站")
return True
async def analyze_prices_background(service, station_name, image_paths):

View File

@@ -1,6 +1,8 @@
import os
import logging
import re
import asyncio
import functools
from sqlalchemy.ext.asyncio import create_async_engine, AsyncSession
from sqlalchemy.orm import sessionmaker
from sqlalchemy.sql import text
@@ -14,6 +16,42 @@ logging.basicConfig(level=logging.INFO, format='%(asctime)s - %(name)s - %(level
logger = logging.getLogger(__name__)
def db_retry(max_retries=3, delay=1):
"""
数据库操作重试装饰器
Args:
max_retries: 最大重试次数默认为3次
delay: 重试间隔时间默认为1秒
"""
def decorator(func):
@functools.wraps(func)
async def wrapper(*args, **kwargs):
# 检查是否提供了 session如果提供了 session通常意味着是在一个更大的事务中不建议在这里重试
session = kwargs.get('session')
if session is not None:
return await func(*args, **kwargs)
last_exception = None
for attempt in range(max_retries):
try:
return await func(*args, **kwargs)
except Exception as e:
last_exception = e
# 只有在还有重试机会时才打印警告
if attempt < max_retries - 1:
logger.warning(f"数据库操作失败 (尝试 {attempt + 1}/{max_retries}),正在进行重试: {str(e)}")
await asyncio.sleep(delay)
else:
logger.error(f"数据库操作在 {max_retries} 次尝试后仍然失败: {str(e)}")
# 如果循环结束仍未返回,说明最后一次尝试也失败了,抛出异常
if last_exception:
raise last_exception
return wrapper
return decorator
class Db:
"""通用数据库操作封装类,提供数据库连接和操作功能"""
# 单例实例
@@ -341,6 +379,7 @@ class Db:
return processed_sql
@db_retry()
async def find(self, sql, params=None, session=None):
"""
执行SQL查询并返回结果异步版本
@@ -783,6 +822,7 @@ class Db:
# 如果所有策略都失败,返回默认查询
return "SELECT COUNT(*)"
@db_retry()
async def execute_update(self, sql, params=None, session=None):
"""执行SQL更新操作插入、更新、删除异步版本
@@ -890,6 +930,7 @@ class Db:
if is_own_session and session:
await session.close()
@db_retry()
async def save(self, table_name, data, primary_key, session=None):
"""插入数据到指定表,并返回插入后的主键值(异步版本)
@@ -948,6 +989,7 @@ class Db:
if is_own_session and session:
await session.close()
@db_retry()
async def update(self, table_name, data, primary_key, session=None):
"""根据主键更新指定表中的数据(异步版本)
@@ -1013,6 +1055,7 @@ class Db:
if is_own_session and session:
await session.close()
@db_retry()
async def batch_insert(self, table_name, data_list, primary_key=None, session=None):
"""
批量插入数据到指定表(异步版本)
@@ -1092,6 +1135,7 @@ class Db:
if is_own_session and session:
await session.close()
@db_retry()
async def batch_update(self, table_name, data_list, primary_key, session=None):
"""
批量更新指定表中的数据(异步版本)