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

Django多进程写入PostgreSQL为何会触发“too many clients already”错误?

Django WebSocket数据批量写入PostgreSQL的优化指南

看你这段代码,核心是通过WebSocket持续捞数据,攒够3000条就用进程池异步写入PostgreSQL对吧?我给你拆解下当前实现的关键点,还有几个能提升稳定性和效率的优化建议:

一、当前逻辑的核心思路

你的MyWebsocketClient里的process_message方法是这么干活的:

  • 把收到的WebSocket消息往data_buffer里塞
  • 缓冲区满3000条时,新建一个CPU核心数2/3大小的进程池
  • 用pool.map_async调用store_data逐条创建BinanceTick对象
  • 清空缓冲区继续接数据

二、几个值得优化的点

1. 别每次都新建进程池!

每次缓冲区满了就创建新进程池,这会浪费不少系统资源——进程创建和销毁都是有开销的。咱们可以提前初始化一个全局进程池,复用它:

from multiprocessing import Pool, cpu_count
import threading

class MyWebsocketClient:
    def __init__(self):
        self.data_buffer = []
        # 初始化一次进程池,后续复用
        self.pool = Pool(cpu_count() * 2 // 3)
        # 处理线程安全的锁
        self.buffer_lock = threading.Lock()
    
    def process_message(self, msg):
        with self.buffer_lock:
            self.data_buffer.append(msg)
            if len(self.data_buffer) > 3000:
                # 复制缓冲区数据,避免异步处理时被覆盖
                buffer_copy = self.data_buffer.copy()
                self.data_buffer = []
        
        # 用已有的进程池提交批量入库任务
        self.pool.apply_async(store_data_batch, args=(buffer_copy,))
    
    def cleanup(self):
        # 程序退出前务必关闭进程池,等待所有任务完成
        self.pool.close()
        self.pool.join()

2. 解决Django ORM多进程兼容问题

直接在子进程里用BinanceTick.objects.create容易出问题——Django的数据库连接是进程隔离的,子进程复用父进程的连接会导致各种奇怪的错误。所以要在子进程里重新初始化数据库连接:

def store_data_batch(data_buffer):
    from django.db import connection
    # 子进程重新建立独立的数据库连接
    connection.connect()
    
    # 将缓冲区数据转换为BinanceTick对象列表
    tick_objs = [
        BinanceTick(
            field1=item['field1'],
            field2=item['field2'],
            # 其他字段按需映射
        ) for item in data_buffer
    ]
    # 批量插入,比逐条create效率提升数倍
    BinanceTick.objects.bulk_create(tick_objs, batch_size=1000)
    
    # 子进程用完关闭连接
    connection.close()

3. 把逐条插入改成批量插入

你原来用map_async是每条数据调用一次create,这太浪费数据库IO了!用Django的bulk_create一次插入一批数据,能把单条插入的网络开销、事务开销摊薄,效率提升非常明显。

4. 缓冲区的线程安全要重视

如果你的WebSocket客户端是在多线程环境下运行(比如大部分WebSocket库的回调都是在单独线程里触发的),那data_buffer的读写就会有线程安全问题——多个线程同时往里面加数据、清空数据,很容易丢数据或者重复处理。所以一定要给缓冲区加个锁,就像上面代码里的buffer_lock那样,用with self.buffer_lock:把缓冲区操作包起来。

三、额外的小提醒

  • 监控任务状态:可以给apply_async加回调函数,方便排查问题:
    def on_task_finish(_):
        print("批量入库任务完成")
    
    def on_task_error(error):
        print(f"入库失败:{str(error)}")
    
    self.pool.apply_async(store_data_batch, args=(buffer_copy,), callback=on_task_finish, error_callback=on_task_error)
    
  • 控制进程池大小:PostgreSQL默认最大连接数是100,进程池大小如果超过这个数,会导致数据库连接被占满,新任务卡住。建议根据数据库配置调整,比如CPU核心数的1/2或2/3都可以。
  • 优雅退出:程序结束前一定要调用cleanup方法,等进程池里的所有任务都完成再退出,避免正在入库的数据丢失。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.21 03:27:37