如何使用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
相关产品推荐
相关产品推荐

