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 >
