forked from brianc/node-pg-query-stream
/
index.js
65 lines (60 loc) · 1.56 KB
/
index.js
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
var util = require('util')
var Cursor = require('pg-cursor')
var Readable = require('readable-stream').Readable
var QueryStream = module.exports = function(text, values, options) {
var self = this
this._reading = false
this._closing = false
options = options || { }
Cursor.call(this, text, values)
Readable.call(this, {
objectMode: true,
highWaterMark: options.highWaterMark || 1000
})
this.batchSize = options.batchSize || 100
this.once('end', function() {
process.nextTick(function() {
self.emit('close')
})
})
}
util.inherits(QueryStream, Readable)
//copy cursor prototype to QueryStream
//so we can handle all the events emitted by the connection
for(var key in Cursor.prototype) {
if(key == 'read') {
QueryStream.prototype._fetch = Cursor.prototype.read
} else {
QueryStream.prototype[key] = Cursor.prototype[key]
}
}
QueryStream.prototype.close = function() {
this._closing = true
var self = this
Cursor.prototype.close.call(this, function(err) {
if(err) return self.emit('error', err)
process.nextTick(function() {
self.push(null)
})
})
}
QueryStream.prototype._read = function(n) {
if(this._reading || this._closing) return false
this._reading = true
var self = this
this._fetch(this.batchSize, function(err, rows) {
if(err) {
return self.emit('error', err)
}
if(!rows.length) {
process.nextTick(function() {
self.push(null)
})
return
}
self._reading = false
for(var i = 0; i < rows.length; i++) {
self.push(rows[i])
}
})
}