You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何结合队列使用multithreading加速Python网站爬虫?

How to Implement Multithreaded Web Scraping with Dynamic Task Queues in Python

Hey there! I totally get why you're stuck—most multithreading guides fixate on fixed loops like for i in range(5), but dynamic tasks (like finding new URLs mid-scrape) need a different approach. Let's break down how to use thread-safe queues to make your scraper faster while handling those evolving tasks.

Core Idea: Use a Thread-Safe Queue for Dynamic Tasks

Python's built-in queue.Queue is perfect here—it handles all the locking and synchronization behind the scenes, so you don't have to worry about race conditions when multiple threads add or pull tasks. Think of it as a "to-do list" that all your worker threads can safely access, even as you add new items to it while scraping.

Step-by-Step Implementation

1. Basic Setup: Queue + Worker Threads

Here's a straightforward example that uses a queue to manage URLs, with worker threads that keep processing tasks until everything's done (including new URLs found mid-scrape):

import threading
from queue import Queue
import requests
from bs4 import BeautifulSoup  # Adjust based on your parsing tool

# Initialize our task queue and a set to track crawled URLs (to avoid duplicates)
task_queue = Queue()
crawled_urls = set()
# Lock to protect the crawled_urls set (prevents race conditions)
crawl_lock = threading.Lock()

def worker():
    """Worker thread that pulls URLs from the queue and processes them"""
    while True:
        # Get a task from the queue; blocks if the queue is empty
        url = task_queue.get()
        
        # Use None as a signal to exit the thread
        if url is None:
            break
            
        try:
            # Your existing scraping logic here
            response = requests.get(url, timeout=10)
            response.raise_for_status()  # Raise error for bad HTTP status codes
            soup = BeautifulSoup(response.text, 'html.parser')
            
            # Extract new URLs and add them to the queue (dynamic task addition)
            with crawl_lock:
                # Filter for internal company links we haven't crawled yet
                new_links = [
                    link['href'] for link in soup.find_all('a', href=True)
                    if link['href'].startswith('https://your-company-site.com')
                    and link['href'] not in crawled_urls
                ]
                
                for link in new_links:
                    crawled_urls.add(link)
                    task_queue.put(link)
            
            # Process the scraped data (save to DB, file, etc.)
            process_scraped_data(soup, url)
            
        except Exception as e:
            print(f"Failed to crawl {url}: {str(e)}")
        finally:
            # Tell the queue this task is finished (critical for queue.join())
            task_queue.task_done()

def process_scraped_data(soup, url):
    """Example data processing function—customize this to your needs"""
    page_title = soup.title.string if soup.title else "No Title"
    print(f"Successfully processed: {page_title} ({url})")

if __name__ == "__main__":
    # Start with your initial URL(s)
    starting_url = "https://your-company-site.com"
    crawled_urls.add(starting_url)
    task_queue.put(starting_url)
    
    # Number of worker threads—adjust based on your site's capacity
    num_workers = 5
    threads = []
    
    # Create and start worker threads
    for _ in range(num_workers):
        thread = threading.Thread(target=worker)
        thread.start()
        threads.append(thread)
    
    # Wait until all tasks in the queue are completed
    task_queue.join()
    
    # Send termination signals to all workers
    for _ in range(num_workers):
        task_queue.put(None)
    
    # Wait for all threads to finish
    for thread in threads:
        thread.join()
    
    print(f"Scraping complete! Crawled {len(crawled_urls)} pages.")

2. Simplify with concurrent.futures.ThreadPoolExecutor

If you don't want to manage threads manually, ThreadPoolExecutor can handle thread lifecycle for you. Here's how to pair it with a queue for dynamic tasks:

from concurrent.futures import ThreadPoolExecutor
from queue import Queue
import threading
import requests
from bs4 import BeautifulSoup

task_queue = Queue()
crawled_urls = set()
crawl_lock = threading.Lock()

def scrape_url(url):
    """Core scraping logic—called by the executor"""
    try:
        response = requests.get(url, timeout=10)
        response.raise_for_status()
        soup = BeautifulSoup(response.text, 'html.parser')
        
        # Add new tasks to the queue
        with crawl_lock:
            new_links = [
                link['href'] for link in soup.find_all('a', href=True)
                if link['href'].startswith('https://your-company-site.com')
                and link['href'] not in crawled_urls
            ]
            for link in new_links:
                crawled_urls.add(link)
                task_queue.put(link)
        
        process_scraped_data(soup, url)
    except Exception as e:
        print(f"Error scraping {url}: {str(e)}")

def queue_worker(executor):
    """Thread that pulls tasks from the queue and submits them to the executor"""
    while True:
        url = task_queue.get()
        if url is None:
            break
        executor.submit(scrape_url, url)
        task_queue.task_done()

if __name__ == "__main__":
    starting_url = "https://your-company-site.com"
    crawled_urls.add(starting_url)
    task_queue.put(starting_url)
    
    # Use a thread pool with 5 workers
    with ThreadPoolExecutor(max_workers=5) as executor:
        # Start the queue worker thread
        worker_thread = threading.Thread(target=queue_worker, args=(executor,))
        worker_thread.start()
        
        # Wait for all tasks to finish
        task_queue.join()
        
        # Send termination signal
        task_queue.put(None)
        worker_thread.join()
    
    print(f"Scraping finished! Total pages crawled: {len(crawled_urls)}")

Key Things to Remember

  • Thread Safety: Always use locks around shared resources like crawled_urls—without them, multiple threads might try to modify the set at the same time, leading to bugs.
  • Avoid Duplicates: The crawled_urls set is crucial to prevent re-scraping the same page over and over.
  • Control Concurrency: Don't crank up the number of threads too high—you could overwhelm your company's website or trigger anti-scraping measures. Start with 3-5 threads and adjust based on performance.
  • Error Handling: Never skip exception handling—network issues, broken pages, or timeouts are common, and you don't want a single failure to take down an entire worker thread.

This setup will handle dynamic tasks seamlessly—any new URLs you find while scraping get added to the queue, and your worker threads keep chugging away until everything's processed.

内容的提问来源于stack exchange,提问作者E. Jaep

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.19 04:21:09