scheduler.py 4.0 KB

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