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