基于Multiprocessing与PySftp的并行下载成功记录入库方案问询
Got it, let's break down how to solve your problem step by step. The core issue here is that Python's multiprocessing isolates memory between processes, so a simple return from your download function won't make its way back to the main process unless you use the right inter-process communication tools. Here's a practical, code-driven solution:
1. Fix Return Values from Parallel Downloads
We'll use multiprocessing.Pool (it handles result collection out of the box) and rewrite the download function to return a structured dictionary with all the info we need.
First, the download function with structured return data:
import pysftp from datetime import datetime def download_file(sftp_config, remote_path, local_path): # Initialize result dict with default values result = { "filename": remote_path.split("/")[-1], "download_time": datetime.now().isoformat(), "success": False, "error": None } try: with pysftp.Connection(**sftp_config) as sftp: sftp.get(remote_path, local_path) result["success"] = True except Exception as e: result["error"] = str(e) return result
Then use Pool to run downloads in parallel and collect results:
from multiprocessing import Pool def main(): # Your SFTP connection config sftp_config = { "host": "your-sftp-host", "username": "your-username", "password": "your-password" # Swap to private_key if using key-based authentication } # List of (remote_file_path, local_save_path) tasks download_tasks = [ ("/remote/docs/report.pdf", "./local/report.pdf"), ("/remote/images/photo.png", "./local/photo.png"), # Add more files here ] # Create a pool of worker processes (adjust count based on your system) with Pool(processes=4) as pool: # Use starmap to pass multiple arguments to the download function results = pool.starmap(download_file, [(sftp_config,) + task for task in download_tasks]) # Process the collected results process_download_results(results) if __name__ == "__main__": main()
2. Print Successful Download Details
Now that we have all results in a list, we can filter and print the successful ones (plus errors for debugging):
def process_download_results(results): print("\n=== Successful Downloads ===") for res in results: if res["success"]: print(f"✅ File: *{res['filename']}* | Downloaded at: {res['download_time']}") else: print(f"❌ File: *{res['filename']}* | Failed with error: {res['error']}") # Send results to database for logging log_to_database(results)
3. Log to Database
We'll use SQLite for simplicity (swap to psycopg2 for PostgreSQL or mysql-connector for MySQL if needed). First, initialize the database table:
import sqlite3 def init_download_log_db(): conn = sqlite3.connect("download_logs.db") cursor = conn.cursor() # Create logs table if it doesn't exist cursor.execute(""" CREATE TABLE IF NOT EXISTS download_logs ( id INTEGER PRIMARY KEY AUTOINCREMENT, filename TEXT NOT NULL, download_time TEXT NOT NULL, success BOOLEAN NOT NULL, error TEXT ) """) conn.commit() conn.close() # Run this once at the start of your program init_download_log_db()
Then write the function to insert logs into the database:
def log_to_database(results): conn = sqlite3.connect("download_logs.db") cursor = conn.cursor() # Prepare bulk insert data for efficiency log_entries = [ (res["filename"], res["download_time"], res["success"], res["error"]) for res in results ] # Bulk insert all entries cursor.executemany(""" INSERT INTO download_logs (filename, download_time, success, error) VALUES (?, ?, ?, ?) """, log_entries) conn.commit() conn.close() print(f"\nLogged {len(log_entries)} total entries to database.")
Why This Works
multiprocessing.Pool.starmapautomatically collects return values from each worker process and returns them in the same order as your input tasks.- The structured result dictionary lets you easily filter successful downloads, debug failures, and log all necessary data to the database.
- Bulk inserts (
executemany) are far more efficient than inserting one entry at a time, especially for large batches of files.
内容的提问来源于stack exchange,提问作者Thiago Matsui

