Thanks Aleksandr and Lukasz! These are really helpful. I'll take a look at the BigQueryIO and see if it's easy to implement something similar in our repo. :)
But at the minimum we can always copy and paste DataStoreV1. :p Cheers, Derek On Mon, Oct 16, 2017 at 11:21 AM, Lukasz Cwik <[email protected]> wrote: > Your idea makes total sense and has been brought up before mirroring the > concept of a dead letter queue [1]. > > Your best bet would be to copy and modify the Datstore code in Apache Beam > and add support for such a policy which outputs failed inserts into > something like a dead letter queue. You can use the ideas from BigqueryIO > since it supports a user configurable retry policy [2]. > > If you get this working, I'm sure other users would like to have this > feature and it would be great if you were able to contribute it back to the > community. > > [1]: https://lists.apache.org/thread.html/aefaf81c202ebcf4a5a5c0114a79df > 7c0d9c5fe7bf872878fe3d8e59@%3Cdev.beam.apache.org%3E > [2]: https://github.com/apache/beam/blob/014614b695bac0b636aae662977dd3 > a3fa3b8a1e/sdks/java/io/google-cloud-platform/src/ > main/java/org/apache/beam/sdk/io/gcp/bigquery/BigQueryIO.java#L1263 > > > > On Mon, Oct 16, 2017 at 10:47 AM, Derek Hao Hu <[email protected]> > wrote: > >> 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/go >> ogle-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. >> > > -- Derek Hao Hu Software Engineer | Snapchat Snap Inc.
