StreamConverters.asInputStream not workable when using broadcast with ByteString source?
#1807
Unanswered
mdedetrich
asked this question in
Q&A
Replies: 1 comment 12 replies
|
|
12 replies
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Uh oh!
There was an error while loading. Please reload this page.
Uh oh!
There was an error while loading. Please reload this page.
Currently I have this task where I have a
Source(which is aByteString) and which I want to fan out via aBroadcastto the other flows. One flow goes directly into a Sink (which I will call an upload sink) and the other flow needs to be needs to convert theByteStringSourceto an JavaInputStreamsince its going to an external API that only acceptsInputStream, specifically this is Apache tika which has a detect function as follows(note the fact that this function does NOT close the
InputStream, you have to do this yourself)Pekko Streams provides a
StreamConverters.asInputStream()function which is aSink[ByteString, InputStream]. As discussed at akka/akka-core#23187 (comment) and noted in the docs, you need to be careful when using the materializedInputStreamas its a blocking interface, which means that if you do something like thisIt will create a deadlock as materialized values need to return immediately and the
Tika().detect(inputStream)call will block. Due to this its advised that you materialize theStreamConverters.asInputStream()stream immediately to get theInputStreamand then use theInputStreamas desired.My current predicament is that while this works if you have a single trivial flow, when you have a broadcast style fanout it causes issues one way or another. For example one way to solve this issue would be as follows
The premise here is to use a Java
CompletableFutureto wrap theTika().detect(inputStream)computation so that.mapMaterializedValuereturns immediately. Doing this however still deadlocks, to avoid this one can changeto
The
inputStream.close()appears to force demand on the stream, fixing the deadlock however doing this causes other issues, namely that this exception ends up being thrown.With this error the logic of the stream still works, but it appears that Pekko Streams cannot properly handle eagerly closing theInputStream(even though callinginputStream.close()multiple times is perfectly valid). More importantly, doinginputStream.close()like this is a workaround/hack because as described in the docsPekko Streams is meant to handle the resource cleanup of the underlying stream and so you shouldn't have to call
.close()like this.The thing is I have tried various different implementations of this (i.e. using
preMaterialize, custom graph withGraphDSL+Broadcast,Flow.fromSinkAndSource,fromMaterializer,Source.lazyCompletionStagevsSource.completionStageetc etc) and nothing seems to have solve the underlying issue. Either you have a deadlock (if you don't doinputStream.close()) or if you do haveinputStream.close()then you get theSubscriptionWithCancelException$NoMoreElementsNeededexception (note that I haven't been able to reproduce this exception with a trivial example locally but it occurs in production whencomputeMediaTypeAndUploadis composed as part of a bigger stream).Of course one way to "solve" the problem is to just do
to avoid using
InputStreamentirely but this basically kills the streaming part of the flow into Apache Tika, as you just evaluate the entire stream into a single largeStringwhere as my aim is to make sure that both flows inBroadCastare streaming based.EDIT: I stated earlier that with the
SubscriptionWithCancelException$NoMoreElementsNeededthe logic still worked, this isn't actually the case.All reactions