【问题标题】:How to "multicast" an async iterable?如何“多播”异步迭代?
【发布时间】:2020-12-12 01:20:28
【问题描述】:

async 生成器能否以某种方式广播或多播,以便其所有迭代器(“消费者”?订阅者?)接收所有值?

考虑这个例子:

const fetchMock = () => "Example. Imagine real fetch";
async function* gen() {
  for (let i = 1; i <= 6; i++) {
    const res = await fetchMock();
    yield res.slice(0, 2) + i;
  }
}
const ait = gen();

(async() => {
  // first "consumer"
  for await (const e of ait) console.log('e', e);
})();
(async() => {
  // second...
  for await (const é of ait) console.log('é', é);
})();

迭代“消耗”一个值,所以只有一个或另一个得到它。 如果可以以某种方式创建这样的生成器,我希望他们两个(以及以后的任何一个)都能获得每个 yielded 值。 (类似于Observable。)

【问题讨论】:

  • 您是否可以更改设计以使gen 接受应用于每个已生成事物的回调函数?您是否有特定原因想要/需要使用生成器?
  • @kingkupps 是的gen 可以接受回调。我正在使用异步生成器,因为它们是 JS 的内置功能,否则我会为此使用 Observable,但我认为它们几乎做同样的事情(除了“多播”的能力)

标签: javascript async-await observable ecmascript-2016


【解决方案1】:

这并不容易。您需要明确地tee 它。这类似于the situation for synchronous iterators,只是稍微复杂一点:

const AsyncIteratorProto = Object.getPrototypeOf(Object.getPrototypeOf(async function*(){}.prototype));
function teeAsync(iterable) {
    const iterator = iterable[Symbol.asyncIterator]();
    const buffers = [[], []];
    function makeIterator(buffer, i) {
        return Object.assign(Object.create(AsyncIteratorProto), {
            next() {
                if (!buffer) return Promise.resolve({done: true, value: undefined});
                if (buffer.length) return buffer.shift();
                const res = iterator.next();
                if (buffers[i^1]) buffers[i^1].push(res);
                return res;
            },
            async return() {
                if (buffer) {
                    buffer = buffers[i] = null;
                    if (!buffers[i^1]) await iterator.return();
                }
                return {done: true, value: undefined};
            },
        });
    }
    return buffers.map(makeIterator);
}

您应该确保两个迭代器以大致相同的速率被消耗,这样缓冲区就不会变得太大。

【讨论】:

  • 哇,谢谢,确实不简单。希望找到一个好的库,其中包含异步迭代器的这个和其他抽象。
  • @ᆼᆺᆼ 我猜你会在带有频道的 CSP 库中找到类似的东西
【解决方案2】:

这是一个使用Highland 作为中介的解决方案。来自文档:

一个分叉到多个消费者的流将一次一个地从其源中提取值,只要最慢的消费者可以处理它们。

import _ from 'lodash'
import H from 'highland'

export function fork<T>(generator: AsyncGenerator<T>): [
    AsyncGenerator<T>,
    AsyncGenerator<T>
] {
    const source = asyncGeneratorToHighlandStream(generator).map(x => _.cloneDeep(x));
    return [
        highlandStreamToAsyncGenerator<T>(source.fork()),
        highlandStreamToAsyncGenerator<T>(source.fork()),
    ];
}

async function* highlandStreamToAsyncGenerator<T>(
    stream: Highland.Stream<T>
): AsyncGenerator<T> {
    for await (const row of stream.toNodeStream({ objectMode: true })) {
        yield row as unknown as T;
    }
}

function asyncGeneratorToHighlandStream<T>(
    generator: AsyncGenerator<T>
): Highland.Stream<T> {
    return H(async (push, next) => {
        try {
            const result = await generator.next();
            if (result.done) return push(null, H.nil);
            push(null, result.value);
            next();
        } catch (error) {
            return push(error);
        }
    });
}

用法:

const [copy1, copy2] = fork(source);

在 Node 中工作,浏览器未经测试。

【讨论】:

  • 这在浏览器中有效还是仅在 Node 上有效?
  • 我假设只有 Node,因为 highlandStreamToAsyncGenerator 是如何实现的。
【解决方案3】:

你是说这个吗?

async function* gen() {
  for (let i = 1; i <= 6; i++) yield i;
}
//const ait = gen();


(async() => {
  // first iteration
  const ait = gen()
  for await (const e of ait) console.log(1, e);
})();
(async() => {
  // second...
  const ait = gen()
  for await (const é of ait) console.log(2, é);
})();

【讨论】:

  • 是的 :) 但对不起,也许我的例子过于简单了。假设gen 例如订阅了一个服务器,所以我不想订阅两次。我将尝试更改示例以反映这一点。
猜你喜欢
  • 2016-07-05
  • 2023-01-16
  • 2014-02-04
  • 2021-12-05
  • 2019-08-03
  • 2021-10-22
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多