jun 6 일 전
부모
커밋
87ab93b4fc
4개의 변경된 파일123개의 추가작업 그리고 64개의 파일을 삭제
  1. BIN
      spides/.gitignore
  2. BIN
      spides/element_screenshot.png
  3. 0 64
      spides/snapshot_jd.py
  4. 123 0
      spides/snapshot_jd.py.bak

BIN
spides/.gitignore


BIN
spides/element_screenshot.png


+ 0 - 64
spides/snapshot_jd.py

@@ -14,7 +14,6 @@ platform_name = "京东"
 class JdMain:
     def __init__(self):
         # self.db_online = MySQLPool39()
-        self.db_online = MySQLPoolOnline()
         self.crawl_count = ""
         self.task_id = ""
         self.task_dict = None
@@ -24,69 +23,6 @@ class JdMain:
         self.cumulative_stored = 0
         self.cumulative_skipped = 0
 
-    def get_status(self, status):
-        if status not in (2, 3, 4):
-            logger.warning(f"未知状态值: {status}, 跳过状态上报")
-            return
-        if status == 2:
-            parmas = {
-                "collect_task_allocate_id": self.task_id, "status": status, "finish_status": 0,
-                "start_time": int(time.time())
-            }
-        if status == 3:
-            parmas = {
-                "collect_task_allocate_id": self.task_id, "status": status, "finish_status": 1,
-                "real_count": self.crawl_count, "end_time": int(time.time()),
-            }
-        if status == 4:
-            parmas = {
-                "collect_task_allocate_id": self.task_id, "status": status, "finish_status": 0,
-                "end_time": int(time.time())}
-
-        # url = "http://scheduletest.dfwy.tech/api/collect_equipment_execute/result_report"
-        url = "http://scheduleapi.findit.ltd/api/collect_equipment_execute/result_report"
-
-        try:
-            res = requests.get(url, params=parmas, timeout=20)
-            res.raise_for_status()
-            logger.info("状态上报: %s", res.text)
-        except Exception as e:
-            logger.warning(f"状态上报失败: {e}")
-
-    def get_task(self):
-        """获取当前设备绑定的京东待执行快照任务。"""
-        sql = """
-            SELECT t.*
-            FROM `retrieve_collect_task_allocate` t
-            INNER JOIN `retrieve_collect_equipment_account` a
-                ON t.`collect_equipment_account_id` = a.`id`
-            WHERE t.`platform` = 2
-              AND t.`status` = 1
-              AND t.`snapshot_collect_status` = 0
-            LIMIT 1
-        """
-        task_list = self.db_online.select_data(sql)
-        print(task_list)
-        if not task_list:
-            return {}
-
-        task_dict = task_list[0]
-        self.task_id = task_dict["id"]
-        print(task_dict)
-        return task_dict
-
-    def heartbeat_task(self):
-        url = "https://scheduleapi.findit.ltd/api/collect_equipment_execute/heartbeat"
-
-        params = {
-            "collect_task_allocate_id": self.task_id,
-        }
-
-        try:
-            res = requests.get(url, params=params, timeout=20)
-            logger.info("心跳任务上报成功")
-        except Exception as e:
-            logger.info("心跳任务上报失败")
 
     def run(self):
 

+ 123 - 0
spides/snapshot_jd.py.bak

@@ -0,0 +1,123 @@
+import json
+import time
+import requests
+from commons.Logger import logger
+from commons.conn_mysql import MySQLPoolOnline, MySQLPool39
+from spiders.jd.jd_auto_crawl_snap2 import JdCrawlerV2
+from commons.scheduler import CrawlerScheduler
+from commons.feishu_webhook import send_text
+import random
+from commons.config import JD_DEVICE_ID
+platform_name = "京东"
+
+
+class JdMain:
+    def __init__(self):
+        # self.db_online = MySQLPool39()
+        self.db_online = MySQLPoolOnline()
+        self.crawl_count = ""
+        self.task_id = ""
+        self.task_dict = None
+        self.driver = None
+        self.cumulative_pages = 0
+        self.cumulative_items = 0
+        self.cumulative_stored = 0
+        self.cumulative_skipped = 0
+
+    def get_status(self, status):
+        if status not in (2, 3, 4):
+            logger.warning(f"未知状态值: {status}, 跳过状态上报")
+            return
+        if status == 2:
+            parmas = {
+                "collect_task_allocate_id": self.task_id, "status": status, "finish_status": 0,
+                "start_time": int(time.time())
+            }
+        if status == 3:
+            parmas = {
+                "collect_task_allocate_id": self.task_id, "status": status, "finish_status": 1,
+                "real_count": self.crawl_count, "end_time": int(time.time()),
+            }
+        if status == 4:
+            parmas = {
+                "collect_task_allocate_id": self.task_id, "status": status, "finish_status": 0,
+                "end_time": int(time.time())}
+
+        # url = "http://scheduletest.dfwy.tech/api/collect_equipment_execute/result_report"
+        url = "http://scheduleapi.findit.ltd/api/collect_equipment_execute/result_report"
+
+        try:
+            res = requests.get(url, params=parmas, timeout=20)
+            res.raise_for_status()
+            logger.info("状态上报: %s", res.text)
+        except Exception as e:
+            logger.warning(f"状态上报失败: {e}")
+
+    def get_task(self):
+        """获取当前设备绑定的京东待执行快照任务。"""
+        sql = """
+            SELECT t.*
+            FROM `retrieve_collect_task_allocate` t
+            INNER JOIN `retrieve_collect_equipment_account` a
+                ON t.`collect_equipment_account_id` = a.`id`
+            WHERE t.`platform` = 2
+              AND t.`status` = 1
+              AND t.`snapshot_collect_status` = 0
+            LIMIT 1
+        """
+        task_list = self.db_online.select_data(sql)
+        print(task_list)
+        if not task_list:
+            return {}
+
+        task_dict = task_list[0]
+        self.task_id = task_dict["id"]
+        print(task_dict)
+        return task_dict
+
+    def heartbeat_task(self):
+        url = "https://scheduleapi.findit.ltd/api/collect_equipment_execute/heartbeat"
+
+        params = {
+            "collect_task_allocate_id": self.task_id,
+        }
+
+        try:
+            res = requests.get(url, params=params, timeout=20)
+            logger.info("心跳任务上报成功")
+        except Exception as e:
+            logger.info("心跳任务上报失败")
+
+    def run(self):
+
+        spider_schedule = CrawlerScheduler(JD_DEVICE_ID, 2)
+        spider_schedule.start()
+        time.sleep(3)
+
+        while 1:
+            if not spider_schedule.end:
+                self.task_dict = spider_schedule.get_task()
+
+                if not self.task_dict:
+                    logger.info(f"{platform_name}暂无任务")
+                    time.sleep(35)
+                    continue
+
+                self.task_id = self.task_dict.get("id", "")
+                self.crawl_count, is_success, self.driver, self.cumulative_pages, self.cumulative_items, self.cumulative_stored, self.cumulative_skipped = JdCrawlerV2(self.task_dict, spider_schedule, self.driver, self.cumulative_pages, self.cumulative_items, self.cumulative_stored, self.cumulative_skipped).run()
+                spider_schedule.stop()
+
+            else:
+                spider_schedule.start()
+                print('休息')
+                time.sleep(35)
+
+            time.sleep(30)
+
+
+if __name__ == '__main__':
+    while True:
+        JdMain().run()
+        interval_time = random.randint(1200, 1800)
+        logger.info(f"程序睡眠{interval_time}秒后继续执行")
+        time.sleep(interval_time)