zengwei 1 month ago
parent
commit
3ddd4f0039
3 changed files with 154 additions and 6 deletions
  1. 2 2
      mt_V2/db.py
  2. 12 4
      mt_V2/main.py
  3. 140 0
      mt_V2/scheduler.py

+ 2 - 2
mt_V2/db.py

@@ -13,8 +13,8 @@ def get_connection():
     return pymysql.connect(
         host="120.24.26.108",
         port=3307,
-        user="root",
-        password="zhijiayun123456",
+        user="collect_user",
+        password="collect123456",
         database="drug_retrieve",
         charset="utf8mb4",
         autocommit=False,

+ 12 - 4
mt_V2/main.py

@@ -8,6 +8,7 @@ import subprocess
 import re
 import random
 import datetime
+import sys
 import json
 import unicodedata
 from aip import AipOcr
@@ -772,7 +773,7 @@ class SpiderMonitor(threading.Thread):
 
 
 class MTScreenshot:
-    def __init__(self, d, oss_config, search_key, title_key, scroll_times=4, compress_quality=7, resize_ratio=0.8,device_id=None,
+    def __init__(self, d, oss_config, search_key, title_key, scroll_times=1, compress_quality=7, resize_ratio=0.8,device_id=None,
                  monitor=None):
         self.device_id = device_id
         # 接收外部已连接好的u2设备实例
@@ -3433,7 +3434,8 @@ class MT:
         print(f"列表当前商品店铺名称:{shop_name}")
 
         if price == '' or shop_name == '':
-            print("列表当前商品价格或店铺名称不存在")
+            print("列表当前商品价格或店铺名称不存在","价格:",price,"店铺名:",shop_name)
+
             return "continue"
 
         scrape_date = self.get_current_date()
@@ -3688,8 +3690,8 @@ def fetch_task_from_scheduler(scheduler, device_id):
     task = scheduler.get_task()
     if not task:
         return None
-    # start_offset = task.get("current_page", 0)  # 移动端起始偏移量,0=从头开始
-    start_offset = 0  # 移动端起始偏移量,0=从头开始
+    start_offset = task.get("current_page", 0)  # 移动端起始偏移量,0=从头开始
+    # start_offset = 0  # 移动端起始偏移量,0=从头开始
     start_page = start_offset if start_offset > 0 else 1
     end_page = task.get("end_page", 0)
     # 转成 page_range 格式,MT.open_product_list_page → move_to_page_range_start 会跳页
@@ -3895,6 +3897,12 @@ def main():
         format='%(asctime)s [%(threadName)s] %(levelname)s: %(message)s'
     )
 
+    # 终端传入设备 ID:python main.py T4VK4LM7AAUOV8AY
+    global DEVICE_ID
+    if len(sys.argv) > 1 and sys.argv[1].strip():
+        DEVICE_ID = sys.argv[1].strip()
+        logging.info(f"使用终端传入设备: {DEVICE_ID}")
+
     # 自动模式:全局只创建一个调度器 + 一个心跳线程,所有任务复用
     scheduler = None
     if not MANUAL_MODE:

+ 140 - 0
mt_V2/scheduler.py

@@ -0,0 +1,140 @@
+import time
+import requests
+from commons.Logger import get_spider_logger
+import threading
+import urllib3
+
+# 禁用SSL警告
+urllib3.disable_warnings(urllib3.exceptions.InsecureRequestWarning)
+
+logger = get_spider_logger('scheduler')
+
+
+class CrawlerScheduler:
+    """爬虫任务调度器"""
+
+    def __init__(self, DEVICE_ID, platform, heartbeat_interval=30):
+        """
+        初始化调度器
+
+        Args:
+            platform: 平台名称
+            heartbeat_url: 心跳上报URL
+            heartbeat_interval: 心跳间隔时间(秒)
+        """
+
+        self.username = DEVICE_ID
+        self.platform = platform
+
+        self._lock = threading.Lock()
+
+        # ⭐ 关键修改:将 http:// 改为 https://
+        self.heartbeat_url = 'https://120.24.26.108:8083/api/collect_task/heartbeat'
+        self.heartbeat_interval = heartbeat_interval
+        self.end = False
+        self.heartbeat_thread = None
+
+    def _heartbeat_reporter(self):
+        """守护线程:只负责上报心跳"""
+        headers = {'X-Crawler-Token': 'zhijiayun_crawler_2026'}
+        while True:
+            try:
+                # 使用 HTTPS,并忽略证书验证
+                response = requests.post(
+                    self.heartbeat_url,
+                    json={
+                        "platform": self.platform,
+                        "username": self.username
+                    },
+                    headers=headers,
+                    timeout=3,
+                    verify=False   # 忽略SSL证书验证
+                )
+
+                print(f"[心跳] 发送成功: {response.status_code}")
+                result = response.json()
+
+                if self.end or result.get('code') != 'success':
+                    logger.error(f'心跳回传:{result}')
+                    self.set_flag(True)
+                    break
+
+                logger.info(f'心跳回传:{result}')
+                self.set_flag(False)
+
+                time.sleep(self.heartbeat_interval)
+
+            except Exception as e:
+                logger.error(e)
+                print(e)
+                time.sleep(5)
+
+    def start(self):
+        """启动调度器"""
+        self.set_flag(False)
+        self.heartbeat_thread = threading.Thread(
+            target=self._heartbeat_reporter,
+            daemon=True
+        )
+        self.heartbeat_thread.start()
+
+    def stop(self):
+        """停止调度器"""
+        self.set_flag(True)
+        logger.info('心跳停止')
+
+    def get_task(self):
+        try:
+            # ⭐ 关键修改:将 http:// 改为 https://
+            task_api = "https://120.24.26.108:8083/api/collect_task/pull"
+            headers = {'X-Crawler-Token': 'zhijiayun_crawler_2026'}
+            params = {
+                'platform': self.platform,
+                'username': self.username
+            }
+            response = requests.get(
+                task_api,
+                params=params,
+                headers=headers,
+                timeout=5,
+                verify=False   # 忽略SSL证书验证
+            )
+            result = response.json()
+            logger.info(f'拉取任务:{result}')
+            if result.get('code') == 'success':
+                return result.get('data').get('task')
+
+            print('拉取任务返回', result)
+
+        except Exception as e:
+            logger.error('获取任务报错', e)
+
+    def post_report(self, data):
+        try:
+            # ⭐ 关键修改:将 http:// 改为 https://
+            url = "https://120.24.26.108:8083/api/collect_task/report"
+            print('传给返回接口的数据', data)
+            logger.info(f'report上传数据:{data}')
+            headers = {'X-Crawler-Token': 'zhijiayun_crawler_2026'}
+            response = requests.post(
+                url,
+                json=data,
+                headers=headers,
+                timeout=5,
+                verify=False   # 忽略SSL证书验证
+            )
+            result = response.json()
+            logger.info(result)
+
+            if result.get('code') != 'success':
+                logger.error(f'翻页回传结果不成功:{result}')
+                self.stop()
+            print(f'任务进度上传 {result}')
+            return result
+        except Exception as e:
+            logger.error(e)
+            return None
+
+    def set_flag(self, value):
+        with self._lock:
+            self.end = value