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

多进程使用PyMongo时抛出WinError 10048错误的解决方法

多进程使用PyMongo触发WinError 10048错误的解决方法

问题场景

在多进程环境中使用PyMongo,通过6个进程执行计算并将结果插入MongoDB。代码中每个子进程都创建了独立的MongoClient,但运行20-30秒后抛出如下错误:

error! AutoReconnect localhost:27017: [WinError 10048] Only one usage of each socket address (protocol/network address/port) is normally permitted (configured timeouts: connectTimeoutMS: 20000.0ms)

原代码如下:

import multiprocessing
import sys
import pymongo
import datetime

from multiprocessing import Pool


def query_records_by_date(ticker, start_of_day):
    """
    根据参数从MongoDB查询记录并返回列表
    """
    mongo = pymongo.MongoClient("mongodb://localhost:27017/")
    ...


def get_all_post_dates(ticker):
    """
    通过聚合查询获取日期列表,返回datetime.datetime对象列表
    """
    mongo = pymongo.MongoClient("mongodb://localhost:27017/")
    ...


def process_ticker_attention(ticker, finish_num, lock):
    mongo_p = pymongo.MongoClient("mongodb://localhost:27017/", connectTimeoutMS=600000)
    try:
        ticker_dates = get_all_post_dates(ticker)
        for date in ticker_dates:
            data = query_records_by_date(ticker, date)
            mongo_p['stock']['attention'].update_one({'date': date}, {
                '$set': {
                    'post_numbers': {
                        ticker: len(data)
                    }
                }
            })
        lock.acquire()
        finish_num.value += 1
        sys.stdout.write(f'\rfinish_count: {finish_num.value}')
        sys.stdout.flush()
        lock.release()
    except Exception as e:
        print('error!', e.__class__.__name__, e)


if __name__ == '__main__':
    mongo = pymongo.MongoClient("mongodb://localhost:27017/")
    database_name = 'StockForum'
    lock = multiprocessing.Lock()
    finish_num = multiprocessing.Manager().Value('i', 0)
    tickers = mongo[database_name].list_collection_names()
    p = Pool(6)
    for t in tickers:
        p.apply_async(process_ticker_attention, args=(t, finish_num, lock))
    p.close()
    p.join()

问题本质

错误是短时间内创建了过多TCP连接,导致本地端口耗尽:

  • 每个process_ticker_attention进程创建1个MongoClient
  • 每次调用get_all_post_dates和query_records_by_date都会新建独立的MongoClient,而ticker_dates包含数千个元素,意味着每个进程会额外创建数千个客户端
  • 每个MongoClient默认维护连接池,频繁创建销毁客户端会产生大量TIME_WAIT状态的连接,占用本地端口资源,最终触发WinError 10048

解决方案

1. 复用MongoClient,避免频繁创建

每个子进程只创建一次MongoClient,传递给需要数据库操作的函数,而非每次调用新建:

def query_records_by_date(client, ticker, start_of_day):
    """
    复用传入的MongoClient查询记录
    """
    ...  # 使用client操作数据库


def get_all_post_dates(client, ticker):
    """
    复用传入的MongoClient执行聚合查询
    """
    ...  # 使用client操作数据库


def process_ticker_attention(ticker, finish_num, lock):
    # 子进程内仅创建一次客户端
    mongo_p = pymongo.MongoClient("mongodb://localhost:27017/", connectTimeoutMS=600000)
    try:
        # 传递客户端给函数复用
        ticker_dates = get_all_post_dates(mongo_p, ticker)
        for date in ticker_dates:
            data = query_records_by_date(mongo_p, ticker, date)
            mongo_p['stock']['attention'].update_one({'date': date}, {
                '$set': {
                    'post_numbers': {
                        ticker: len(data)
                    }
                }
            })
        # 操作完成后关闭客户端
        mongo_p.close()
        lock.acquire()
        finish_num.value += 1
        sys.stdout.write(f'\rfinish_count: {finish_num.value}')
        sys.stdout.flush()
        lock.release()
    except Exception as e:
        print('error!', e.__class__.__name__, e)
        # 异常时确保关闭客户端
        mongo_p.close()

2. 限制连接池大小

通过MongoClient参数限制单客户端的连接数,避免连接过载:

mongo_p = pymongo.MongoClient(
    "mongodb://localhost:27017/",
    connectTimeoutMS=600000,
    maxPoolSize=10,  # 限制连接池最大连接数
    minPoolSize=2    # 保持最小空闲连接数
)

3. 批量操作减少请求次数

将循环中的单次update_one改为批量更新,大幅减少连接使用频次:

def process_ticker_attention(ticker, finish_num, lock):
    mongo_p = pymongo.MongoClient("mongodb://localhost:27017/", connectTimeoutMS=600000, maxPoolSize=10)
    try:
        ticker_dates = get_all_post_dates(mongo_p, ticker)
        bulk_ops = []
        for date in ticker_dates:
            data = query_records_by_date(mongo_p, ticker, date)
            bulk_ops.append(
                pymongo.UpdateOne(
                    {'date': date},
                    {'$set': {'post_numbers': {ticker: len(data)}}}
                )
            )
        # 批量执行更新操作
        if bulk_ops:
            mongo_p['stock']['attention'].bulk_write(bulk_ops)
        mongo_p.close()
        lock.acquire()
        finish_num.value += 1
        sys.stdout.write(f'\rfinish_count: {finish_num.value}')
        sys.stdout.flush()
        lock.release()
    except Exception as e:
        print('error!', e.__class__.__name__, e)
        mongo_p.close()

4. 调整Windows系统TCP参数(可选)

若代码优化后仍有问题,可修改系统端口配置:

  • 打开注册表编辑器,定位到HKEY_LOCAL_MACHINE\SYSTEM\CurrentControlSet\Services\Tcpip\Parameters
  • 添加或修改以下DWORD值:
    • MaxUserPort:设置为65534(允许使用的最大端口号)
    • TcpTimedWaitDelay:设置为30(缩短TIME_WAIT状态连接的保留时间)
  • 修改后重启系统生效

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.16 00:41:10