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

如何在Python中调度concurrent.futures.ThreadPoolExecutor的map()函数

问题解决方案

核心问题排查

你的问题大概率出自两个原因:

  1. 线程池异常被静默吞掉:ThreadPoolExecutor的submit或map方法如果触发异常,不会主动抛出,必须手动获取任务结果才能捕获错误。
  2. 调度逻辑与线程执行冲突:while循环的时间等待逻辑可能阻塞了线程池任务,或是没有正确触发任务执行流程。

修正后的完整代码与执行流程

下面是整合时间调度与并发抓取的完整实现,每一步逻辑清晰:

1. 导入依赖库

import time
import datetime
import concurrent.futures
import requests
import pyodbc  # 用pyodbc连接SQL Server,根据你的实际库调整

2. 单个URL抓取与入库函数

负责单URL的抓取、解析、入库,强制包含异常捕获:

def fetch_and_save_single_url(url):
    try:
        # 替换为你的网页抓取逻辑
        response = requests.get(url, timeout=12)
        response.raise_for_status()  # 触发HTTP状态码错误
        parsed_data = parse_response_content(response.text)  # 替换为你的数据解析逻辑

        # 替换为你的SQL Server入库逻辑
        conn = pyodbc.connect('DRIVER={SQL Server};SERVER=你的服务器地址;DATABASE=你的库名;UID=账号;PWD=密码')
        cursor = conn.cursor()
        insert_sql = "INSERT INTO 你的表名 (字段1, 字段2) VALUES (?, ?)"
        cursor.execute(insert_sql, parsed_data['字段1'], parsed_data['字段2'])
        conn.commit()
        cursor.close()
        conn.close()

        return f"成功处理: {url}"
    except Exception as e:
        # 捕获所有异常并返回,避免线程池吞掉错误信息
        return f"处理失败 {url}: {str(e)}"

3. 批量并发处理函数

显式获取线程池任务结果,确保异常能被输出:

def batch_process_urls(url_list):
    # 线程池大小建议设为10-20,根据服务器性能和目标网站反爬规则调整
    with concurrent.futures.ThreadPoolExecutor(max_workers=15) as executor:
        # 提交所有任务并同步获取结果
        task_results = list(executor.map(fetch_and_save_single_url, url_list))
    
    # 打印所有任务的执行结果
    for result in task_results:
        print(result)

4. 时间调度逻辑(while循环实现)

严格控制工作日9:15-15:30期间每3分钟执行一次:

def run_scheduler(url_list):
    while True:
        current_time = datetime.datetime.now()
        # 检查是否为工作日(weekday()返回0=周一,4=周五)
        if current_time.weekday() < 5:
            # 定义任务执行的时间范围
            daily_start = current_time.replace(hour=9, minute=15, second=0, microsecond=0)
            daily_end = current_time.replace(hour=15, minute=30, second=0, microsecond=0)
            
            if daily_start <= current_time <= daily_end:
                print(f"{current_time.strftime('%Y-%m-%d %H:%M:%S')} 启动批量抓取任务")
                # 调用并发处理函数
                batch_process_urls(url_list)
                # 执行完任务后等待3分钟,避免任务未完成就进入下一轮
                time.sleep(180)
            else:
                # 非目标时间段,每分钟检查一次时间
                time.sleep(60)
        else:
            # 周末,每小时检查一次时间
            time.sleep(3600)

5. 主程序入口

if __name__ == "__main__":
    # 替换为你的190个URL列表
    target_urls = ["https://example.com/page1", "https://example.com/page2", ...]
    run_scheduler(target_urls)

关键注意事项

  • 异常必捕获:fetch_and_save_single_url必须捕获所有异常,batch_process_urls必须获取任务结果,否则线程池内的错误会完全静默,导致你看不到任何输出。
  • 线程池大小合理:不要设置过大,否则可能触发目标网站反爬机制,或导致本地资源耗尽。
  • 数据库连接隔离:每个线程单独创建数据库连接,禁止多线程共享同一个连接,避免线程安全问题。
  • 等待逻辑优化:如果单次任务执行时间超过3分钟,可改为计算下一次执行的时间点(比如next_run = current_time + datetime.timedelta(minutes=3)),再计算等待时长,避免任务重叠。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.15 01:25:25