【问题标题】:How can I get last item of Flux without collapsing it with reduce() or last()如何在不使用 reduce() 或 last() 折叠 Flux 的情况下获得最后一项
【发布时间】:2019-11-03 22:03:03
【问题描述】:

如何在不使用 reduce() 或 last() 折叠 Flux 的情况下获得最后一项?这是我的用例:

1) 我有基于状态生成Flux<T> 的生成器。 2) 当内部Flux 完成时,它会改变影响我在生成器中发出的下一个Flux 对象的状态。

大概是这个样子

static class State {
    int secret = 2;
    int iteration = 0;
}

Random rand = new Random(1024);
Flux<Integer> stream = Flux.<Flux<Integer>, State>generate(State::new, (state, sink) -> {

    System.out.println(String.format("Generate: %d", state.secret));
    Flux<Integer> inner = Flux.range(1, rand.nextInt(10));

    sink.next(inner.doOnComplete(() -> {
        // How do I get last item of `inner` here ?
        // For example I'd like to decrement `state.secret` by last value of `inner`
    }));

    return state;
}).flatMap(Function.identity());

UPD:我没有标记我的答案,因为黑客被证明是不可靠的。 .generate() 可能会在之前的 Flux 被完全消耗之前被调用,因此导致 last 中的值不正确。

【问题讨论】:

  • 您可以在您的内部流量上调用share() 以允许多个订阅者,这样您就可以毫无问题地调用last() 吗?
  • @MichaelBerry 会不会因为源现在很热而增加流开始时丢失项目的风险?
  • 可能,是的。如果你可以cache() 整个事情是可行的,那么它并不重要,但如果不是(并且缺少元素很重要),那么这不是我推荐的方法。
  • 因为我现在没有东西可以测试它。你想试试流的物化吗?因为它会将其转换为信号+上下文它可能具有“完成”信号中的最后一个值。

标签: java project-reactor reactive-streams


【解决方案1】:

第一个版本不可靠。我又黑了一个:

static <T> Flux<T> expandOnLastItem(Supplier<Flux<T>> seed, Function<T, Flux<T>> generator) {
    return Flux.just(new AtomicReference<T>())
            .flatMap(last -> Flux.just(seed.get().materialize())
                    .flatMap(Function.identity())
                    .expand(v -> {
                        if (v.hasValue()) {
                            last.set(v.get());
                        } else if (v.isOnComplete() && last.get() != null) {
                            Flux<T> res = generator.apply(last.get());
                            last.set(null);
                            return res.materialize();
                        }
                        return Flux.empty();
                    })
                    .filter(s -> !s.isOnComplete())
                    .dematerialize());
}

可以用作

static Flux<Integer> getPage(int pageId, int size) {
    return Flux.defer(() -> {
        if (pageId < 3) {
            System.out.println("Returning data for pageId: " + pageId);
            return Flux.range(pageId * 100, size);
        } else {
            System.out.println("Returning empty for pageId: " + pageId);
            return Flux.empty();
        }
    });
}

expandOnLastItem(
        () -> getPage(0, 5),
        lastId -> {
            System.out.println("  Expanding. Last item: " + lastId);
            int curPage = lastId / 100;
            return getPage(curPage + 1, 5);
        })
        .reduce(0L, (count, value) -> {
            System.out.println("==> " + value);
            return count + 1;
        })
        .block();

所以我通过改变生成器中的状态变量来破解它。它可以工作,但不是很实用。如果其他人可以提出替代方案,我将不胜感激。

Random rand = new Random(1024);
Flux.<Flux<String>, State>generate(State::new, (state, sink) -> {

    if (state.iteration < 4) {
        final int count = rand.nextInt(10) + 1;
        System.out.println(String.format("*** Generate %d: start %d (count %d)", state.iteration, state.secret, count));
        Flux<Integer> inner = Flux.range(state.secret, count);

        final int[] last = {Integer.MIN_VALUE};
        sink.next(
                inner
                        .doOnNext(value -> {
                            last[0] = value;
                        })
                        .map(value -> String.format("Iter %d value %d", state.iteration, value))
                        .doOnComplete(() -> {
                            System.out.println(String.format("Inner complete (last item was %d)", last[0]));
                            state.secret = last[0];
                            state.iteration += 1;
                        }));
    } else {
        System.out.println("Generate complete");
        sink.complete();
    }

    return state;
})
        .flatMap(Function.identity())
        .map(value -> {
            System.out.println(String.format("Ext map: %s", value));
            return value;
        })
        .buffer(5)
        .flatMapIterable(Function.identity())
        .subscribe(value -> System.out.println(String.format("  ---> %s", value)));

System.out.println("Exiting");

【讨论】:

  • 变异变量是我想到的唯一“快速”方式,但我不想回答这个问题,因为它... hacky,正如你所展示的!至少我会在这里使用AtomicInteger 而不是原始数组。您可以通过将“hacky”内容包装在一个单独的处理器中来使其更好一点,该处理器跟踪最后一个发出的值,但它仍然会做同样的事情。如果我想到更好的解决方案,我会发布它。
  • @MichaelBerry 你能想到它坏掉的场景吗?例如。下游请求的项目是否可能比内部流中可用的项目多,从而在第一个内部流完成之前触发新的生成调用(并导致last[0] 没有最后一个元素)?
猜你喜欢
  • 1970-01-01
  • 2017-09-16
  • 2011-03-30
  • 2015-07-08
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2016-08-26
  • 2021-10-27
相关资源
最近更新 更多