snapshot_taobao.py 4.3 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122
  1. import random
  2. import time
  3. import requests
  4. from commons.Logger import logger
  5. from commons.conn_mysql import MySQLPoolOn2
  6. from spiders.taobao.snapshot_taobao_crawl2 import TaobaoCrawl
  7. from commons.scheduler import CrawlerScheduler
  8. from commons.config import TB_DEVICE_ID, LEGACY_REPORT_URL, LEGACY_HEARTBEAT_URL
  9. platform_name = "淘宝"
  10. class TaobaoMain:
  11. def __init__(self):
  12. self.db_online = MySQLPoolOn2()
  13. self.crawl_count = ""
  14. self.task_id = ""
  15. self.task_dict = None
  16. self.driver = None
  17. self.cumulative_pages = 0
  18. self.cumulative_items = 0
  19. self.cumulative_stored = 0
  20. self.cumulative_skipped = 0
  21. def get_status(self, status):
  22. if status not in (2, 3, 4):
  23. logger.warning(f"未知状态值: {status}, 跳过状态上报")
  24. return
  25. if status == 2:
  26. parmas = {
  27. "collect_task_allocate_id": self.task_id, "status": status, "finish_status": 0,
  28. "start_time": int(time.time())
  29. }
  30. if status == 3:
  31. parmas = {
  32. "collect_task_allocate_id": self.task_id, "status": status, "finish_status": 1,
  33. "real_count": self.crawl_count, "end_time": int(time.time()),
  34. }
  35. if status == 4:
  36. parmas = {
  37. "collect_task_allocate_id": self.task_id, "status": status, "finish_status": 0,
  38. "end_time": int(time.time())}
  39. # url = "http://scheduletest.dfwy.tech/api/collect_equipment_execute/result_report"
  40. url = LEGACY_REPORT_URL
  41. try:
  42. res = requests.get(url, params=parmas, timeout=20)
  43. res.raise_for_status()
  44. logger.info("状态上报: %s", res.text[:500])
  45. except Exception as e:
  46. logger.warning(f"状态上报失败: {e}")
  47. def heartbeat_task(self):
  48. url = LEGACY_HEARTBEAT_URL
  49. params = {
  50. "collect_task_allocate_id": self.task_id,
  51. }
  52. try:
  53. res = requests.get(url, params=params, timeout=20)
  54. logger.info("心跳任务上报成功")
  55. except Exception as e:
  56. logger.info(f"心跳任务上报失败{str(e)}")
  57. def get_task(self):
  58. """获取当前设备绑定的淘宝待执行任务(快照版)。"""
  59. sql = """
  60. SELECT t.*
  61. FROM `retrieve_collect_task_allocate` t
  62. INNER JOIN `retrieve_collect_equipment_account` a
  63. ON t.`collect_equipment_account_id` = a.`id`
  64. WHERE a.`device_id` = %s
  65. AND t.`platform` = 1
  66. AND t.`status` = 1
  67. AND t.`snapshot_collect_status` = 0
  68. LIMIT 1
  69. """
  70. task_list = self.db_online.select_data(sql, (TB_DEVICE_ID,))
  71. if not task_list:
  72. return {}
  73. task_dict = task_list[0]
  74. self.task_id = task_dict["id"]
  75. return task_dict
  76. def run(self):
  77. spider_schedule = CrawlerScheduler(TB_DEVICE_ID, 1)
  78. spider_schedule.start()
  79. time.sleep(3)
  80. while 1:
  81. print(f'not spider_schedule.end:{not spider_schedule.end}')
  82. if not spider_schedule.end:
  83. print('进入工作')
  84. self.task_dict = spider_schedule.get_task()
  85. if not self.task_dict:
  86. logger.info(f"{platform_name}暂无任务")
  87. time.sleep(35)
  88. continue
  89. print(self.task_dict)
  90. self.task_dict.update({'collect_equipment_account_id': TB_DEVICE_ID})
  91. self.crawl_count, is_success, self.driver, self.cumulative_pages, self.cumulative_items, self.cumulative_stored, self.cumulative_skipped = TaobaoCrawl(self.task_dict, spider_schedule, self.driver, self.cumulative_pages, self.cumulative_items, self.cumulative_stored, self.cumulative_skipped).run()
  92. spider_schedule.stop()
  93. else:
  94. spider_schedule.start()
  95. print('休息')
  96. time.sleep(35)
  97. time.sleep(30)
  98. if __name__ == '__main__':
  99. # 每10分钟执行一次
  100. while True:
  101. TaobaoMain().run()
  102. interval_time = random.randint(1200, 1800)
  103. logger.info(f"程序睡眠{interval_time}秒后继续执行")
  104. time.sleep(interval_time)