Showing posts with label Big Data. Show all posts
Showing posts with label Big Data. Show all posts

Tuesday, June 21, 2016

Count the Wor(l)d with Apache Beam

Count the wor(l)d with Apache Beam

About Beam : http://blog.nanthrax.net/2016/01/introducing-apache-dataflow/

On a  brief note:

Apache Beam is an open source, a unified programming model that you can use to create a data processing PipelineA Pipeline is a chain of processes on the data. Jobs are described in Pipeline. User has to write only this Pipeline part.


The Apache Beam SDK is composed by four parts:
  • Pipelines are the streaming and processing logic that you want to implement. It’s a chain of processes. Basically, in a pipeline, you read data from a source, you apply transformations on the data, and eventually send the data to a destination (named sink in Dataflow wording).
  • PCollection is the object transported inside a pipeline. It’s the container of the data, flowing between each step of the pipeline.
  • Transform is a step of a pipeline. It takes an incoming PCollection and creates an outcoming PCollection. You can implement your own transform function.
  • Sink and Source are used to retrieve data as input (first step) of a pipeline, and eventually send the data outside of the pipeline.
Here is a Word Count sample ( it works with Jdk1.8 only)



package com.test;


import org.apache.beam.sdk.Pipeline;
import org.apache.beam.sdk.io.TextIO;
import org.apache.beam.sdk.options.PipelineOptions;
import org.apache.beam.sdk.options.PipelineOptionsFactory;
import org.apache.beam.sdk.transforms.Count;
import org.apache.beam.sdk.transforms.Filter;
import org.apache.beam.sdk.transforms.FlatMapElements;
import org.apache.beam.sdk.transforms.MapElements;
import org.apache.beam.sdk.values.KV;
import org.apache.beam.sdk.values.TypeDescriptor;

import java.util.Arrays;

/**
 * Created by Siva on 6/16/16.
 */
public class WordCount {
    public static void main(String[] args) {
        // Step1: Create a Pipeline
        PipelineOptions options = PipelineOptionsFactory.create();
        // Step2: set filter options
        Pipeline p = Pipeline.create(options);
        // Step2 & 3: Apply Filters.
        p.apply(TextIO.Read.from("/Users/Siva/beam-files/count.txt"))
                .apply(FlatMapElements.via((String word) -> Arrays.asList(word.split("[^a-zA-Z']+")))
                        .withOutputType(new TypeDescriptor() {
                        }))
                .apply(Filter.byPredicate((String word) -> !word.isEmpty()))
                .apply(Count.perElement())
                .apply(MapElements
                        .via((KV wordCount) -> wordCount.getKey() + ": " + wordCount.getValue())
                        .withOutputType(new TypeDescriptor() {
                        }))
         // Step4: Write output
                .apply(TextIO.Write.to("/Users/Siva/beam-files/count-output.txt"));

        p.run();

        System.out.println(">>> Done");

    }

}


And the pom.xml should look like below.

pom.xml

<?xml version="1.0" encoding="UTF-8"?>
<project xmlns="http://maven.apache.org/POM/4.0.0"
         xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
         xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsd">
    <modelVersion>4.0.0</modelVersion>

    <groupId>groupId</groupId>
    <artifactId>Temp</artifactId>
    <version>1.0-SNAPSHOT</version>

    <dependencies>
        <dependency>
            <groupId>org.apache.beam</groupId>
            <artifactId>beam-sdks-java-core</artifactId>
            <version>0.1.0-incubating</version>
        </dependency>
        <dependency>
            <groupId>org.apache.beam</groupId>
            <artifactId>beam-runners-direct-java</artifactId>
            <version>0.1.0-incubating</version>
            <scope>runtime</scope>
        </dependency>
    </dependencies>

    <build>
        <plugins>
            <plugin>
                <groupId>org.apache.maven.plugins</groupId>
                <artifactId>maven-compiler-plugin</artifactId>
                <configuration>
                    <source>1.8</source>
                    <target>1.8</target>
                </configuration>
            </plugin>
        </plugins>
    </build>
</project>