ProxyRefreshSchedule.py 5.8 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185
  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 os
  16. import random
  17. import sys
  18. import time
  19. from multiprocessing import Process
  20. import requests
  21. from apscheduler.schedulers.blocking import BlockingScheduler
  22. sys.path.append('../')
  23. from DB.DbClient import DbClient
  24. from Manager.ProxyManager import ProxyManager
  25. __author__ = 'JHao'
  26. log = logging.getLogger('apscheduler')
  27. class ProxyRefreshSchedule(ProxyManager):
  28. """
  29. 代理定时刷新
  30. """
  31. def __init__(self):
  32. ProxyManager.__init__(self)
  33. def valid_proxy(self):
  34. logger_raw = logging.getLogger('apscheduler.raw_check-{}'.format(os.getpid()))
  35. fmt = logging.Formatter('%(levelname)s:%(name)s:%(message)s')
  36. fh = logging.FileHandler(filename='../log/log.txt')
  37. fh.setFormatter(fmt=fmt)
  38. logger_raw.addHandler(fh)
  39. self.db.changeTable(self.raw_proxy_queue)
  40. raw_proxy = self.db.pop()
  41. while raw_proxy:
  42. logger_raw.debug('[*] check raw proxy {} ...'.format(raw_proxy))
  43. # print '[*] check raw proxy {} ...'.format(raw_proxy)
  44. proxies = {"http": "http://{proxy}".format(proxy=raw_proxy),
  45. "https": "https://{proxy}".format(proxy=raw_proxy)}
  46. try:
  47. # 超过30秒的代理就不要了
  48. r = requests.get('https://www.baidu.com/', proxies=proxies, timeout=30, verify=False)
  49. if r.status_code == 200:
  50. self.db.changeTable(self.useful_proxy_queue)
  51. self.db.put(raw_proxy)
  52. logger_raw.debug('[+] raw proxy {} succeed validating...'.format(raw_proxy))
  53. # print '[+] raw proxy {} succeed validating...'.format(raw_proxy)
  54. except Exception, e:
  55. print e
  56. pass
  57. self.db.changeTable(self.raw_proxy_queue)
  58. raw_proxy = self.db.pop()
  59. logger_raw.debug('[-] raw proxy {} invalid'.format(raw_proxy))
  60. # print '[-] raw proxy {} invalid'.format(raw_proxy)
  61. def validate_useful_proxy(self, proxy_list):
  62. logger_avail = logging.getLogger('apscheduler.avail_check-{}'.format(os.getpid()))
  63. fmt = logging.Formatter('%(levelname)s:%(name)s:%(message)s')
  64. fh = logging.FileHandler(filename='../log/log.txt')
  65. fh.setFormatter(fmt=fmt)
  66. logger_avail.addHandler(fh)
  67. len_proxy = len(proxy_list)
  68. self.db.changeTable(self.useful_proxy_queue)
  69. while proxy_list:
  70. proxy = proxy_list.pop()
  71. logger_avail.debug('[*] check available proxy : {} ...({} remained)'.format(proxy, len(proxy_list)))
  72. # print '[*] check available proxy : {} ...({} remained)'.format(proxy, len(proxy_list))
  73. proxies = {"http": "http://{proxy}".format(proxy=proxy),
  74. "https": "https://{proxy}".format(proxy=proxy)}
  75. try:
  76. r = requests.get('https://www.baidu.com/', proxies=proxies, timeout=30, verify=False)
  77. if r.status_code == 200:
  78. logger_avail.debug('[+] proxy {} is still available'.format(proxy))
  79. continue
  80. except Exception, e:
  81. self.db.delete(proxy)
  82. logger_avail.debug('[-] A checked proxy {} has been removed from useful_proxy_queue'.format(proxy))
  83. # print '[-] A checked proxy {} has been removed from useful_proxy_queue'.format(proxy)
  84. logger_avail.debug('Process {} finished checking {} proxies.'.format(os.getpid(), len_proxy))
  85. def refresh_pool():
  86. pp = ProxyRefreshSchedule()
  87. pp.valid_proxy()
  88. def validate_user_proxy(proxy_list):
  89. pp = ProxyRefreshSchedule()
  90. pp.validate_useful_proxy(proxy_list)
  91. def main(process_num=10):
  92. p = ProxyRefreshSchedule()
  93. p.refresh()
  94. pl = []
  95. for num in range(process_num):
  96. proc = Process(target=refresh_pool, args=())
  97. # proc.daemon = True
  98. pl.append(proc)
  99. for num in range(process_num):
  100. pl[num].start()
  101. for num in range(process_num):
  102. pl[num].join()
  103. log.debug('Process main completed.')
  104. def main_check(process_num=10):
  105. db = DbClient()
  106. useful_proxy_queue = 'useful_proxy_queue'
  107. db.changeTable(useful_proxy_queue)
  108. proxy_list = db.getAll()
  109. uncheck_list = [list() for i in xrange(process_num)]
  110. for proxy in proxy_list:
  111. uncheck_list[random.randint(0, process_num - 1)].append(proxy)
  112. pl = []
  113. for num in range(process_num):
  114. proc = Process(target=validate_user_proxy, args=(uncheck_list[num],))
  115. # proc.daemon = True
  116. pl.append(proc)
  117. for num in range(process_num):
  118. pl[num].start()
  119. for num in range(process_num):
  120. pl[num].join()
  121. log.debug('Process main_check completed.')
  122. to_time = time.time()
  123. def test():
  124. print 'start test', time.time() - to_time
  125. time.sleep(20)
  126. print 'end test', time.time() - to_time
  127. def test2():
  128. print 'start test2', time.time() - to_time
  129. time.sleep(40)
  130. print 'end test2', time.time() - to_time
  131. if __name__ == '__main__':
  132. log.setLevel(logging.DEBUG) # DEBUG
  133. fmt = logging.Formatter('%(levelname)s:%(name)s:%(message)s')
  134. h = logging.StreamHandler()
  135. h.setFormatter(fmt)
  136. log.addHandler(h)
  137. fh = logging.FileHandler(filename='../log/log.txt')
  138. fh.setFormatter(fmt=fmt)
  139. log.addHandler(fh)
  140. # main()
  141. sched = BlockingScheduler()
  142. sched.add_job(main, 'interval', minutes=10)
  143. sched.add_job(main_check, 'interval', minutes=15)
  144. sched.start()