Hello,
Take a look flushBatch
https://github.com/apache/beam/blob/master/sdks/java/io/
google-cloud-platform/src/main/java/org/apache/beam/sdk/
io/gcp/datastore/DatastoreV1.java
You can create your own implementation based on that code (do not throw
exception there)


16. okt 2017 8:48 PM kirjutas kuupäeval "Derek Hao Hu" <
[email protected]>:

> I see. Thanks Lukasz.
>
> In that case, do you think there is an easy / clean way to implement the
> behavior I explained: "fail after a certain number of retries and then
> write the failed data to an external datasource"? I'm not sure using the
> [google-cloud-java](https://googlecloudplatform.github.io/
> google-cloud-java/0.26.0/index.html) in a DoFn is the right way to go.
>
> Let me also briefly describe why I feel this is a necessary thing to have:
> I'm using a streaming pipeline taking Pub/Sub input, therefore we need to
> acknowledge each bundle of elements. Suppose we are writing some
> non-critical data to Datastore but if writing to Datastore fails, the whole
> pipeline would just stall. Ideally, it would be better for us to say - hey,
> we can try to write to Datastore up to five times, if it still fails, maybe
> we'll just write it to logs and then redo this write later.
>
> Do you think that makes sense?
>
> Thanks,
>
> Derek
>
> On Mon, Oct 16, 2017 at 10:31 AM, Lukasz Cwik <[email protected]> wrote:
>
>> That source is not available to you as it is part of the Dataflow service.
>>
>> On Mon, Oct 16, 2017 at 10:25 AM, Derek Hao Hu <[email protected]>
>> wrote:
>>
>>> Thanks Lukasz!
>>>
>>> "For an unbounded (streaming) pipeline, Dataflow will retry the bundle
>>> forever until the pipeline is cancelled by the user."
>>>
>>> Can you help point out where this behavior is implemented? I'd like to
>>> take a look and see if it is possible to modify it a bit (e.g. write this
>>> to an external source if it fails after a certain number of retries.)
>>>
>>> Derek
>>>
>>> On Mon, Oct 16, 2017 at 9:40 AM, Lukasz Cwik <[email protected]> wrote:
>>>
>>>> It depends on the runner but that exception that is thrown is per
>>>> bundle processed and it is up to the runner to choose what to do with
>>>> bundles that fail.
>>>>
>>>> For a bounded (batch) pipeline, Dataflow will fail the pipeline after a
>>>> fixed number of retries of each bundle.
>>>> For an unbounded (streaming) pipeline, Dataflow will retry the bundle
>>>> forever until the pipeline is cancelled by the user.
>>>>
>>>> On Sun, Oct 15, 2017 at 11:33 PM, Derek Hao Hu <[email protected]>
>>>> wrote:
>>>>
>>>>> Hi,
>>>>>
>>>>> ​I'm using DatastoreV1 API to write data into Datastore. I've briefly
>>>>> gone through the implementation and it seems the Write transform will 
>>>>> throw
>>>>> a DatastoreException (https://github.com/apache/bea
>>>>> m/blob/1bd17d1b95a6b27331626fa9bdbaa723969b710d/sdks/java/io
>>>>> /google-cloud-platform/src/main/java/org/apache/beam/sdk/io/
>>>>> gcp/datastore/DatastoreV1.java#L1346) when it fails to commit.​
>>>>>
>>>>> In a streaming pipeline, I think it might be possible that some
>>>>> commits might occasionally fail even after five retries, what is the
>>>>> expected behavior here? Is there a way to catch these failed mutations and
>>>>> then save them somewhere? I'm not sure what is the recommended approach
>>>>> since `Write` itself is already a transform which means there seems to be
>>>>> no easy way to catch which mutations / commits actually failed.
>>>>>
>>>>> Could someone help explain what the best approach here is? Right now
>>>>> I'm thinking of writing my own DoFn which just writes each entity to
>>>>> Datastore without even batching - by doing this it seems I'll be easy to
>>>>> catch or write the failed commits to logs. It doesn't seem to be the right
>>>>> approach though - considering there is already a significant amount of
>>>>> effort for Beam to provide it's own DatastoreIO.
>>>>>
>>>>> Thanks,
>>>>> --
>>>>> Derek Hao Hu
>>>>>
>>>>> Software Engineer | Snapchat
>>>>> Snap Inc.
>>>>>
>>>>
>>>>
>>>
>>>
>>> --
>>> Derek Hao Hu
>>>
>>> Software Engineer | Snapchat
>>> Snap Inc.
>>>
>>
>>
>
>
> --
> Derek Hao Hu
>
> Software Engineer | Snapchat
> Snap Inc.
>

Reply via email to