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

如何使用Prisma实现主从数据源读写查询自动分流

NestJS + Prisma 3.15 自动读写分流实现

你提到的重写getter的思路完全可行,没必要手动在每次调用时选择数据源,用一层通用代理就能自动识别操作类型完成路由,不用逐个给每个model写重复的getter逻辑。

首先先修正你现有代码的初始化问题:你现在在构造函数参数里声明了readClient和writeClient又在构造函数内重新赋值,属于冗余代码,另外记得补上Prisma连接建立、断开的生命周期逻辑,避免连接泄漏。

核心实现逻辑是把Prisma的操作分成读、写两类,通过ES6 Proxy拦截所有model的方法调用,读操作路由到只读副本客户端,写操作路由到主库客户端,同时处理事务、原生SQL这类特殊场景。

完整实现代码

import { Injectable, OnModuleDestroy, OnModuleInit } from '@nestjs/common';
import { PrismaClient } from '@prisma/client';

// 所有走只读库的读操作方法,后续Prisma新增读方法加在这里就行
const READ_METHODS = new Set([
  'findUnique',
  'findUniqueOrThrow',
  'findFirst',
  'findFirstOrThrow',
  'findMany',
  'count',
  'aggregate',
  'groupBy',
]);

@Injectable()
export class PrismaService implements OnModuleInit, OnModuleDestroy {
  private readClient: PrismaClient;
  private writeClient: PrismaClient;
  // 缓存model代理,避免每次访问都重新生成实例
  private proxyCache = new Map<string, any>();

  constructor() {
    this.readClient = new PrismaClient({
      datasources: { db: { url: process.env.PRISMA_READ_DB_URL } },
      // 按需开启日志
      // log: ['error', 'warn'],
    });
    this.writeClient = new PrismaClient({
      datasources: { db: { url: process.env.PRISMA_WRITE_DB_URL } },
    });

    // 代理整个服务实例,拦截所有属性访问
    return new Proxy(this, {
      get: (target, prop: string) => {
        // 优先返回服务本身定义的属性、方法
        if (Object.prototype.hasOwnProperty.call(target, prop)) {
          return target[prop];
        }

        // 处理$开头的Prisma顶级API
        if (prop.startsWith('$')) {
          return this.routeTopLevelMethod(prop);
        }

        // 拦截Prisma model的访问
        if (prop in this.writeClient) {
          if (this.proxyCache.has(prop)) {
            return this.proxyCache.get(prop);
          }
          const modelProxy = this.createModelProxy(prop);
          this.proxyCache.set(prop, modelProxy);
          return modelProxy;
        }
      },
    });
  }

  async onModuleInit() {
    await Promise.all([
      this.readClient.$connect(),
      this.writeClient.$connect(),
    ]);
  }

  async onModuleDestroy() {
    await Promise.all([
      this.readClient.$disconnect(),
      this.writeClient.$disconnect(),
    ]);
  }

  /**
   * 强制走主库的入口,解决主从延迟场景下的即时读需求
   */
  get master() {
    return this.writeClient;
  }

  /**
   * 路由Prisma顶级$开头的API
   */
  private routeTopLevelMethod(method: string) {
    // 事务全走主库,保证一致性
    if (method === '$transaction') {
      return this.writeClient.$transaction.bind(this.writeClient);
    }
    // 写类原生操作全走主库
    if (method === '$executeRaw' || method === '$runCommandRaw') {
      return this.writeClient[method].bind(this.writeClient);
    }
    // 原生查询默认判断SELECT语句走从库
    if (method === '$queryRaw') {
      return (...args: any[]) => {
        const sql = args[0]?.toString?.()?.trim() || '';
        return /^select/i.test(sql)
          ? this.readClient.$queryRaw.apply(this.readClient, args)
          : this.writeClient.$queryRaw.apply(this.writeClient, args);
      };
    }
    // 连接、中间件类操作两个客户端同时执行
    if (['$connect', '$disconnect', '$on', '$use'].includes(method)) {
      return (...args: any[]) =>
        Promise.all([
          this.readClient[method].apply(this.readClient, args),
          this.writeClient[method].apply(this.writeClient, args),
        ]);
    }
    // 其他未识别的顶级方法默认走主库
    return this.writeClient[method]?.bind(this.writeClient);
  }

  /**
   * 生成单个model的路由代理
   */
  private createModelProxy(modelName: string) {
    const readModel = this.readClient[modelName];
    const writeModel = this.writeClient[modelName];

    return new Proxy(readModel, {
      get: (_, method: string) => {
        // 读操作走从库,其余全走主库
        return READ_METHODS.has(method)
          ? readModel[method].bind(readModel)
          : writeModel[method]?.bind(writeModel);
      },
    });
  }
}

使用说明

  • 常规业务代码不需要任何改动,注入PrismaService后正常调用即可:findMany/count这类读操作自动走只读副本,create/update/delete这类写操作自动走主库
  • 刚写完数据需要立即查询、避免主从同步延迟的场景,直接调用this.prisma.master.xxx.findXxx()强制走主库
  • 如果使用自定义Prisma中间件,需要分别给readClient和writeClient注册,两个是独立的客户端实例
  • 复杂原生SQL如果自动判断路由不符合预期,直接通过master入口手动指定走主库即可
  • 所有事务操作默认全走主库,不需要额外处理,避免从库延迟导致的数据一致性问题

你之前想的单独给每个model写getter的方式也能跑,但维护成本太高,每次新增model都要改服务代码,用通用代理的方案一次写完之后,新增model不需要做任何调整。Prisma 3.15版本的Client结构已经稳定,这个方案可以直接用。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.30 03:12:24