snapshot_jd.py.bak 4.2 KB

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