Repository navigation
for await & Readable #29428
Description
Activity
I think this is probably correct (although probably surprising) since readable is supposed to be greedy.
In order to get the behaviour you want you need to set
highWaterMarkto1, e.g.for await (const d of Readable.from(generate(), { highWaterMark: 1 })) { console.log(d) }
Reacted by Benjamin Gruenbaum and Naor Tedgi (Abu Emma)@mcollina: Thoughts? In the given example I believe the behaviour is correct since the generator is
async.However, using a sync generator has the same behavior and I'm not sure that is correct?
- addedstreamIssues and PRs related to Node.js streams.Issues and PRs related to Node.js streams.
on Sep 5, 2019 We might need to think about it a bit more. Look at the following:
async function* generate() { yield 1 yield 2 throw new Error('Boum') } ;(async () => { try { for await (const d of generate()) { console.log(d) } } catch (e) { console.log(e) } })()
The result is:
1 2 BoumI think this might be something we can fix. Specifically, the following is making the iterator exit eagerly:
node/lib/internal/streams/async_iterator.js
Lines 61 to 86 in 63b056d
// If we have detected an error in the meanwhile // reject straight away. const error = this[kError]; if (error !== null) { return Promise.reject(error); } if (this[kEnded]) { return Promise.resolve(createIterResult(undefined, true)); } if (this[kStream].destroyed) { // We need to defer via nextTick because if .destroy(err) is // called, the error will be emitted via nextTick, and // we cannot guarantee that there is no error lingering around // waiting to be emitted. return new Promise((resolve, reject) => { process.nextTick(() => { if (this[kError]) { reject(this[kError]); } else { resolve(createIterResult(undefined, true)); } }); }); } I think we might want to put the promise rejection in the
kLastPromiseproperty queue instead:.node/lib/internal/streams/async_iterator.js
Line 155 in 63b056d
iterator[kLastPromise] = null; @benjamingr what do you think?
On second thought, I might be wrong, and this might be due for something else, or something that we cannot really fix.
Streams are greedy into emitting
'error'. They emit an error as soon as it happens. In, we see that we callLines 1223 to 1236 in 63b056d
async function next() { try { const { value, done } = await iterator.next(); if (done) { readable.push(null); } else if (readable.push(await value)) { next(); } else { reading = false; } } catch (err) { readable.destroy(err); } } destroyeagerly and this is what triggers it.As a short term measure, we should definitely document it.
@benjamingr what do you think?
That streams do watermarking and that they are greedy :]
To avoid the stream behaviour one can simple not convert their AsyncIterator to a stream.
Now I do believe we can "fix" this (on the async iterator side) by only emitting the error and destroying after we are done emitting data on the iterator but honestly I am not sure that would be better from a stream consumer PoV. In fact it would likely be worse.
So +1 on a docs change, -0.5 on changing the
Symbol.asyncIterator.Now I do believe we can "fix" this (on the async iterator side) by only emitting the error and destroying after we are done emitting data on the iterator
I think I’d be in favour of that, fwiw.
I think I’d be in favour of that, fwiw.
That was my first intuition too but I am having a hard time coming up with a real use-case in which this (buffering) behaviour is better. Namely, in all cases I considered I would rather recover from the error sooner rather than later. The cases I came up with were:
- Persistent database change stream (I would rather re-connect and sync which I have to do anyway to listen to new changes).
- Query results iterated with a cursor (if the query failed on some results, I would rather retry the query sooner rather than later).
- An http body stream (I would need to recover and marginally find the incomplete data useful).
Edit: the theme of these is having to reconnect to establish synchronisation with the producer later on anyway.
On the other hand if I look at async iterators without streams, such use cases are easy to come up with and are abundant in reactive systems (like the whole server flow being accepting requests from clients in a for...await )
The biggest issue I see here is that there is missing data that we are dropping that the user might be interested in. I am also not sure how else we can expose it. Basically the constraints are:
- As a consumer if there is an error downstream I want to know about it as soon as possible.
- As a consumer I never want to miss any data on the stream.
The biggest issue I see here is that there is missing data that we are dropping that the user might be interested in. I am also not sure how else we can expose it. Basically the constraints are:
- As a consumer if there is an error downstream I want to know about it as soon as possible.
- As a consumer I never want to miss any data on the stream.
Yeah, that’s what I’m thinking too. I’d also consider it more expected for (async) iterators to throw the exception only after providing data that was produced before the error, because in the typical (async) generator function approach, that’s when that exception would be have been created.
Note that we can consider this to be a problem for our
Readable.from()implementation rather than our async iterator support. Essentially, defer callingdestroy()until all chunks have been emitted. This can be easily achieved by settinghighWaterMark: 1, and we may just set that as the default.Reacted by Robert Nagy and Benjamin GruenbaumNote that running this other snippet produces a different output
function makeStreamV2() { return new Readable({ objectMode: true, read() { if (this.reading) { return } this.reading = true this.push(1) this.push(2) this.emit('error', 'boum') } }) } ;(async () => { try { for await (const d of makeStreamV2()) { console.log(d) } } catch (e) { console.log(e) } })()
1 boum@mcollina did we resolve this or is this something someone needs to further look into?
This needs some decision to be made, and possibly some docs to be added, or the behavior changed.
Reacted by Robert Nagy- added a commit that references this issue
on May 7, 2020 - added a commit that references this issue
on Jun 7, 2020 - added a commit that references this issue
on Jul 27, 2026
Hi, I am trying to consume a readable with a
for awaitloop.When I create the simple stream below, the console output the error straight away instead of logging the data events first. Is it the expected behavior
Version: v12.9.1
Platform: Darwin Kernel Version 18.7.0
the output is
instead of the expected