Skip to content

docs: fill stream operator doc gaps and add missing overloads - #3590

Merged
He-Pin merged 5 commits into
apache:mainfrom
pjfanning:docs-stream-operator-gaps
Oct 9, 2026
Merged

He-Pin merged 5 commits into
apache:mainfrom
pjfanning:docs-stream-operator-gaps

Conversation

@pjfanning

@pjfanning pjfanning commented Oct 8, 2026 •

Copy link
Copy Markdown
Member

Motivation

Some stream operator pages have empty sections or stubs, several pages
list only some of an operator's overloads, and the operator index
generator still lists removed operators as pending.

Modification

  • Unzip.md, UnzipWith.md: fill the empty Signature sections
  • Balance.md: add the missing Signature section
  • Source/items.md: remove the empty Examples heading (no snippet exists)
  • stream-flows-and-basics.md, stream-dynamic.md: remove empty
    Introduction headings
  • stream-refs.md: replace the two-line Bulk Stream References stub with a
    "Transferring large amounts of data" subsection under Stream References.
    It explains using the existing SourceRef/SinkRef for large transfers
    (there is no separate bulk API) and lists the caveats: back-pressure
    across the network, single-use refs, no automatic resumption, the
    remoting frame size limit and the subscription timeout
  • StreamOperatorsIndexGenerator.scala: drop the removed throttleEven,
    actorPublisher and actorSubscriber from the pending lists
  • Add missing overloads and describe their extra parameters:
    dropRepeated(p), flatMapConcat(parallelism, f),
    groupedWeightedWithin(maxWeight, maxNumber, d)(costFn),
    takeWhile(p, inclusive), zipLatestWith(..., eagerComplete),
    Compression.gzip/deflate level and autoFlush variants, and the
    onErrorContinue predicate and Java variants

Also fix stream docs whose text contradicts the code (often text copied from a neighbouring operator):

  • Guides: Unzip name and MergeSequence output (stream-graphs), pull() precondition (stream-customize), KillSwitch is a BidiFlow (stream-dynamic), Java failure wording (stream-error), Scala vs Java framing truncation (stream-io), NotUsed package (stream-quickstart), drop the "founding member" claim (stream-introduction, stream-flows-and-basics)
  • Source pages: asSubscriber/fromPublisher signatures and links, iterate parameter names, lazySource, maybe (completionStage), queue completion
  • Sink pages: queue backpressure, collect, completionStageSink, futureSink, lazyCompletionStageSink, lazyFutureSink (sink, not flow)
  • Flow pages: flattenOptional description; remove the wrong second "completes" line from completionStageFlow, futureFlow, lazyFlow, lazyFutureFlow, lazyCompletionStageFlow and add missing "cancels" lines
  • asJavaStream: elements come from upstream
  • Source-or-Flow pages: alsoToAll attaches Sinks, batch has no costFn, buffer fail case, delay/delayWith take DelayOverflowStrategy, foldWhile completion, interleave/interleaveAll backpressure when downstream backpressures, limitWeighted accumulated cost, mergeLatest example and eagerComplete, onErrorContinue fails on unhandled errors, wireTap Sink variant signatures, zipLatestWith combine function, Source.sliding signature

