【问题标题】:Continuously probing multiple ULRs with fewer threads, how to control threads用更少的线程连续探测多个ULR,如何控制线程
【发布时间】:2018-08-23 02:33:19
【问题描述】:

背景:

我想监控 100 个 URL(拍摄快照,如果内容与以前不同,则将其存储),我的计划是使用 urllib.request 每 x 分钟扫描一次,比如 x=5,不间断。

所以我不能使用单个 for 循环和睡眠,因为我想启动对 ULR1 的检测,然后几乎同时启动 URL2。

while TRUE:
  for url in urlList:
    do_detection()
    time.sleep(sleepLength)

因此我应该使用池?但我应该将线程限制在我的 CPU 可以处理的少量(如果我有 100 个 ULR,则不能设置为 100 个线程)

我的问题:

即使我可以将列表中的 100 个 URL 发送到具有四个线程的 ThreadPool(4),我应该如何设计来控制每个线程来处理 100/4=25 个 URL,因此线程会探测 URL1,sleep(300)在下一次探测到 URL1 之前,然后执行 URL2.... ULR25 并返回到 URL1...?我不想等待 5 分钟 * 25 个完整的周期。

伪代码或示例将有很大帮助!我找不到或想办法让 looper() 和detector() 按需要运行?

(我认为How to scrap multiple html page in parallel with beautifulsoup in python? 这是接近但不准确的答案)

也许每个线程都这样?我现在将尝试解决如何将 100 个项目拆分到每个线程。使用 pool.map(func, iterable[, chunksize]) 需要一个列表,我可以将 chunksize 设置为 25。

def one_thread(Url):

    For url in Url[0:24]:
          CurrentDetect(url)
    if 300-timelapsed>0:
        remain_sleeping=300-timtlapsed
    else:
        remain_sleeping=0


    sleep (remain_sleeping)

    For url in Url[0:24]:
          NextDetect()

我正在尝试编写的非工作代码:

import urllib.request as req
import time
def url_reader(url = "http://stackoverflow.com"):

    try
        f = req.urlopen(url)
        print (f.read())

    except Exception as err
        print (err)

def save_state():
    pass
    return []

def looper (sleepLength=720,urlList):
    for url in urlList: #initial save
        Latest_saved.append(save_state(url_reader(url))) # return a list
    while TRUE:
        pool = ThreadPool(4) 


        results = pool.map(urllib2.urlopen, urls)
        time.sleep(sleepLength)  # how to parallel this? if we have 100 urls, then takes 100*20 min to loop?
        detector(urlList) #? use last saved status returned to compare?

def detector (urlList):




    for url in urlList:
            contentFirst=url_reader(url)

            contentNext=url_reader(url)

            if contentFirst!=contentNext:
                save_state(contentFirst)
                save_state(contentNext)

【问题讨论】:

标签: python multithreading


【解决方案1】:

您需要安装requests

pip install requests

如果你想使用下面的代码:

# -*- coding: utf-8 -*-

import concurrent.futures
import requests
import queue
import threading

# URL Pool
URLS = [
    # Put your urls here
]

# Time interval (in seconds)
INTERVAL = 5 * 60

# The number of worker threads
MAX_WORKERS = 4

# You should set up request headers
# if you want to better evade anti-spider programs
HEADERS = {
    'Accept': '*/*',
    'Accept-Encoding': 'gzip, deflate',
    'Accept-Language': 'en-US,en;q=0.9',
    'Cache-Control': 'max-age=0',
    'Connection': 'keep-alive',
    #'Host': None,
    'If-Modified-Since': '0',
    #'Referer': None,
    'User-Agent': 'Mozilla/5.0 (Windows NT 6.1) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/67.0.3396.62 Safari/537.36',
}

############################

def handle_response(response):
    # TODO implement your logics here !!!
    raise RuntimeError('Please implement function `handle_response`!')

# Retrieve a single page and report the URL and contents
def load_url(session, url):
    #print('load_url(session, url={})'.format(url))
    response = session.get(url)
    if response.status_code == 200:
        # You can refactor this part and
        # make it run in another thread
        # devoted to handling local IO tasks,
        # to reduce the burden of Net IO worker threads
        return handle_response(response)

def ThreadPoolExecutor():
    return concurrent.futures.ThreadPoolExecutor(max_workers=MAX_WORKERS)

# Generate a session object
def Session():
    session = requests.Session()
    session.headers.update(HEADERS)
    return session

# We can use a with statement to ensure threads are cleaned up promptly
with ThreadPoolExecutor() as executor, Session() as session:
    if not URLS:
        raise RuntimeError('Please fill in the array `URLS` to start probing!')

    tasks = queue.Queue()

    for url in URLS:
        tasks.put_nowait(url)

    def wind_up(url):
        #print('wind_up(url={})'.format(url))
        tasks.put(url)

    while True:
        url = tasks.get()

        # Work
        executor.submit(load_url, session, url)

        threading.Timer(interval=INTERVAL, function=wind_up, args=(url,)).start()

【讨论】:

  • 非常感谢 KaiserKatze!我尝试了代码,但它卡在了 load_url() 中,没有完成 session.get(url)。我会进行一些调试,可能会使用更简单的操作而不是读取 url 来删除一些锁定或测试,因为我可以看到前 2 个工作人员启动但立即其他 url 也提交,可能阻塞队列或其他什么?
  • 很可能,1个会话对象不能被2个线程同时使用?
  • 下面是立即输出(立即,然后挂起) url1 popleft OK load_url 函数已启动#(但未完成 session.request,我还尝试使用自己的调用 init_request 对每个负载进行新的session object) url2 popleftOK load_url 函数启动 url3 load_url 函数启动 url4 url​​5 load_url url6 load_url url7 load_url url8 load_url url9 load_url load_url
  • @LanSi 哦,我刚刚调试了我的代码。你来了。
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2016-08-31
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多