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