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

基于Multiprocessing与PySftp的并行下载成功记录入库方案问询

Solution for Capturing Parallel Download Results & Logging to Database

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.starmap automatically 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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.27 04:16:43