Every signature was checked against the scaladsl/javadsl sources. The new @apidoc anchors are first checked by the CI docs build. This PR merges cleanly with the other open docs PRs (#3583, #3587, #3588 and others; checked with git merge-tree).

Adding reference pages for operators that have none (Sink.foldAsync, Sink.contramap, Flow.fromFunction, Flow.join, the processor methods and the for-comprehension support) is tracked in #3589.

Result

Operator pages list all overloads, have no empty sections, and describe what the code does.

Tests

  • Not run - docs only (plus removing stale names from the docs index
    generator's pending lists)

References

Refs #3589

@pjfanning pjfanning added this to the 2.0.0-M5 milestone Oct 8, 2026
pjfanning added a commit to pjfanning/incubator-pekko that referenced this pull request Oct 8, 2026
Motivation:
The mapAsync page says "Up to `n` elements can be processed
concurrently", but the parameter is called parallelism. apache#3590 could
not fix it because this PR changes the line before it.

Modification:
Say `parallelism` instead of `n`.

Result:
The description matches the mapAsync(parallelism)(f) signature.

Tests:
- Not run - docs only

References:
Refs apache#3590

@He-Pin He-Pin left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I verified all 27 documented overload entries against the real scaladsl and javadsl sources and they all match, including the groupedWeightedWithin argument order and the @since 1.2.0 tag - no mismatches, and the pendingSourceOrFlow deletions are genuinely dead (throttleEven removal is confirmed by the 2.0.x MiMa excludes). Two things before this can go in: (1) Check / Tests (scala3) is red here, so it cannot be merged until that job is re-run and green - I do not want to assume it is just the 6h timeout; (2) the iterate.md page below is only half aligned. Note also that dropRepeated.md:8 still has the malformed scala="#dropRepeated):FlowOps..." anchor, which I believe #3587 already fixes - worth a look so the two do not conflict.


The same `seed` value will be used for every materialization of the `Source` so it is **mandatory** that the state is immutable. For example a `java.util.Iterator`, `Array` or Java standard library collection would not be safe as the fold operation could mutate the value. If you must use a mutable value, combining with @ref:[Source.lazySource](lazySource.md) to make sure a new mutable `zero` value is created for each materialization is one solution.
The same `seed` value will be used for every materialization of the `Source` so it is **mandatory** that the state is immutable. For example a `java.util.Iterator`, `Array` or Java standard library collection would not be safe as the `next` function could mutate the value. If you must use a mutable value, combining with @ref:[Source.lazySource](lazySource.md) to make sure a new mutable `seed` value is created for each materialization is one solution.

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

The signature fix here is right - the real parameters are p and f (Source.scala:559 def iterate[T](seed: T)(p: T => Boolean, f: T => T)) - but this sentence you just edited now names next, a parameter the new signature says does not exist. The page ends up using four names for two parameters: the Description above still says hasNext predicate and next function, and the Reactive Streams semantics below still say `next` function returns and `haxNext` predicate (the latter is also misspelled). Could you align the whole page to p/f while you are in here?

He-Pin pushed a commit that referenced this pull request Oct 9, 2026
* docs: fix typos and grammar

Motivation:
A review of the docs found many typos, grammar mistakes and stray
markup characters.

Modification:
Fix typos, grammar, stray backticks/brackets/asterisks, wrong labels
in link text, a duplicated snippet block and an H1 heading in the
middle of a page, across the top-level, common, additional, general,
project, discovery, durable-state, includes, typed and stream docs.
No facts or link targets are changed.

Result:
The docs read correctly.

Tests:
- Not run - docs only

References:
None - found during a review of the docs

* docs: leave discovery/index.md line 238 to #3583

Motivation:
#3583 fixes the neighbouring line in discovery/index.md; fixing line
238 here would conflict with it.

Modification:
Drop the line 238 typo fix from this PR; #3583 fixes it together with
line 239.

Result:
This PR and #3583 merge cleanly.

Tests:
- Not run - docs only

References:
Refs #3583

* docs: mapAsync processes up to parallelism elements concurrently

Motivation:
The mapAsync page says "Up to `n` elements can be processed
concurrently", but the parameter is called parallelism. #3590 could
not fix it because this PR changes the line before it.

Modification:
Say `parallelism` instead of `n`.

Result:
The description matches the mapAsync(parallelism)(f) signature.

Tests:
- Not run - docs only

References:
Refs #3590
@He-Pin

He-Pin commented Oct 9, 2026

Copy link
Copy Markdown
Member

Pushed b06f22e aligning the whole iterate.md page with the real parameter names p/f (Description, warning, RS semantics), fixing the haxNext typo and completing the unfinished "The next example shows how to craet" sentence. The previous Tests (scala3) run was cancelled at the 6h timeout; this push re-triggers CI - will merge once it is green.

@He-Pin

He-Pin commented Oct 9, 2026

Copy link
Copy Markdown
Member

Merged current main into the branch to clear the conflict with #3588 (it had meanwhile fixed haxNext -> hasNext and the unfinished Examples sentence on main). Resolution: kept main's more precise "counts from 1 to n" example sentence, kept this PR's p/f naming everywhere (matches Source.scala:559), and #3587's dropRepeated.md anchor fix survived the auto-merge. CI is re-running; will merge when green.

pjfanning and others added 5 commits October 9, 2026 13:40
Motivation:
Some stream operator pages have empty sections or stubs, several pages
list only some of an operator's overloads, and the operator index
generator still lists removed operators as pending.

Modification:
- Unzip.md, UnzipWith.md: fill the empty Signature sections
- Source/items.md: remove the empty Examples heading (no snippet exists)
- stream-flows-and-basics.md, stream-dynamic.md: remove empty
  Introduction headings
- stream-refs.md: remove the Bulk Stream References stub; no such API
  exists
- StreamOperatorsIndexGenerator.scala: drop the removed throttleEven,
  actorPublisher and actorSubscriber from the pending lists
- Add missing overloads and describe their extra parameters:
  dropRepeated(p), flatMapConcat(parallelism, f),
  groupedWeightedWithin(maxWeight, maxNumber, d)(costFn),
  takeWhile(p, inclusive), zipLatestWith(..., eagerComplete),
  Compression.gzip/deflate level and autoFlush variants, and the
  onErrorContinue predicate and Java variants

Result:
Operator pages list all overloads and have no empty sections.

Tests:
- Not run - docs only (plus removing stale names from the docs index
  generator's pending lists)

References:
Refs apache#3589
Motivation:
operators/Balance.md has no Signature section, unlike the other
fan-out operator pages.

Modification:
Add a Signature section linking stream.*.Balance, in the same form as
Broadcast.md.

Result:
The Balance page links to its API docs.

Tests:
- Not run - docs only

References:
Refs apache#3589
Motivation:
A number of stream guide and operator pages describe behaviour,
parameters or types that do not match the code, often because text was
copied from a neighbouring operator.

Modification:
- Guides: Unzip name and MergeSequence output (stream-graphs), pull()
  precondition (stream-customize), KillSwitch is a BidiFlow
  (stream-dynamic), Java failure wording (stream-error), Scala vs Java
  framing truncation (stream-io), NotUsed package (stream-quickstart),
  drop the "founding member" claim (stream-introduction,
  stream-flows-and-basics)
- Source pages: asSubscriber/fromPublisher signatures and links,
  iterate parameter names, lazySource, maybe (completionStage), queue
  completion
- Sink pages: queue backpressure, collect, completionStageSink,
  futureSink, lazyCompletionStageSink, lazyFutureSink (sink not flow)
- Flow pages: flattenOptional description; remove the wrong second
  "completes" line from completionStageFlow, futureFlow, lazyFlow,
  lazyFutureFlow and lazyCompletionStageFlow, add missing "cancels"
- asJavaStream: elements come from upstream
- Source-or-Flow pages: alsoToAll attaches Sinks, batch has no costFn,
  buffer fail case, delay(With) DelayOverflowStrategy, foldWhile
  completion, interleave(All) downstream backpressure, limitWeighted
  accumulated cost, mergeLatest example and eagerComplete,
  onErrorContinue fails on unhandled errors, wireTap Sink variant
  signatures, zipLatestWith combine function, Source.sliding signature

Result:
These pages match the code.

Tests:
- Not run - docs only

References:
Refs apache#3589
Motivation:
An earlier commit in this PR removed the two-line "Bulk Stream
References" section on the grounds that no bulk stream ref API
exists. The section never described a separate API: it described
using SourceRef/SinkRef as a side-channel for large amounts of data,
which is a valid use case, but it was an unexplained stub.

Modification:
Add a "Transferring large amounts of data" subsection under Stream
References that explains the use case with the existing
StreamRefs.sourceRef()/sinkRef() API and lists the points that matter
for large transfers: back-pressure across the network, single-use
refs, no persistence or automatic resumption, elements must fit in
the remoting frame size (pekko.remote.artery.advanced.maximum-frame-size),
and the subscription timeout.

Result:
The use case is documented instead of being a stub or removed.

Tests:
- Not run - docs only

References:
Refs apache#3589
@He-Pin
He-Pin force-pushed the docs-stream-operator-gaps branch from ef347ef to 620727c Compare October 9, 2026 05:40
@He-Pin

He-Pin commented Oct 9, 2026

Copy link
Copy Markdown
Member

Rebased onto current main (8a012b9) instead of the merge commit. Only iterate.md conflicted (with #3588's meanwhile-merged fixes); resolution keeps main's "counts from 1 to n" example sentence and this PR's p/f naming everywhere else. Branch tip is now 620727c; CI re-running.

He-Pin pushed a commit that referenced this pull request Oct 9, 2026
Motivation:
The Scaladoc of limit and limitWeighted does not match the code:
- the javadsl Flow, Source, SubFlow and SubSource versions were copied
  from take: they say the stream completes when "the defined number of
  elements has been taken" and completes without producing elements
  when n is zero or negative, but these operators fail the stream with
  StreamLimitReachedException once the limit is exceeded
- the scaladsl limitWeighted "Completes when" line talks about the
  number of emitted elements rather than the accumulated cost
- all of them name the exception StreamLimitException, which does not
  exist; LimitWeighted fails with StreamLimitReachedException

Modification:
Describe the actual behaviour in the Emits/Completes/Errors/Cancels
lines of limit and limitWeighted in scaladsl Flow and javadsl Flow,
Source, SubFlow and SubSource, drop the copied zero-or-negative
sentence, add the missing Errors lines and name
StreamLimitReachedException.

Result:
The Scaladoc matches the LimitWeighted stage (limit is limitWeighted
with a cost of 1 per element).

Tests:
- Not run - Scaladoc only

References:
Refs #3590

@He-Pin He-Pin left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

lgtm

@He-Pin
He-Pin merged commit 9474592 into apache:main Oct 9, 2026
10 checks passed
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants