Jacob,

Just curious if you think extra GBK (through Reshuffle transform) is
prohibitively expensive in your case, or it is more on the lines of
minimizing waste. Is  this a streaming job? There is already an implicit
shuffle for pubsub dedup.

On Thu, Dec 7, 2017 at 2:24 PM, Jacob Marble <[email protected]> wrote:

> re asynchronicity, I agree that simply having 1000 threads blocking on IO
> is also not ideal.
>
> Ben, you reminded me, I think Flume has a PCollection property doNotFuse()
> or something like that. For example:
>
> pipeline
>   .apply(TextIO.read()...)
>   .doNotFuse()
>   .apply(NextTransformation)
>
> An annotation may not be ideal for fusion because a given DoFn may behave
> better with fusion in some circumstances, but not others.
>
> Jacob
>
> On Thu, Dec 7, 2017 at 1:15 PM, Robert Bradshaw <[email protected]>
> wrote:
>
>> Preventing fusion is probably two low-level, and one rarely want to
>> disable all fusion. The Beam SDK has a "reshuffle" operation which ie
>> preferable to a GBK for decoupling the parallelism and/or work
>> distribution of two stages which is what you want here.
>>
>> In this particular case I might wonder if adding asynchronicity into
>> the DoFn itself is not a bad option (if indeed a single worker can do
>> the job of thousands, it's wasteful to have thousands of workers
>> sitting around doing nothing but waiting). Ideally an intelligent
>> runner could possibly pack more worker threads onto a single machine
>> in this case of course, but until then you could still do this
>> manually but factor out the code cleanly by having a DoFn that takes a
>> callable (with the actual logic) as input and does the work. (Note
>> that one must take care to handle windows and bundle closing property;
>> consider using the batching transforms that ship with the Beam SDK to
>> turn a PCollection<T> into a PCollection<List<T>> and then writing a
>> DoFn that processes batches, amortizing communication costs across
>> elements.)
>>
>> - Robert
>>
>>
>>
>>
>> On Thu, Dec 7, 2017 at 12:33 PM, Jacob Marble <[email protected]>
>> wrote:
>> > Is there a long-term plan for preventing fusion in Dataflow pipelines?
>> Maybe
>> > a simple flag --disableFusion ?
>> >
>> > I have read a few of the discussions in the Beam mailing lists, and I
>> > haven't found any sentiment that something should be changed about
>> Dataflow,
>> > only that Dataflow users should work around this with sometimes-costly
>> GBKs.
>> >
>> > My example is a streaming pipeline:
>> > PubSub => memcached read => http request => memcached write => GCS write
>> >
>> > The three middle steps are implemented as few lines of readable code.
>> That's
>> > how pipeline steps/transformations should be written, agreed? But, since
>> > they are fused, they are all slow.
>> >
>> > I have tried some hacks successfully. For example, a DoFn handles an
>> element
>> > by (1) start an async network request then (2) output any network
>> responses
>> > responses waiting in a response queue. One DoFn instance can handle
>> > thousands of network requests concurrently. Also, this feels like I'm
>> using
>> > Beam incorrectly.
>> >
>> > Jacob
>>
>
>

Reply via email to