scheduler.py 4.3 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140
  1. import time
  2. import requests
  3. from commons.Logger import get_spider_logger
  4. import threading
  5. import urllib3
  6. # 禁用SSL警告
  7. urllib3.disable_warnings(urllib3.exceptions.InsecureRequestWarning)
  8. logger = get_spider_logger('scheduler')
  9. class CrawlerScheduler:
  10. """爬虫任务调度器"""
  11. def __init__(self, DEVICE_ID, platform, heartbeat_interval=30):
  12. """
  13. 初始化调度器
  14. Args:
  15. platform: 平台名称
  16. heartbeat_url: 心跳上报URL
  17. heartbeat_interval: 心跳间隔时间(秒)
  18. """
  19. self.username = DEVICE_ID
  20. self.platform = platform
  21. self._lock = threading.Lock()
  22. # ⭐ 关键修改:将 http:// 改为 https://
  23. self.heartbeat_url = 'https://120.24.26.108:8083/api/collect_task/heartbeat'
  24. self.heartbeat_interval = heartbeat_interval
  25. self.end = False
  26. self.heartbeat_thread = None
  27. def _heartbeat_reporter(self):
  28. """守护线程:只负责上报心跳"""
  29. headers = {'X-Crawler-Token': 'zhijiayun_crawler_2026'}
  30. while True:
  31. try:
  32. # 使用 HTTPS,并忽略证书验证
  33. response = requests.post(
  34. self.heartbeat_url,
  35. json={
  36. "platform": self.platform,
  37. "username": self.username
  38. },
  39. headers=headers,
  40. timeout=3,
  41. verify=False # 忽略SSL证书验证
  42. )
  43. print(f"[心跳] 发送成功: {response.status_code}")
  44. result = response.json()
  45. if self.end or result.get('code') != 'success':
  46. logger.error(f'心跳回传:{result}')
  47. self.set_flag(True)
  48. break
  49. logger.info(f'心跳回传:{result}')
  50. self.set_flag(False)
  51. time.sleep(self.heartbeat_interval)
  52. except Exception as e:
  53. logger.error(e)
  54. print(e)
  55. time.sleep(5)
  56. def start(self):
  57. """启动调度器"""
  58. self.set_flag(False)
  59. self.heartbeat_thread = threading.Thread(
  60. target=self._heartbeat_reporter,
  61. daemon=True
  62. )
  63. self.heartbeat_thread.start()
  64. def stop(self):
  65. """停止调度器"""
  66. self.set_flag(True)
  67. logger.info('心跳停止')
  68. def get_task(self):
  69. try:
  70. # ⭐ 关键修改:将 http:// 改为 https://
  71. task_api = "https://120.24.26.108:8083/api/collect_task/pull"
  72. headers = {'X-Crawler-Token': 'zhijiayun_crawler_2026'}
  73. params = {
  74. 'platform': self.platform,
  75. 'username': self.username
  76. }
  77. response = requests.get(
  78. task_api,
  79. params=params,
  80. headers=headers,
  81. timeout=5,
  82. verify=False # 忽略SSL证书验证
  83. )
  84. result = response.json()
  85. logger.info(f'拉取任务:{result}')
  86. if result.get('code') == 'success':
  87. return result.get('data').get('task')
  88. print('拉取任务返回', result)
  89. except Exception as e:
  90. logger.error('获取任务报错', e)
  91. def post_report(self, data):
  92. try:
  93. # ⭐ 关键修改:将 http:// 改为 https://
  94. url = "https://120.24.26.108:8083/api/collect_task/report"
  95. print('传给返回接口的数据', data)
  96. logger.info(f'report上传数据:{data}')
  97. headers = {'X-Crawler-Token': 'zhijiayun_crawler_2026'}
  98. response = requests.post(
  99. url,
  100. json=data,
  101. headers=headers,
  102. timeout=5,
  103. verify=False # 忽略SSL证书验证
  104. )
  105. result = response.json()
  106. logger.info(result)
  107. if result.get('code') != 'success':
  108. logger.error(f'翻页回传结果不成功:{result}')
  109. self.stop()
  110. print(f'任务进度上传 {result}')
  111. return result
  112. except Exception as e:
  113. logger.error(e)
  114. return None
  115. def set_flag(self, value):
  116. with self._lock:
  117. self.end = value