apache / apache/pulsar-client-node
Interruptible Reader.readNext()?
- Dominant language
- C++
- Stars
- 164
- Forks
- 98
- PR merge metrics
- No merged PRs in 30d
Description
Hi,
I have a use case where I expose a Pulsar topic over HTTP via Server-Sent Events. Basically, when a client connects over HTTP, I do this:
```js
const reader = await client.createReader({
topic: request.params.topic,
startMessageId: request.headers['last-event-id'] ?
Pulsar.MessageId.deserialize(Buffer.from(request.headers['last-event-id'], 'base64')) :
Pulsar.MessageId.earliest()
});
```
Then I use a loop that reads messages as they come and sends them to the client:
```js
while (!clientGoneAway) {
let message;
message = await reader.readNext();
reply.raw.write(formatServerSentEvent(
message.getMessageId().serialize().toString('base64'),
message.getData().toString('utf-8')
));
}
```
Additionally, the server detects the `close` event on the HTTP requests, closes the reader and prevents further iteration:
```js
request.raw.once('close', async function() {
clientGoneAway = true;
await reader.close();
reply.raw.end();
})
```
(`request.raw` and `reply.raw` are Node.js req and res objects, respectively - they're just wrapped like this in Fastify.js)
Now, my problem is that even if I call `reader.close()`, the `reader.readNext();` never resolves nor rejects. It's not just a Promise problem - it seems like it's keeping a thread busy, because then all other operations hang: things like `fs.createReadStream`, as well as creating new readers, hang forever until I completely restart the Node.js process.
I know I can use a timeout with `reader.readNext(timeoutMS)`, but this has 2 major disadvantages:
* It turns the reader into a kind of poller
* It still does not vacate the thread - so it's possible to trivially saturate the thread pool by creating more readers than the pool size within the timeout period (so for example 4 readers in 1 second, when using a timeout of 1000 ms)
Is there any way to have the reader immediately abort all reads when closed?
Contributor guide
No contributing guide indexed for this repository
Research direction
Start at the Reader API entry points used by reader.readNext() and reader.close(), then reproduce the close-while-read case with the Node.js client. Trace whether closing the reader interrupts an outstanding read and check the related asynchronous behavior. Done means a blocked read resolves or rejects promptly after close without preventing other Node.js operations.
Written by the indexing model from the issue text.
Assessment
- Tech stack
- javascript, node.js
- Domain
- backend
- Issue type
- Bug
- Difficulty
- 4/5
- Estimated time
- 3-5 days
- Activity status
- Stale
- Clarity
- Mostly clear
- Newbie friendliness
- 38/100