Bladeren bron

fix(zj_new_daily_spider): 修复类型错误并优化任务循环逻辑

- 修正resp_json中obj_order_rating_goods访问方式,从列表改为字典
- 替换inspect模块为sys._getframe()获取当前函数名,简化代码
- 增加MySQL连接池健康检查,提升数据库连接稳定性
- 优化主循环逻辑,分批次处理state=0任务,防止死循环
- 添加日志记录每批处理数量及异常信息,便于监控
- 暂停定时调度任务,改为直接调用主函数执行
charley 2 weken geleden
bovenliggende
commit
9a8bd7c44c
1 gewijzigde bestanden met toevoegingen van 27 en 23 verwijderingen
  1. 27 23
      zhongjian_spider/zj_new_daily_spider.py

+ 27 - 23
zhongjian_spider/zj_new_daily_spider.py

@@ -2,9 +2,9 @@
 # Author : Charley
 # Python : 3.10.8
 # Date   : 2025/6/9 15:56
+import sys
 import time
 import json
-import inspect
 import requests
 import schedule
 import hashlib
@@ -130,7 +130,7 @@ def parse_data(resp_json, sql_pool):
     order_no = resp_json.get('obj_order_rating_goods', {}).get('order_no')
     tag_no = resp_json.get('obj_order_rating_goods', {}).get('tag_no')  # 标签号/查询的号码
 
-    images = resp_json.get('obj_order_rating_goods', []).get('images')
+    images = resp_json.get('obj_order_rating_goods', {}).get('images')
     card_create_time = resp_json.get('obj_order_rating_goods', {}).get('create_time')
     card_update_time = resp_json.get('obj_order_rating_goods', {}).get('update_time')
     score = resp_json.get('obj_order_rating_goods', {}).get('score')  # 中检评分
@@ -192,7 +192,7 @@ def loop_rating_no(log, sql_pool, sql_ra_no_list):
                 log.warning(f"other warning, please check ......................................")
                 sql_pool.update_one('update zhongjian_task set state = 3 where tag_no = %s', (rating_no_,))
         except Exception as e:
-            log.warning(f"{inspect.currentframe().f_code.co_name} error: {e}")
+            log.warning(f"{sys._getframe().f_code.co_name} error: {e}")
             sql_pool.update_one('update zhongjian_task set state = 3 where tag_no = %s', (rating_no_,))
             continue
 
@@ -204,34 +204,37 @@ def zhongjian_main(log):
     :param log:
     """
     log.info(
-        f'开始运行 {inspect.currentframe().f_code.co_name} 爬虫任务....................................................')
+        f'开始运行 {sys._getframe().f_code.co_name} 爬虫任务....................................................')
 
     # 配置 MySQL 连接池
     sql_pool = MySQLConnectionPool(log=log)
-    if not sql_pool:
+    if not sql_pool.check_pool_health():
         log.error("MySQL数据库连接失败")
         raise Exception("MySQL数据库连接失败")
 
     try:
-        # while True:
-        # sql_ra_no_list = sql_pool.select_all('select tag_no from zhongjian_task where state = 0 limit 10000')
-        sql_ra_no_list = sql_pool.select_all(
-            "select tag_no from zhongjian_task where tag_no like '529%' and state = 0 limit 50000")
-        # sql_ra_no_list = sql_pool.select_all("select tag_no from zhongjian_task where tag_no > '519354131' and state != 1 limit 10000")
-        sql_ra_no_list = [i[0] for i in sql_ra_no_list]
-        if not sql_ra_no_list:
-            log.info(f'没有需要处理的数据,等待下一轮处理........................................................')
-            # break
-            return
-
-        try:
-            loop_rating_no(log, sql_pool, sql_ra_no_list)
-        except  Exception as e:
-            log.error(f'{inspect.currentframe().f_code.co_name} error: {e}')
+        # 单轮内循环:每批取 5 万 state=0,处理完再取下一批,直到本轮把 state=0 抓空为止。
+        # loop_rating_no 对每条都会把 state 置为 1/2/3(必然离开 state=0),故循环单调收敛、不会死循环。
+        batch_no = 0
+        while True:
+            # 519/529/539/549 新入队的号码段均在此,走 index_state 索引
+            sql_ra_no_list = sql_pool.select_all(
+                "select tag_no from zhongjian_task where state = 0 limit 50000")
+            sql_ra_no_list = [i[0] for i in sql_ra_no_list]
+            if not sql_ra_no_list:
+                log.info(f'state=0 任务已全部处理完,本轮结束,等待下一轮........................................')
+                break
+
+            batch_no += 1
+            log.info(f'第 {batch_no} 批:取到 {len(sql_ra_no_list)} 条待处理,开始采集............')
+            try:
+                loop_rating_no(log, sql_pool, sql_ra_no_list)
+            except Exception as e:
+                log.error(f'{sys._getframe().f_code.co_name} 第 {batch_no} 批 error: {e}')
     except Exception as e:
-        log.error(f'{inspect.currentframe().f_code.co_name} error: {e}')
+        log.error(f'{sys._getframe().f_code.co_name} error: {e}')
     finally:
-        log.info(f'爬虫程序 {inspect.currentframe().f_code.co_name} 运行结束,等待下一轮的采集任务............')
+        log.info(f'爬虫程序 {sys._getframe().f_code.co_name} 运行结束,等待下一轮的采集任务............')
 
 
 def schedule_task():
@@ -251,4 +254,5 @@ def schedule_task():
 
 
 if __name__ == '__main__':
-    schedule_task()
+    # schedule_task()
+    zhongjian_main(log=logger)