Explorar o código

[update] helper

jhao %!s(int64=6) %!d(string=hai) anos
pai
achega
eef8b9bc97
Modificáronse 4 ficheiros con 51 adicións e 43 borrados
  1. 22 1
      helper/check.py
  2. 5 0
      helper/fetch.py
  3. 23 41
      helper/scheduler.py
  4. 1 1
      proxyPool.py

+ 22 - 1
helper/check.py

@@ -52,6 +52,9 @@ def proxyCheck(proxy_obj):
 
 
 class Checker(Thread):
+    """
+    多线程检测代理是否可用
+    """
 
     def __init__(self, check_type, queue, thread_name):
         Thread.__init__(self, name=thread_name)
@@ -66,7 +69,7 @@ class Checker(Thread):
             try:
                 proxy_json = self.queue.get(block=False)
             except Empty:
-                self.log.info("ProxyCheck - {}  : exit".format(self.name))
+                self.log.info("ProxyCheck - {}  : complete".format(self.name))
                 break
 
             proxy = Proxy.createFromJson(proxy_json)
@@ -83,3 +86,21 @@ class Checker(Thread):
             else:
                 pass
             self.queue.task_done()
+
+
+def runChecker(tp, queue):
+    """
+    run Checker
+    :param tp: raw/use
+    :param queue: Proxy Queue
+    :return:
+    """
+    thread_list = list()
+    for index in range(20):
+        thread_list.append(Checker(tp, queue, "thread_%s" % str(index).zfill(2)))
+
+    for thread in thread_list:
+        thread.start()
+
+    for thread in thread_list:
+        thread.join()

+ 5 - 0
helper/fetch.py

@@ -55,4 +55,9 @@ class Fetcher(object):
             except Exception as e:
                 self.log.error("ProxyFetch - {func}: error".format(func=fetch_name))
                 self.log.error(str(e))
+        self.log.info("ProxyFetch - all complete!")
         return proxy_set
+
+
+def runFetcher():
+    return Fetcher().fetch()

+ 23 - 41
helper/scheduler.py

@@ -13,71 +13,53 @@
 __author__ = 'JHao'
 
 from apscheduler.schedulers.blocking import BlockingScheduler
+from apscheduler.executors.pool import ProcessPoolExecutor
 
 from util.six import Queue
-from helper.fetch import Fetcher
-from helper.check import Checker
+from helper.fetch import runFetcher
+from helper.check import runChecker
 from helper.proxy import Proxy
 from handler.logHandler import LogHandler
 from handler.proxyHandler import ProxyHandler
 
 
-def doProxyFetch():
+def runProxyFetch():
     proxy_queue = Queue()
 
-    fetcher = Fetcher()
-    for proxy in fetcher.fetch():
+    for proxy in runFetcher():
         proxy_queue.put(Proxy(proxy).to_json)
 
-    thread_list = list()
-    for index in range(20):
-        thread_list.append(Checker("raw", proxy_queue, "thread_%s" % str(index).zfill(2)))
+    runChecker("raw", proxy_queue)
 
-    for thread in thread_list:
-        thread.start()
 
-    for thread in thread_list:
-        thread.join()
-
-
-def doProxyCheck():
+def runProxyCheck():
     proxy_queue = Queue()
 
-    proxy_handler = ProxyHandler()
-    for proxy in proxy_handler.getAll():
+    for proxy in ProxyHandler().getAll():
         proxy_queue.put(proxy.to_json)
 
-
-# class DoFetchProxy(ProxyManager):
-#     """ fetch proxy"""
-#
-#     def __init__(self):
-#         ProxyManager.__init__(self)
-#         self.log = LogHandler('fetch_proxy')
-#
-#     def main(self):
-#         self.log.info("start fetch proxy")
-#         self.fetch()
-#         self.log.info("finish fetch proxy")
-#
-#
-# def rawProxyScheduler():
-#     DoFetchProxy().main()
-#     doRawProxyCheck()
-#
-#
-# def usefulProxyScheduler():
-#     doUsefulProxyCheck()
+    runChecker("use", proxy_queue)
 
 
 def runScheduler():
-    doProxyFetch()
+    runProxyFetch()
 
     scheduler_log = LogHandler("scheduler")
     scheduler = BlockingScheduler(logger=scheduler_log)
 
-    scheduler.add_job(doProxyFetch, 'interval', minutes=5, id="proxy_fetch", name="proxy采集")
-    # scheduler.add_job(usefulProxyScheduler, 'interval', minutes=1, id="useful_proxy_check", name="useful_proxy定时检查")
+    scheduler.add_job(runProxyFetch, 'interval', minutes=4, id="proxy_fetch", name="proxy采集")
+    scheduler.add_job(runProxyCheck, 'interval', minutes=2, id="proxy_check", name="proxy检查")
+
+    executors = {
+        'default': {'type': 'threadpool', 'max_workers': 20},
+        'processpool': ProcessPoolExecutor(max_workers=5)
+    }
+    job_defaults = {
+        'coalesce': False,
+        'max_instances': 10
+    }
+
+    scheduler.configure(executors=executors, job_defaults=job_defaults)
 
     scheduler.start()
 

+ 1 - 1
proxyPool.py

@@ -16,7 +16,7 @@ import click
 
 from config.setting import BANNER
 
-from helper.proxyScheduler import runScheduler
+from helper.scheduler import runScheduler
 from api.proxyApi import runFlask
 
 CONTEXT_SETTINGS = dict(help_option_names=['-h', '--help'])