From 0d75714c55e6fda0a2b1f548db0b7aafee73b0d5 Mon Sep 17 00:00:00 2001 From: Ben Date: Sat, 30 Sep 2023 08:40:22 -0400 Subject: [PATCH] batch inserts --- src/ORM/client/client.service.ts | 67 +++++++++++++++++++++++++----- src/services/stratum-v1.service.ts | 2 +- 2 files changed, 58 insertions(+), 11 deletions(-) diff --git a/src/ORM/client/client.service.ts b/src/ORM/client/client.service.ts index 0399dd1..c952843 100644 --- a/src/ORM/client/client.service.ts +++ b/src/ORM/client/client.service.ts @@ -1,6 +1,8 @@ import { Injectable } from '@nestjs/common'; import { InjectRepository } from '@nestjs/typeorm'; -import { Repository } from 'typeorm'; +import { BehaviorSubject, lastValueFrom } from 'rxjs'; +import { FindOptionsWhere, ObjectId, Repository } from 'typeorm'; +import { QueryDeepPartialEntity } from 'typeorm/query-builder/QueryPartialEntity'; import { ClientEntity } from './client.entity'; @@ -9,13 +11,54 @@ import { ClientEntity } from './client.entity'; @Injectable() export class ClientService { - + private insertQueue: { client: Partial, subject: BehaviorSubject }[] = []; + private heartbeatQueue: { + criteria: string | string[] | number | number[] | Date | Date[] | ObjectId | ObjectId[] | FindOptionsWhere, + partialEntity: QueryDeepPartialEntity, + subject: BehaviorSubject + }[] = []; constructor( @InjectRepository(ClientEntity) private clientRepository: Repository ) { + setInterval(async () => { + if (this.heartbeatQueue.length < 1) { + return; + } + const heartbeatToInsert = [...this.heartbeatQueue]; + this.heartbeatQueue = []; + + + await this.clientRepository.manager.transaction(async transactionalEntityManager => { + for (let i = 0; i < heartbeatToInsert.length; i++) { + await this.clientRepository.update(heartbeatToInsert[i].criteria, heartbeatToInsert[i].partialEntity); + } + }); + + + heartbeatToInsert.forEach((val, index) => { + heartbeatToInsert[index].subject.next(null); + heartbeatToInsert[index].subject.complete(); + }); + + }, 5000); + + setInterval(async () => { + if (this.insertQueue.length < 1) { + return; + } + const clientsToInsert = [...this.insertQueue]; + this.insertQueue = []; + const results = await this.clientRepository.insert(clientsToInsert.map(i => i.client)); + + results.generatedMaps.forEach((val, index) => { + clientsToInsert[index].subject.next({ ...clientsToInsert[index].client, ...val }); + clientsToInsert[index].subject.complete(); + }); + + }, 5000); } @@ -31,7 +74,15 @@ export class ClientService { } public async heartbeat(address: string, clientName: string, sessionId: string, hashRate: number) { - return await this.clientRepository.update({ address, clientName, sessionId }, { hashRate, deletedAt: null, updatedAt: new Date() }); + + const subject$ = new BehaviorSubject(null); + this.heartbeatQueue.push({ + criteria: { address, clientName, sessionId }, + partialEntity: { hashRate, deletedAt: null, updatedAt: new Date() }, + subject: subject$ + } + ); + return lastValueFrom(subject$); } // public async save(client: Partial) { @@ -40,14 +91,10 @@ export class ClientService { public async insert(partialClient: Partial): Promise { - const insertResult = await this.clientRepository.insert(partialClient); - const client = { - ...partialClient, - ...insertResult.generatedMaps[0] - }; - - return client as ClientEntity; + const subject$ = new BehaviorSubject(null); + this.insertQueue.push({ client: partialClient, subject: subject$ }); + return lastValueFrom(subject$); } public async delete(sessionId: string) { diff --git a/src/services/stratum-v1.service.ts b/src/services/stratum-v1.service.ts index edc2d60..68f529e 100644 --- a/src/services/stratum-v1.service.ts +++ b/src/services/stratum-v1.service.ts @@ -32,7 +32,7 @@ export class StratumV1Service implements OnModuleInit { if (process.env.ENABLE_SOLO == 'true') { //await this.clientStatisticsService.deleteAll(); - // await this.clientService.deleteAll(); + await this.clientService.deleteAll(); this.startSocketServer(); } }