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

如何以字典为输入实现Python多线程任务调度?

Hey there! Let's get your multithreaded API request setup working properly so those 10 worker threads automatically grab the next task once they're done. Here's a step-by-step fix and explanation:

Key Issues in Your Current Code

  • You're creating too many threads (one per dictionary key instead of just 10 workers)
  • Tasks aren't being added to the queue, so threads have no work to pick up
  • The opthist_job function is hardcoded to use d_threading[0] instead of dynamic tasks
  • Function parameters don't match between definition and call
  • You're not properly handling DataFrame thread-safety when writing results

Revised Working Code

import threading
from queue import Queue
import requests
import pandas as pd

# Your original dictionary of request parameters
d_thread = {
    0: (('instrumentType', 'OPTSTK'), ('symbol', 'INFY'), ('expiryDate', 'select'), ('optionType', 'PE'), ('strikePrice', '2800'), ('dateRange', 'day'), ('fromDate', '11-04-2012'),('toDate', '12-04-2012'), ('segmentLink', '9'), ('symbolCount', '')),
    12: (('instrumentType', 'OPTSTK'), ('symbol', 'INFY'), ('expiryDate', 'select'), ('optionType', 'PE'), ('strikePrice', '2400'), ('dateRange', 'day'), ('fromDate', '27-04-2012'), ('toDate', '28-04-2012'), ('segmentLink', '9'), ('symbolCount', ''))
    # Add your remaining 498 entries here
}

# Lock for thread-safe DataFrame operations
hist_lock = threading.Lock()
# Global DataFrame to store all combined results
final_df = pd.DataFrame()

def opthist_job(params):
    global final_df
    headers = {
        'Pragma': 'no-cache',
        'Accept-Encoding': 'gzip, deflate, br',
        'Accept-Language': 'en-US,en;q=0.9',
        'User-Agent': 'Mozilla/5.0 (X11; Linux x86_64) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/71.0.3578.98 Safari/537.36',
        'Accept': '*/*',
        'Referer': 'https://www.nseindia.com/products/content/derivatives/equities/historical_fo.htm',
        'X-Requested-With': 'XMLHttpRequest',
        'Connection': 'keep-alive',
        'Cache-Control': 'no-cache',
    }
    
    try:
        # Convert your tuple-based params to a dictionary for requests
        params_dict = dict(params)
        response = requests.get('https://www.nseindia.com/products/dynaContent/common/productsSymbolMapping.jsp', headers=headers, params=params_dict)
        response.raise_for_status()  # Catch HTTP errors like 404/500
        
        # Parse NSE's HTML table response into a DataFrame
        df = pd.read_html(response.text)[0]
        
        # Use lock to safely append to the global DataFrame (prevents race conditions)
        with hist_lock:
            final_df = pd.concat([final_df, df], ignore_index=True)
            
    except Exception as e:
        print(f"Failed to process request: {str(e)}")

def threader(queue):
    while True:
        # Grab the next task from the queue
        params = queue.get()
        # Run the API job
        opthist_job(params)
        # Mark the task as completed so the queue knows to wait for all work
        queue.task_done()

if __name__ == "__main__":
    # Initialize a queue to hold all your API tasks
    task_queue = Queue()
    
    # Add all request parameter sets from your dictionary to the queue
    for params in d_thread.values():
        task_queue.put(params)
    
    # Create exactly 10 worker threads
    for _ in range(10):
        worker = threading.Thread(target=threader, args=(task_queue,))
        worker.daemon = True  # Daemon threads exit when the main thread finishes
        worker.start()
    
    # Wait until all tasks in the queue are finished
    task_queue.join()
    
    # Use your final merged DataFrame
    print("All tasks completed successfully!")
    print(final_df.head())

How It Works

  1. Task Queue: The Queue acts as a task pool—workers pull the next available parameter set automatically once they finish their current job.
  2. Worker Threads: We spin up 10 persistent daemon threads that run until the queue is empty.
  3. Thread-Safe Data Writes: The hist_lock ensures only one thread modifies the final_df at a time, avoiding data corruption from concurrent writes.
  4. Error Resilience: Basic error handling prevents one failed request from crashing all threads.
  5. Parameter Conversion: Converts your tuple-based parameters to a dictionary that the requests library can use.

Quick Tips

  • For simpler thread management, you could use concurrent.futures.ThreadPoolExecutor (it handles the queue/threads behind the scenes).
  • Add a progress counter (with a lock) if you want to track how many tasks are completed.
  • Be mindful of NSE's API rate limits to avoid getting blocked—consider adding small delays if needed.

内容的提问来源于stack exchange,提问作者Sid

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.12 03:54:54