正如 TJ 所建议的,您不能用 single 承诺来表示订阅。我使用了一种类似于 TJ 链接中描述的技术,但在抽象方面有显着差异。下面的 duplexStream 为调用者提供了一个 read 函数,用于您程序中的任何流观察者,以及一个 write 函数以在您的可观察订阅中使用 -
function duplexStream () {
let t = defer()
async function* read () {
while (true) yield await t.deferred
}
function write (err, value) {
if (err) t.reject(err)
else t.resolve(value)
t = defer()
}
return [read, write]
}
defer 是一个提供外部控制承诺的简单抽象 -
function defer () {
let resolve, reject
return { deferred: new Promise((res, rej) => (resolve = res, reject = rej)), resolve, reject }
}
假设我们有两个 HTML 元素 -
<p id="foo"></p>
<p id="bar"></p>
我们现在可以使用 duplexStream 将 observable 转换为异步迭代器。注意read 可以被多次调用,所有读者都将获得传递给write 的值。 write 也可以赋予任意数量的 observables。 duplexStream 是通用的,不绑定到特定的库或类。通过将read 和write 作为通用函数提供给调用者,可以在普通的subscribe 和unsubscribe 事件期间管理这些处理程序-
const [read, write] = duplexStream()
async function update(elem, it) {
for await (const value of it)
elem.innerHTML += (value + "<br>")
}
update(document.querySelector("#foo"), read()).catch(console.error)
update(document.querySelector("#bar"), read()).catch(console.error)
A.subscribe(write)
对于这个演示,我们将 A.subscribe 编写为一个模拟 observable,它发出三 (3) 个值,然后是一个错误 -
const A = {
subscribe (f) {
setTimeout(f, 1000, null, "one")
setTimeout(f, 2000, null, "two")
setTimeout(f, 3000, null, "three")
setTimeout(f, 4000, Error("something bad"))
}
}
程序完成后,我们会看到-
<p id="foo">
one<br>
two<br>
three<br>
</p>
<p id="bar">
one<br>
two<br>
three<br>
</p>
并且错误将在console.error 中记录两次,因为每个流阅读器catched -
Error: something bad
Error: something bad
展开下面的sn-p,在自己的浏览器中运行程序验证结果-
function duplexStream () {
let t = defer()
async function* read () {
while (true) yield await t.deferred
}
function write (err, value) {
if (err) t.reject(err)
else t.resolve(value)
t = defer()
}
return [read, write]
}
function defer () {
let resolve, reject
return { deferred: new Promise((res, rej) => (resolve = res, reject = rej)), resolve, reject }
}
const A = {
subscribe (f) {
setTimeout(f, 1000, null, "one")
setTimeout(f, 2000, null, "two")
setTimeout(f, 3000, null, "three")
setTimeout(f, 4000, Error("something bad"))
}
}
const [read, write] = duplexStream()
async function update(elem, it) {
for await (const value of it)
elem.innerHTML += (value + "<br>")
}
update(document.querySelector("#foo"), read()).catch(console.error)
update(document.querySelector("#bar"), read()).catch(console.error)
A.subscribe(write)
.as-console-wrapper { max-height: 33% !important; }
<p id="foo"></p>
<p id="bar"></p>