When you read from pubsub, there is an implicit shuffle of the messages
read from PubSub in order to handle messages that PubSub might send more
than once (presumably due to any delays in acking first message back to
PubSub, though that is not the only reason).

Often in such cases where bottleneck for a pipeline are blocking calls to
an external service, simplicity of extra parallelism through reshuffle far
out weight cost of reshuffle. But YMMV. It does take a bit more CPU, but
does not add much to latency.

On Tue, Dec 12, 2017 at 5:38 PM, Jacob Marble <[email protected]> wrote:

> Sorry, can you tell me about the implicit PubSub dedupe?
>
> Jacob
>
> On Tue, Dec 12, 2017 at 5:38 PM, Jacob Marble <[email protected]> wrote:
>
>> This is a very simple streaming job. GBK would be overkill.
>>
>> Jacob
>>
>> On Fri, Dec 8, 2017 at 10:58 AM, Raghu Angadi <[email protected]> wrote:
>>
>>> 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