Repository navigation
docs: fill stream operator doc gaps and add missing overloads - #3590
Conversation
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
left a comment
There was a problem hiding this comment.
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. | ||
|
|
There was a problem hiding this comment.
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?
* 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
|
Pushed b06f22e aligning the whole |
|
Merged current main into the branch to clear the conflict with #3588 (it had meanwhile fixed |
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
ef347ef to
620727c
Compare
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
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
Introduction headings
"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
actorPublisher and actorSubscriber from the pending lists
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):
Unzipname andMergeSequenceoutput (stream-graphs),pull()precondition (stream-customize), KillSwitch is aBidiFlow(stream-dynamic), Java failure wording (stream-error), Scala vs Java framing truncation (stream-io),NotUsedpackage (stream-quickstart), drop the "founding member" claim (stream-introduction, stream-flows-and-basics)asSubscriber/fromPublishersignatures and links,iterateparameter names,lazySource,maybe(completionStage),queuecompletionqueuebackpressure,collect,completionStageSink,futureSink,lazyCompletionStageSink,lazyFutureSink(sink, not flow)flattenOptionaldescription; remove the wrong second "completes" line fromcompletionStageFlow,futureFlow,lazyFlow,lazyFutureFlow,lazyCompletionStageFlowand add missing "cancels" linesasJavaStream: elements come from upstreamalsoToAllattaches Sinks,batchhas nocostFn,bufferfail case,delay/delayWithtakeDelayOverflowStrategy,foldWhilecompletion,interleave/interleaveAllbackpressure when downstream backpressures,limitWeightedaccumulated cost,mergeLatestexample andeagerComplete,onErrorContinuefails on unhandled errors,wireTapSink variant signatures,zipLatestWithcombine function,Source.slidingsignatureEvery signature was checked against the scaladsl/javadsl sources. The new
@apidocanchors are first checked by the CI docs build. This PR merges cleanly with the other open docs PRs (#3583, #3587, #3588 and others; checked withgit 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
generator's pending lists)
References
Refs #3589