Apache Beam 3rd Party Java Extensions
These are some of the 3rd party Java libraries that may be useful for specific applications.
Parsing HTTPD/NGINX access logs.
Summary
The Apache HTTPD webserver creates logfiles that contain valuable information about the requests that have been done to the webserver. The format of these log files is a configuration option in the Apache HTTPD server so parsing this into useful data elements is normally very hard to do.
To solve this problem in an easy way a library was created that works in combination with Apache Beam and is capable of doing this for both the Apache HTTPD and NGINX.
The basic idea is that the logformat specification is the schema used to create the line. This parser is simply initialized with this schema and the list of fields you want to extract.
Project page
https://github.com/nielsbasjes/logparser
License
Apache License 2.0
Download
<dependency>
<groupId>nl.basjes.parse.httpdlog</groupId>
<artifactId>httpdlog-parser</artifactId>
<version>5.0</version>
</dependency>
Code example
Assuming a WebEvent class that has a setters setIP, setQueryImg and setQueryStringValues
PCollection<WebEvent> filledWebEvents = input
.apply("Extract Elements from logline",
ParDo.of(new DoFn<String, WebEvent>() {
private Parser<WebEvent> parser;
@Setup
public void setup() throws NoSuchMethodException {
parser = new HttpdLoglineParser<>(WebEvent.class,
"%h %l %u %t \"%r\" %>s %b \"%{Referer}i\" \"%{User-Agent}i\" \"%{Cookie}i\"");
parser.addParseTarget("setIP", "IP:connection.client.host");
parser.addParseTarget("setQueryImg", "STRING:request.firstline.uri.query.img");
parser.addParseTarget("setQueryStringValues", "STRING:request.firstline.uri.query.*");
}
@ProcessElement
public void processElement(ProcessContext c) throws InvalidDissectorException, MissingDissectorsException, DissectionFailure {
c.output(parser.parse(c.element()));
}
})
);
Analyzing the Useragent string
Summary
Parse and analyze the useragent string and extract as many relevant attributes as possible.
Project page
https://github.com/nielsbasjes/yauaa
License
Apache License 2.0
Download
<dependency>
<groupId>nl.basjes.parse.useragent</groupId>
<artifactId>yauaa-beam</artifactId>
<version>4.2</version>
</dependency>
Code example
PCollection<WebEvent> filledWebEvents = input
.apply("Extract Elements from Useragent",
ParDo.of(new UserAgentAnalysisDoFn<WebEvent>() {
@Override
public String getUserAgentString(WebEvent record) {
return record.useragent;
}
@YauaaField("DeviceClass")
public void setDC(WebEvent record, String value) {
record.deviceClass = value;
}
@YauaaField("AgentNameVersion")
public void setANV(WebEvent record, String value) {
record.agentNameVersion = value;
}
}));
Error handling and dead letter queues
Summary
Asgarde simplifies error handling in the transformation steps of Beam pipelines. The Beam
ErrorHandler
aggregates bad records into a single dead letter queue, but each transformation step still has to catch its own errors
(a try/catch block and a BadRecordRouter in each DoFn, or exceptionsInto/exceptionsVia). Asgarde keeps the
fluent style of the apply chain: each step catches its errors as Failure objects (step name, input element and
exception), gathered for the whole flow.
It accepts the Beam MapElements and FlatMapElements, and provides DoFn classes with built-in error handling
(MapElementFn, FlatMapElementFn, FilterFn…) supporting side inputs and the DoFn lifecycle. It can also keep,
in the failures, the element that entered the flow, to replay a failure from the start, and counts the failures per
step with Beam metrics. The Asgarde failures can be converted to BadRecords and added to an ErrorHandler, for a
single dead letter queue together with the Beam IOs. Beam is a provided dependency: the library isn’t tied to a Beam
version. Kotlin extensions are included, and a Python version is available on PyPI
(pasgarde).
Project page
https://github.com/tosun-si/asgarde
Documentation: https://tosun-si.github.io/asgarde/
License
MIT License
Download
<dependency>
<groupId>fr.groupbees</groupId>
<artifactId>asgarde</artifactId>
<version>1.4.0</version>
</dependency>
Code example
WithFailures.Result<PCollection<Integer>, Failure> result = CollectionComposer.of(input)
.apply("Trim", MapElements.into(TypeDescriptors.strings()).via((String value) -> value.trim()))
.apply("Parse", MapElementFn.into(TypeDescriptors.integers()).via((String value) -> Integer.parseInt(value)))
.apply("Keep even numbers", FilterFn.by(number -> number % 2 == 0))
.getResult();
PCollection<Integer> output = result.output();
PCollection<Failure> failures = result.failures(); // The failures of all the steps
// Optional: a single dead letter queue with the Beam IOs using the ErrorHandler.
errorHandler.addErrorCollection(failures.apply("To bad records", FailureTransforms.toBadRecords()));

