mirror of
https://gitlab.com/MrFry/mrfrys-node-server
synced 2025-04-01 20:24:18 +02:00
Added worker pools, and very very basic workers
This commit is contained in:
parent
771959ec2b
commit
f5f3b51eee
5 changed files with 211 additions and 5 deletions
141
src/utils/workerPool.ts
Normal file
141
src/utils/workerPool.ts
Normal file
|
@ -0,0 +1,141 @@
|
|||
import { Worker } from 'worker_threads'
|
||||
import genericPool from 'generic-pool'
|
||||
import os from 'os'
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
const workerFile = './src/utils/classes.ts'
|
||||
const workerTimeoutMin = 5 * 1000
|
||||
const workerTimeoutMax = 10 * 1000
|
||||
let pool: any = null
|
||||
let workers: any = null
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
// doALongTask(i, taskObj).then((res) => {
|
||||
// msgAll(res)
|
||||
// console.log('[MSGFROMCLIENT]', res)
|
||||
// })
|
||||
|
||||
export function doALongTask(i: Number, obj: any): Promise<any> {
|
||||
return new Promise((resolve) => {
|
||||
pool
|
||||
.acquire()
|
||||
.then(function(client) {
|
||||
doSomething(client, obj).then((res) => {
|
||||
resolve(res)
|
||||
// TODO: check if result is really a result, and want to release port
|
||||
pool.release(client)
|
||||
console.log('[RELEASE]: #' + client.index)
|
||||
})
|
||||
})
|
||||
.catch(function(err) {
|
||||
console.log('resourcePromise error')
|
||||
console.log(err)
|
||||
// handle error - this is generally a timeout or maxWaitingClients
|
||||
// error
|
||||
})
|
||||
})
|
||||
}
|
||||
|
||||
export function init(): void {
|
||||
if (workers && pool) {
|
||||
console.log('WORKER AND POOL ALREADY EXISTS')
|
||||
return
|
||||
}
|
||||
workers = []
|
||||
const factory = {
|
||||
create: function() {
|
||||
const worker = getAWorker(workers.length)
|
||||
workers.push(worker)
|
||||
return {
|
||||
worker: worker,
|
||||
index: workers.length,
|
||||
}
|
||||
},
|
||||
destroy: function(client) {
|
||||
console.log('[DESTROY]')
|
||||
client.worker.terminate()
|
||||
console.log('[DESTROYED] #' + client.index)
|
||||
},
|
||||
}
|
||||
|
||||
const opts = {
|
||||
min: os.cpus().length - 1, // minimum size of the pool
|
||||
max: os.cpus().length - 1, // maximum size of the pool
|
||||
maxWaitingClients: 999,
|
||||
}
|
||||
|
||||
pool = genericPool.createPool(factory, opts)
|
||||
}
|
||||
|
||||
export function msgAll(data: any): void {
|
||||
workers.forEach((worker) => {
|
||||
worker.postMessage({
|
||||
type: 'update',
|
||||
data,
|
||||
})
|
||||
})
|
||||
}
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
function getAWorker(i) {
|
||||
const worker = workerTs(workerFile, {
|
||||
workerData: {
|
||||
workerIndex: i,
|
||||
workerTimeoutMin,
|
||||
workerTimeoutMax,
|
||||
},
|
||||
})
|
||||
|
||||
worker.setMaxListeners(50)
|
||||
|
||||
// worker.on('message', (msg) => {
|
||||
// console.log(`[MAIN]: Msg from worker #${i}`, msg)
|
||||
// })
|
||||
|
||||
worker.on('online', () => {
|
||||
console.log(`[THREAD #${i}]: Worker ${i} online`)
|
||||
})
|
||||
|
||||
worker.on('error', (err) => {
|
||||
console.log(err)
|
||||
})
|
||||
|
||||
worker.on('exit', (code) => {
|
||||
console.log(`[MAIN]: worker #${i} exit code: `, code)
|
||||
})
|
||||
return worker
|
||||
}
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
function doSomething(client, obj) {
|
||||
const { index, worker } = client
|
||||
return new Promise((resolve) => {
|
||||
console.log('[ACCUIRE]: #' + index)
|
||||
worker.postMessage(obj)
|
||||
worker.once('message', (msg) => {
|
||||
resolve(msg)
|
||||
})
|
||||
})
|
||||
}
|
||||
|
||||
const workerTs = (file: string, wkOpts: any) => {
|
||||
wkOpts.eval = true
|
||||
if (!wkOpts.workerData) {
|
||||
wkOpts.workerData = {}
|
||||
}
|
||||
wkOpts.workerData.__filename = file
|
||||
return new Worker(
|
||||
`
|
||||
const wk = require('worker_threads');
|
||||
require('ts-node').register();
|
||||
let file = wk.workerData.__filename;
|
||||
delete wk.workerData.__filename;
|
||||
require(file);
|
||||
`,
|
||||
wkOpts
|
||||
)
|
||||
}
|
Loading…
Add table
Add a link
Reference in a new issue