tosun-si commented on PR #40268:
URL: https://github.com/apache/beam/pull/40268#issuecomment-5871328521

   Thanks for the review, good point! I updated the summary to mention 
`ErrorHandler`.
   
   To illustrate how they relate, here is the same two-step flow with both:
   
   ```java
   // Beam ErrorHandler: each step is a DoFn catching its own errors
   class ParseFn extends DoFn<String, Order> {
     @ProcessElement
     public void process(@Element String json, MultiOutputReceiver out) throws 
Exception {
       try {
         out.get(ORDERS).output(parse(json));
       } catch (Exception e) {
         BadRecordRouter.RECORDING_ROUTER.route(out, json, 
StringUtf8Coder.of(), e, "Parse");
       }
     }
   }
   // ... and the same for ValidateFn
   
   BadRecordErrorHandler<?> handler = 
pipeline.registerBadRecordErrorHandler(deadLetterSink);
   PCollectionTuple parsed = input.apply("Parse",
       ParDo.of(new ParseFn()).withOutputTags(ORDERS, 
TupleTagList.of(BAD_RECORD_TAG)));
   handler.addErrorCollection(parsed.get(BAD_RECORD_TAG));
   PCollectionTuple validated = parsed.get(ORDERS).apply("Validate",
       ParDo.of(new ValidateFn()).withOutputTags(VALID_ORDERS, 
TupleTagList.of(BAD_RECORD_TAG)));
   handler.addErrorCollection(validated.get(BAD_RECORD_TAG));
   handler.close();
   ```
   
   ```java
   // Asgarde: the same steps, the errors caught for you
   WithFailures.Result<PCollection<Order>, Failure> result = 
CollectionComposer.of(input)
       .apply("Parse", 
MapElementFn.into(TypeDescriptor.of(Order.class)).via(OrderParser::parse))
       .apply("Validate", FilterFn.by(Order::isValid))
       .getResult();
   
   // ... and it plugs into the same ErrorHandler, e.g. with the bad records of 
KafkaIO
   
handler.addErrorCollection(result.failures().apply(FailureTransforms.toBadRecords()));
   ```
   
   So `ErrorHandler` handles the aggregation (and the IO integration), and 
Asgarde removes the per-step error handling code on top of it. Asgarde can also 
keep, in the failures, the element that entered the flow (e.g. the raw 
message), not only the input of the failing step: very useful in production to 
debug and replay a failure from the start. They're complementary, which is what 
the updated summary now says.
   
   The full comparison, and an example with both together (KafkaIO → Asgarde 
steps → BigQueryIO, a single dead letter queue): 
https://tosun-si.github.io/asgarde/getting-started/why-asgarde/#and-the-beam-errorhandler
   
   Let me know if you'd like a different wording!
   


-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]

Reply via email to