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()));