Are you referring to the PubsubIO implementation that is in the Java sdk or the 
native on that gets used when a dataflow is running using the GCP Cloud 
dataflow service?

Sent from my iPhone

> On Dec 12, 2017, at 19:46, Raghu Angadi <[email protected]> wrote:
> 
> 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