如何以字典为输入实现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_jobfunction is hardcoded to used_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
- Task Queue: The
Queueacts as a task pool—workers pull the next available parameter set automatically once they finish their current job. - Worker Threads: We spin up 10 persistent daemon threads that run until the queue is empty.
- Thread-Safe Data Writes: The
hist_lockensures only one thread modifies thefinal_dfat a time, avoiding data corruption from concurrent writes. - Error Resilience: Basic error handling prevents one failed request from crashing all threads.
- Parameter Conversion: Converts your tuple-based parameters to a dictionary that the
requestslibrary 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
相关产品推荐
相关产品推荐

