Blizzard / Blizzard/node-rdkafka

Reading local file through ReadStream is very slow when a lot of consumer streams are created.

Open
#938 0 comments 0 reactions 0 assignees View on GitHub
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:
![image](https://user-images.githubusercontent.com/71441807/154434177-8be4bde6-88ce-40da-8885-4ae2639b903b.png)

when use createMultipleReadStreams, reading file becomes very slow:
![image](https://user-images.githubusercontent.com/71441807/154434091-015234fa-f400-49ef-9d09-10fbae05be58.png)

Contributor guide

Open the contributing 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

Get new issues in your inbox

A short digest of beginner-friendly GitHub issues.