Then maybe I am confused -- where does the difference lie between push and
pull?

You said that streams are push-style. Is that not what I just described? If
there are no more threads in the ForkJoinPool, then the only recourse
besides waiting is for the producer thread itself to jump in and start
doing the work too. Other than that, I am not seeing the difference. And if
that is the only difference, I am fine with doing that instead.

I always thought the difference of push vs pull was who was deciding when
work is to be done. Push being that the producer says do it now, only
stopping when there is not a worker available. Whereas pull is the workers
choosing when they become available to do work, stopping or delaying for
any reason they want.

Please do respond to the above, as I think my lack of understanding is
definitely part of the problem here, so I want to understand. That said,
your ClosableBlockingQueue idea has become more interesting to me.

For me, dropping doesn't make sense, at least as a DIRECT option. But, what
might make more sense would be an overloaded constructor/factory that
provides a Consumer<T> overflowConsumer, so the user can decide what they
want to do upon filling up, whereas the original constructor/factory is for
blocking.

For worker termination, can you give me an example of when a worker would
want to terminate? For me, I see workers purely as threads -- probably from
the FJP or WSP. I don't really see a situation where them "terminating"
makes sense, but that's likely lack of creativity from me.

At best, a worker should be able to communicate that it completed the task
at hand, but that's a given by being a thread.

As for FIFO, I think that first depends on what we actually want to make
here.

If this is going to be another branch of the collections framework, like
TransferQueue, then this would make sense to do as a separate
implementation under the same interface. Call that parent interface a
ClosableQueue, where the producer can close the queue once the queue has
been emptied, and simultaneously, all workers have finished their tasks.
From there, the idea of priority could be introduced, to help decide order,
similar to how the PriorityQueue does it.

But if that is too much, and this should be a single implementation (maybe
even without a new parent interface, just existing parent interfaces), then
providing an overloaded constructor/factory that adds order would do the
job. The non-comparator overloads must communicate that they don't do
priority, but tbh, multiple implementations make more sense at that point
imo. Less error-prone for the developer, lest they assume that no
comparator means natural order.

On Fri, Aug 28, 2026 at 5:32 AM Viktor Klang <[email protected]>
wrote:

> >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