>I want to announce that there is more work, then immediately hand off the task to a worker until I run out of workers, which I believe is known as backpressure. In which case, requests will queue up until a worker is free. Am I understanding this correctly?

This is by its very nature an async hand-off between a producer and consumers. With a push-style stream the availability of a new element causes the immediate processing of said element. It isn't the availability of workers pulling in available work.

On the topic of "closable" BlockingQueues, it is something which I've thought about and experimented with a lot, and there's a ton of nuance there: what available strategies for dealing with a full queue should exist, block, drop-oldest, drop-newest, drop-all, etc. Also, should workers be able to signal termination such that the producer doesn't end up producing to a queue where all workers have terminated? Is it always first-come-first-served in terms of consumption or does it support "broadcast"?

On 2026-08-27 22:42, David Alayachew wrote:
Hmmm, this looks more like async-style, where it's pub-sub. That's not really what I am looking for.

What I am asking for is inherently push-based, assuming that I understand the terminology (basing myself off of this -- https://stackoverflow.com/questions/51254117/what-is-difference-between-push-based-and-pull-based-structures-like-ienumerable <https://urldefense.com/v3/__https://stackoverflow.com/questions/51254117/what-is-difference-between-push-based-and-pull-based-structures-like-ienumerable__;!!ACWV5N9M2RV99hQ!MRNMRsm0f3rRv6WFvB4G2dQcwa_aPT_GGxQ0l0IMKow3rwkzSfQpE7rfQ3bkH2FKjOM2GoSS7daDbQD07Wdq3LBPomQ$>). I want to announce that there is more work, then immediately hand off the task to a worker until I run out of workers, which I believe is known as backpressure. In which case, requests will queue up until a worker is free. Am I understanding this correctly?

And let me ask a way more fundamental question -- why can't I use the Stream framework for an upstream that changes its contents? For me, if I could do that -- toss items onto my upstream collection as I am iterating, and then signal that I am done, then this entire problem goes away. I guess let's start with that, as I feel like me not understanding why that is in place prevents me from understanding.

And I am not at all implying that every collection should be able to do that. Certainly not. But if there was even one collection implementation that had a stream method, and that method allowed me to communicate that there are no more elements, and specifically, that there won't be anymore elements, then this entire problem goes away.


On Thu, Aug 27, 2026 at 4:58 AM Viktor Klang <[email protected]> wrote:

    Thanks for your kind words, David!

    Perhaps I focused too much on the "recursive"-side of your question?

    What you describe seems more amendable to a pull-style stream with
    an async hand-off (allowing the upstream to progress
    independently, up to a limit) in order to allow to hide the
    latency of generating/retriving the upstream values. This specific
    scenario was exactly the motivating reason behind
    /java.util.concurrent.Flow/. There are inherent trade-offs between
    *push*-style streams (/java.util.stream.Stream/) and *pull*-style
    streams (/java.util.concurrent.Flow/) and sometimes either of
    those approaches are preferable given the desired runtime
    characteristics.

    What you could experiment with is the style I shared previously
    and pair it with Gatherers.mapConcurrent() for the subsequent
    processing of each element (Path), or you could try using
    mapConcurrent as inspiration for how to implement
    flatMapConcurrent to concurrently (up to a limit) traverse the
    hierarchy?

    Let me know if that is more aligned with what you're thinking here.

    On 2026-08-27 05:00, David Alayachew wrote:
    Lol, that is so creative. I never would have thought to use
    anonymous classes on the java. util. function interfaces lol.
    That's like choosing the horse when you have a motorcycle right
    in front of you. But lo and behold, it genuinely is the
    Lol, that is so creative.

    I never would have thought to use anonymous classes on the
    java.util.function interfaces lol. That's like choosing the horse
    when you have a motorcycle right in front of you. But lo and
    behold, it genuinely is the tersest way forward using normal Java.

    How did you come up with this? lol

    Anyways, creativeness aside, this one had about the same
    performance as the Gatherer that I came up with. I am wondering
    if there might be too much hand off, or something performance
    wise that I might be doing wrong here. Long story short is that,
    I want to gather the files as fast as possible, but I don't see a
    way to get the full extent of the parallelism out without keeping
    it all on one pipeline. At the end of the day, the solution that
    you made seems to be making all of these instances of Stream over
    and over again, and I worry if that might be making things slower
    than needed. The instance of Gatherer I made doesn't have this
    problem, but that's because it has the opposite problem of not
    being able to parallelize easily. Yours looks like parallelizes
    easily enough, but also does a lot of work to make that happen.

    Maybe this is imaginative of me, but I am thinking of some sort
    of Collection type, which doesn't suffer from
    ConcurrentModificationException, where the Stream can just keep
    pulling from it until it gets told that no more elements are
    coming through, and it can pull and process the elements from
    that Collection type in parallel. This is my dream object, if you
    will. You just keep appending to the end of it as you come across
    directories, and you keep pulling from the start of it until the
    collection object finally becomes empty, in which case, you are
    finally done.

    Does any of that make sense to you?


    On Wed, Aug 26, 2026, 11:26 AM Viktor Klang
    <[email protected]> wrote:

        Hello David,

        You mean something like this?

        void main() throws Exception {
            final Path root = Path.of(System.getProperty("user.home"));

            IO.println(root);
            IO.println(root.toAbsolutePath());

            Stream
               .of(root)
               .flatMap(new Function<Path, Stream<Path>>() {
                  @Override public Stream<Path> apply(Path p) {
                     if (Files.isRegularFile(p)) {
                     return Stream.of(p);
                  } else {
                     try {
                        return Files.list(p).flatMap(this::apply);
                     } catch (IOException _) {
                        return Stream.empty();
                     }
                  }
                 }
               })
               .limit(5)
               .forEach(IO::println);
        }

        On 2026-08-26 16:24, David Alayachew wrote:
        > Hello Viktor and Rémi,
        >
        > First off, sorry for the horrifically delayed response.
        Juggling disasters
        > and emergencies.
        >
        > Rémi, thanks for the example with mapMulti. That has been
        my temporary
        > workaround for now, and while it is still not ideal, it is
        better than
        > where I was before.
        >
        > Viktor, sure, here is a simple example -- traversing a
        directory tree, and
        > only passing the files down the stream.
        >
        > I actually made a post on StackOverflow --
        >
        https://softwareengineering.stackexchange.com/questions/461442/
        
<https://urldefense.com/v3/__https://softwareengineering.stackexchange.com/questions/461442/__;!!ACWV5N9M2RV99hQ!JsL3fQJ1ult_TN0D6BCiRDGewOWQs2G_97ksDE0txV_GsNtGvjC1XLvjm8-rVEgGNQKIOX_5Xz7VRezKpGpwi9bbnPg$>
        >
        > But anyways, when attempting to do this with streams, this
        is where I
        > started.
        >
        >
        > import module java.base;
        >
        > void main() throws Exception
        > {
        >
        >     final Path root = Path.of(System.getProperty("user.home"));
        >
        >     IO.println(root);
        >     IO.println(root.toAbsolutePath());
        >
        >     Stream
        >        .of(root)
        >        .mapMulti(this::recursiveDescent)
        >        .limit(5)
        >        .forEach(IO::println)
        >        ;
        >
        > }
        >
        > private void recursiveDescent(final Path rootPath, final
        Consumer<Path>
        > downstream)
        > {
        >
        >     final Stack<Path> stack = new Stack<>();
        >     stack.push(rootPath);
        >
        >     while (!stack.empty())
        >     {
        >
        >        final Path path = stack.pop();
        >
        >        if (Files.isRegularFile(path))
        >        {
        >
        >           downstream.accept(path);
        >
        >        }
        >
        >        else
        >        {
        >
        >           try (final Stream<Path> folderContents =
        Files.list(path))
        >           {
        >
        > folderContents.forEach(stack::push);
        >
        >           }
        >
        >           catch (final Exception exception)
        >           {
        >
        >              throw new IllegalStateException("Failed for "
        + path,
        > exception);
        >
        >           }
        >
        >        }
        >
        >     }
        >
        > }
        >
        > But this has at least 2 major downsides.
        >
        > 1 - This is not easy to turn parallel (comparatively).
        >
        > 2 - This does not short-circuit when the downstream no
        longer accepts
        > elements.
        >
        > Ok, I can at least solve problem 2 by becoming a Gatherer
        instead.
        >
        > Here is my gatherer attempt.
        >
        >
        > import module java.base;
        >
        > void main() throws Exception
        > {
        >
        >     final Path root = Path.of(System.getProperty("user.home"));
        >
        >     IO.println(root);
        >     IO.println(root.toAbsolutePath());
        >
        >     final Gatherer<Path, Stack<Path>, Path> gatherer =
        >        Gatherer
        >        .of
        >        (
        >        Stack<Path>::new,
        >              (stack, rootPath, downstream) ->
        >              {
        >
        >                 stack.push(rootPath);
        >
        >                 while (!stack.isEmpty())
        >                 {
        >
        >                    final Path path = stack.pop();
        >
        >                    if (Files.isRegularFile(path))
        >                    {
        >
        >                       final boolean acceptingMoreElements =
        > downstream.push(path);
        >
        >                       if (!acceptingMoreElements)
        >                       {
        >
        >                          return false;
        >
        >                       }
        >
        >                    }
        >
        >                    else
        >                    {
        >
        >                       try (final Stream<Path> folderContents =
        > Files.list(path))
        >                       {
        >
        > folderContents.forEach(stack::push);
        >
        >                       }
        >
        >                       catch (final Exception exception)
        >                       {
        >
        >                          throw new
        IllegalStateException("Failed for " +
        > path, exception);
        >
        >                       }
        >
        >                    }
        >                 }
        >
        >                 return true;
        >
        >              },
        >              (s1, s2) ->
        >              {
        >
        >                 s1.addAll(s2);
        >                 return s1;
        >
        >              },
        >              (stack, downstream) ->
        >              {
        >
        >                 for (final Path path : stack)
        >                 {
        >
        >                    if (!downstream.push(path))
        >                    {
        >
        >                       return;
        >
        >                    }
        >
        >                 }
        >
        >              }
        >        )
        >        ;
        >
        >     Stream
        >        .of(root)
        >        .gather(gatherer)
        >        .limit(5)
        >        .forEach(IO::println)
        >        ;
        >
        > }
        >
        > So, problem 2 is solved, but problem 1 is not really. Sure,
        I could turn my
        > stream parallel, but the actual meat of the processing is
        sequential when
        > it really doesn't need to be.
        >
        > Let me know if this makes more sense. And sorry, you might
        have to scroll
        > to read earlier emails in this thread to get the context. I
        know it was
        > several months back.
        >
        >
        > On Fri, Nov 14, 2025 at 5:42 AM Remi Forax
        <[email protected]> wrote:
        >
        >> Hi David,
        >> You can always transform an imperative code to a stream by
        pushing the
        >> element through a consumer.
        >> Internally, a stream uses a push iterator (see
        >> Spliterator.tryAdvance(consumer)).
        >>
        >> As a silly example, this is a way to write fibonacci (the
        recursive form)
        >> with a stream right in the middle.
        >>
        >>
        >> static void fibo(int n, IntConsumer consumer) {
        >>    if (n < 2) {
        >>      consumer.accept(n);
        >>      return;
        >>    }
        >>    var result = Stream.of("")
        >>        .mapMultiToInt((_, consumer2) -> {
        >>          fibo(n - 1, consumer2);
        >>          fibo(n - 2, consumer2);
        >>        })
        >>        .sum();
        >>    consumer.accept(result);
        >> }
        >>
        >> static void main() {
        >>    fibo(7, IO::println);
        >> }
        >>
        >>
        >> Here, I use mapMulti() to convert the imperative code to a
        Stream
        >> (there is no factory method on Stream that takes a
        consumer of consumer).
        >>
        >> If you also want to short-circuit, you can use a gatherer
        instead of
        >> mapMulti but short-circuiting the recursive code will
        require to use an
        >> exception as control flow (it will not be pretty).
        >>
        >> regards,
        >> Rémi
        >>
        >>
        >> ------------------------------
        >>
        >> *From: *"David Alayachew" <[email protected]>
        >> *To: *"core-libs-dev" <[email protected]>
        >> *Sent: *Tuesday, November 11, 2025 4:36:29 AM
        >> *Subject: *Difficulties of recursion with Streams
        >>
        >> Hello @core-libs-dev <[email protected]>,
        >>
        >> When working with streams, I often run into situations
        where I have to
        >> "demote" back to imperative code because I am trying to
        solve a problem
        >> best solved by recursion.
        >>
        >> Consider the common use case of cycling through
        permutations to find all
        >> permutations that satisfy some condition. With recursion,
        the answer is
        >> incredibly simple -- just grab an element from the set,
        then call the
        >> recursive method with a copy of the set minus the grabbed
        element. Once you
        >> reach the empty set, you've reached your terminal condition.
        >>
        >> Use cases like that are not only incredibly common, but
        usually,
        >> embarrassingly parallel. The example above of cycling
        through permutations
        >> is only a few lines of imperative code, but I struggle to
        imagine how I
        >> would do this with Streams.
        >>
        >> I guess let me start by asking -- are there any good ways
        currently to
        >> accomplish the above permutation example with Streams? And
        if not, should
        >> there be?
        >>
        >> Thank you for your time and consideration.
        >> David Alayachew
        >>
        >>
-- Cheers,
        √


        Viktor Klang
        Software Architect, Java Platform Group
        Oracle

-- Cheers,
    √


    Viktor Klang
    Software Architect, Java Platform Group
    Oracle

--
Cheers,
√


Viktor Klang
Software Architect, Java Platform Group
Oracle

Reply via email to