ProxyRefreshSchedule.py 4.3 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147
  1. # -*- coding: utf-8 -*-
  2. # !/usr/bin/env python
  3. """
  4. -------------------------------------------------
  5. File Name: ProxyRefreshSchedule.py
  6. Description : 代理定时刷新
  7. Author : JHao
  8. date: 2016/12/4
  9. -------------------------------------------------
  10. Change Activity:
  11. 2016/12/4: 代理定时刷新
  12. -------------------------------------------------
  13. """
  14. import logging
  15. import random
  16. from multiprocessing import Process
  17. import requests
  18. from apscheduler.schedulers.blocking import BlockingScheduler
  19. import sys
  20. sys.path.append('../')
  21. from DB.DbClient import DbClient
  22. from Manager.ProxyManager import ProxyManager
  23. __author__ = 'JHao'
  24. class ProxyRefreshSchedule(ProxyManager):
  25. """
  26. 代理定时刷新
  27. """
  28. def __init__(self):
  29. ProxyManager.__init__(self)
  30. def valid_proxy(self):
  31. self.db.changeTable(self.raw_proxy_queue)
  32. raw_proxy = self.db.pop()
  33. print '[*]check raw proxy {} ...'.format(raw_proxy)
  34. while raw_proxy:
  35. proxies = {"http": "http://{proxy}".format(proxy=raw_proxy),
  36. "https": "https://{proxy}".format(proxy=raw_proxy)}
  37. try:
  38. # 超过30秒的代理就不要了
  39. r = requests.get('https://www.baidu.com/', proxies=proxies, timeout=30, verify=False)
  40. if r.status_code == 200:
  41. self.db.changeTable(self.useful_proxy_queue)
  42. self.db.put(raw_proxy)
  43. except Exception, e:
  44. print e
  45. pass
  46. self.db.changeTable(self.raw_proxy_queue)
  47. raw_proxy = self.db.pop()
  48. print 'validate a proxy'
  49. def validate_useful_proxy(self, proxy_list):
  50. self.db.changeTable(self.useful_proxy_queue)
  51. for proxy in proxy_list:
  52. print '[*]validating proxy : {} ...({} remained)'.format(proxy, len(proxy_list))
  53. proxies = {"http": "http://{proxy}".format(proxy=proxy),
  54. "https": "https://{proxy}".format(proxy=proxy)}
  55. try:
  56. r = requests.get('https://www.baidu.com/', proxies=proxies, timeout=30, verify=False)
  57. if r.status_code == 200:
  58. continue
  59. except Exception, e:
  60. self.db.delete(proxy)
  61. print '[-]delete proxy {}'.format(proxy)
  62. def refresh_pool():
  63. pp = ProxyRefreshSchedule()
  64. pp.valid_proxy()
  65. def validate_user_proxy(proxy_list):
  66. pp = ProxyRefreshSchedule()
  67. pp.validate_useful_proxy(proxy_list)
  68. def main(process_num=10):
  69. p = ProxyRefreshSchedule()
  70. p.refresh()
  71. pl = []
  72. for num in range(process_num):
  73. proc = Process(target=refresh_pool, args=())
  74. proc.daemon = True
  75. pl.append(proc)
  76. for num in range(process_num):
  77. pl[num].start()
  78. print 'All raw_proxy_crawler sub-processes start.'
  79. for num in range(process_num):
  80. pl[num].join()
  81. def main_check(process_num=10):
  82. db = DbClient()
  83. useful_proxy_queue = 'useful_proxy_queue'
  84. db.changeTable(useful_proxy_queue)
  85. proxy_list = db.getAll()
  86. uncheck_list = [list() for i in xrange(process_num)]
  87. for proxy in proxy_list:
  88. uncheck_list[random.randint(0, process_num - 1)].append(proxy)
  89. pl = []
  90. for num in range(process_num):
  91. proc = Process(target=validate_user_proxy, args=(uncheck_list[num],))
  92. proc.daemon = True
  93. pl.append(proc)
  94. for num in range(process_num):
  95. pl[num].start()
  96. print 'All proxy validator sub-processes start.'
  97. for num in range(process_num):
  98. pl[num].join()
  99. # def main(process_num=100):
  100. # p = ProxyRefreshSchedule()
  101. # p.refresh()
  102. # for num in range(process_num):
  103. # P = Process(target=refreshPool, args=())
  104. # P.daemon = True
  105. # P.start()
  106. # P.join()
  107. # print '{time}: refresh complete!'.format(time=time.ctime())
  108. if __name__ == '__main__':
  109. log = logging.getLogger('apscheduler')
  110. log.setLevel(logging.INFO) # DEBUG
  111. fmt = logging.Formatter('%(levelname)s:%(name)s:%(message)s')
  112. h = logging.StreamHandler()
  113. h.setFormatter(fmt)
  114. log.addHandler(h)
  115. # main()
  116. sched = BlockingScheduler()
  117. # sched.add_job(main, 'interval', seconds=10)
  118. sched.add_job(main_check, 'interval', seconds=15)
  119. sched.start()