scheduler.py 3.9 KB

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