Usamos el siguiente código para transmitir los resultados de una consulta al cliente:
app.get('/events', (req, res) => { try { const stream = db('events') .select('*') .where({ id_user: 'foo' }) .stream() stream.pipe(JSONStream.stringify()).pipe(res) } catch (err) { next(err) } })Si bien el código parece tener un perfil de uso de memoria excelente (uso de memoria estable/bajo), crea tiempos de espera de adquisición de conexión de base de datos aleatorios:
Knex: Tiempo de espera para adquirir una conexión. Probablemente la piscina esté llena. ¿Te estás perdiendo una llamada .transacting(trx)?
Esto sucede en la producción a intervalos aparentemente aleatorios. ¿Alguna idea de por qué?
Esto sucede porque las solicitudes anuladas (es decir, el cliente cierra el navegador a mitad de la solicitud) no liberan la conexión al grupo.
Primero, asegúrese de tener la última versión de knex; o al menos v0.21.3+ que ha introducido correcciones en el manejo de secuencias/grupos .
A partir de ahora tienes un par de opciones:
Utilice stream.pipeline en lugar de stream.pipe , que maneja las solicitudes abortadas correctamente de la siguiente manera:
const { pipeline } = require('stream') app.get('/events', (req, res) => { try { const stream = db('events') .select('*') .where({ id_session: req.query.id_session }) .stream() return pipeline(stream, JSONStream.stringify(), res, err => { if (err) { return console.log(`Pipeline failed with err:`, err) } console.log(`Pipeline ended succesfully`) }) } catch (err) { next(err) } }) o escuche el evento [ close ][close] a req y destruya el flujo de DB usted mismo, así:
app.get('/events', (req, res) => { try { const stream = db('events') .select('*') .where({ id_session: req.query.id_session }) .stream() // Not listening to this event will crash the process if // stream.destroy(err) is called. stream.on('error', () => { console.log('Stream was destroyed') }) req.on('close', () => { // stream.end() does not seem to work, only destroy() stream.destroy('Aborted request') }) stream.pipe(JSONStream.stringify()).pipe(res) } catch (err) { next(err) } })stream.end mencionado aquí no parece funcionar.