The latter. Dataflow native PubsubIO while running on GCP. On Tue, Dec 12, 2017 at 9:28 PM Ankur Chauhan <[email protected]> wrote:
> 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 >>>>>> >>>>> >>>>> >>>> >>> >> >
