Thanks Raghu -
Cross-posted per your request.

On Mon, Nov 20, 2017 at 2:10 PM, Raghu Angadi <[email protected]> wrote:

> Thanks for the report Peter. Generally speaking autoscaling in Dataflow
> streaming application is based on two factors : backlog (e.g. PubSub in
> this case) and CPU used for current throughput. Multiple stages makes it a
> bit more complex, but over it is essentially trying to scale CPU resource
> so that extrapolated throughput handles current backlog, backlog growth.
>
> As you noted, cases where small number of messages trigger large amount of
> work in the pipeline makes backlog estimation tricky. That can cause delays
> in upscaling (depending on how the fan out happens in the application), but
> I would not expect the pipeline to downscale too early when it still has a
> lot of work pending. One possibility is that processing involves blocking
> work which might keep CPU utilization lower (say 40-50%).
>
> I am certainly interested in closer look at your specific job. Would you
> mind asking it on Stackoverflow (https://stackoverflow.com/
> questions/tagged/google-cloud-dataflow)? It is also open to public.
> Please provide job_id to look at.
>
> IFAIK, Dataflow is in the only Beam runner that currently supports
> autoscaling based on changes in load. Others might need user to trigger
> recaling (e.g. Flink).
>
> Raghu.
>
> On Fri, Nov 17, 2017 at 8:02 AM, Peter Mueller <[email protected]> wrote:
>
>> Hi,
>>
>> I'm curious if anyone can point me towards greater visibility into how
>> various Beam Runners manage autoscaling.  We seem to be experiencing
>> hiccups during both the 'spin up' and 'spin down' phases, and we're left
>> wondering what to do about it.  Here's the background of our particular
>> flow:
>>
>> 1- Binary files arrive on gs://, and object notification duly notifies a
>> PubSub topic.
>> 2- Each file requires about 1Min of parsing on a standard VM to emit
>> about 30K records to downstream areas of the Beam DAG.
>> 3- 'Downstream' components include things like inserts to BigQuery,
>> storage in GS:, and various sundry other tasks.
>> 4- The files in step 1 arrive intermittently, usually in batches of
>> 200-300 every hour, making this - we think - an ideal use case for
>> autoscaling.
>>
>> What we're seeing, however, has us a little perplexed:
>>
>> 1- It looks like when 'workers=1', Beam bites off a little more than it
>> can chew, eventually causing some out-of-RAM errors, presumably as the
>> first worker tries to process a few of the PubSub messages which, again,
>> take about 60 seconds/message to complete because the 'message' in this
>> case is that a binary file needs to be deserialized in gs.
>> 2- At some point, the runner (in this case, Dataflow), gets the message
>> additional workers are needed and real work can now get done.  During this
>> phase, errors decrease and throughput rises.   Not only are there more
>> deserializers for step1, but the step3/downstream tasks are evenly spread
>> out.
>> 3-Alas, the previous step gets cut short when Dataflow senses (I'm
>> guessing) that enough of the PubSub messages are 'in flight' to begin
>> cooling down a little. That seems to come a little too soon, and workers
>> are getting pulled as they chew through the PubSub messages themselves -
>> even before the messages are 'ACK'd'.
>>
>> We're still thrilled with Beam, but I'm guessing the less-than-optimal
>> spin-up/spin-down phases are resulting in 50% more VM usage than what is
>> needed.  What do the runners look for beside PubSub consumption?  Do they
>> look at RAM/CPU/etc??? Is there anything a developer can do, beside ACK a
>> PubSub message to provide feedback to the runner that more/less resources
>> are required?
>>
>> Incidentally, in case anyone doubted Google's commitment to open-source,
>> I spoke about this very topic with an employee there yesterday, and she
>> expressed interest in hearing about my use case, especially if it ran on a
>> non-Dataflow runner!  We hadn't yet tried our Beam work on Spark (or
>> elsewhere), but would obviously be interested in hearing if one runner has
>> superior abilities to accept feedback from the workers for THROUGHPUT_BASED
>> work.
>>
>> Thanks in advance,
>>
>> Peter Mueller
>> CTO,
>> ATS, Inc.
>>
>>
>>
>

Reply via email to