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]