Blizzard / Blizzard/node-rdkafka
Reading local file through ReadStream is very slow when a lot of consumer streams are created.
- Dominant language
- JavaScript
- Stars
- 2.2k
- Forks
- 403
- PR merge metrics
- No merged PRs in 30d
Description
Reading local file through ReadStream is very slow when a lot of consumer streams are created, I don't know why.
**Environment Information**
- OS: Windows 11
- Node Version: 16.14.0
- NPM Version: 8.3.1
- C++ Toolchain:
- node-rdkafka version: 2.12.0
**node-rdkafka Configuration Settings**
**Additional context**
**Steps to Reproduce**
package.json
```json
{
"name": "node-rdkafka-test",
"version": "1.0.0",
"main": "index.js",
"license": "MIT",
"dependencies": {
"debug": "^4.3.3",
"node-rdkafka": "^2.12.0"
},
"scripts": {
"start": "node index.js"
}
}
```
logger.js
```javascript
process.env.DEBUG = 'info:*,warn:*,error:*,debug:*,trace:*'
process.env.DEBUG_COLORS = 'true'
const debug = require('debug')
const dInfo = debug('info')
const dErr = debug('error')
const dDebug = debug('debug')
const dWarn = debug('warn')
const dTrace = debug('trace')
module.exports = (ns) => ({
info: dInfo.extend(ns),
err: dErr.extend(ns),
debug: dDebug.extend(ns),
warn: dWarn.extend(ns),
trace: dTrace.extend(ns),
})
```
kafka.js
```javascript
const { KafkaConsumer } = require('node-rdkafka')
const logger = require('./logger')
const consumerLogger = logger('kafka')
const createReadStream = ({ topics, brokerList, groupId }) => new Promise((resolve) => {
const stream = KafkaConsumer.createReadStream(
{ 'metadata.broker.list': brokerList, 'group.id': groupId },
{ 'auto.offset.reset': 'largest' },
{ topics }
)
stream.consumer.on('subscribed', (topics) => {
consumerLogger.info('kafka read stream subscribed on', topics)
resolve(stream)
})
})
const createSingleReadStream = createReadStream
const createMultipleReadStreams = ({ topics, brokerList, groupId }) => {
const configs = topics.map((topic) => ({ topics: [ topic ], brokerList, groupId }))
return Promise.all(configs.map(createReadStream))
}
module.exports = { createSingleReadStream, createMultipleReadStreams }
```
file.js
```javascript
const readline = require("readline")
const fs = require("fs")
const { once } = require("events")
const logger = require('./logger')
const log = logger('file')
const readFile = async (filePath) => {
log.info('reading start')
const rl = readline.createInterface({ input: fs.createReadStream(filePath), crlfDelay: Infinity })
let count = 0
const lines = []
rl.on('line', (line) => {
log.info('reading line', ++count)
lines.push(line)
})
await once(rl, 'close')
log.info('reading end')
return lines
}
module.exports = { readFile }
```
test.log
```log
line 1
line 2
line 3
line 4
line 5
line 6
line 7
line 8
line 9
line 10
```
index.js
```javascript
const topics = [
'kxb-pay',
'xa-mongo-event',
'xa-chain-event',
'kxb-user-got-mobile',
'kxb-merchant-information-submission',
'kxb-merchant-information-approved',
'kxb-merchant-information-rejected',
'kxb-merchant-recharge',
'kxb-merchant-annual-fee-deduction',
'kxb-merchant-trial',
'kxb-video-rejected',
]
const brokerList = '10.88.1.13:9092'
const groupId = 'node-rdkafka-test'
const { createSingleReadStream, createMultipleReadStreams } = require('./kafka')
const { readFile } = require('./file')
Promise.resolve({ topics, brokerList, groupId })
.then(createSingleReadStream)
// .then(createMultipleReadStreams)
.then(() => readFile('./test.log'))
```
when use createSingleReadStream, reading file is fast, note the first `line` event's received time:

when use createMultipleReadStreams, reading file becomes very slow:

Contributor guide
Research direction
Start with kafka.js, especially createReadStream and createMultipleReadStreams, then inspect file.js where readline consumes fs.createReadStream. Reproduce the timing difference using index.js, the listed Node.js and node-rdkafka versions, and test.log; done means identifying and resolving why multiple consumer streams delay the first file line.
Written by the indexing model from the issue text.
Assessment
- Tech stack
- javascript, kafka, node.js
- Domain
- backend, distributed-systems
- Issue type
- Bug
- Difficulty
- 4/5
- Estimated time
- 3-5 days
- Activity status
- Stale
- Clarity
- Needs clarification
- Newbie friendliness
- 35/100