Thursday, October 1, 2015

Building a Twitter Application to collect the tweets on Hadoop

This blog talks about building a BigInsights Application to get the real-time Twitter Feeds.

Step 1: Register the application in Twitter and generate the OAuth Keys

a) Login to and create an Application

b) Open the "Keys and Access Tokens" Tab and get the Consumer Key and Consumer Secret

c) In the same page, click the button "Create my access token" to generate the access token

We will be using the Consumer Key / Consumer Secret and Access Token / Access Token Secret in our BigInsights Applications to get the feeds from Twitter.

Step 2: Building a BigInsights Application to get the Feeds from Twitter.

a) Create a BigInsights Project - TwitterSearchApp.

Refer to install Text Analytics Plugin in Eclipse.

b)  I am using twitter4j API to connect to Twitter. Download twitter4j-core-*.jar from

Put the twitter4j-core-*.jar under /TwitterSearchApp/BIApp/workflow/lib/twitter4j-core-4.0.4.jar

c) Create a java file under src folder in package com.test.reader.

Content of

package com.test.reader;

import java.util.ArrayList;
import java.util.HashSet;
import java.util.List;
import java.util.Set;

import org.apache.hadoop.conf.Configuration;
import org.apache.hadoop.fs.FSDataOutputStream;
import org.apache.hadoop.fs.FileSystem;
import org.apache.hadoop.fs.Path;

import twitter4j.JSONObject;
import twitter4j.Query;
import twitter4j.QueryResult;
import twitter4j.Status;
import twitter4j.Twitter;
import twitter4j.TwitterFactory;
import twitter4j.conf.ConfigurationBuilder;

public class TweetReader {

    Twitter twitter;

     * Get the Twitter Instance
    public Twitter getConnection(String consumerKey, String consumerSecret,
            String accessToken, String accessTokenSecret, String httpProxyHost,
            int httpProxyPort) throws Exception {
        if (twitter == null) {
            ConfigurationBuilder cb = new ConfigurationBuilder();
            TwitterFactory tf = new TwitterFactory(;
            twitter = tf.getInstance();
        return twitter;

    // Get the Tweets based on the search Term
    private List<Status> getTweets(String keyword) throws Exception {
        List<Status> tweets = new ArrayList<Status>();
        List<Status> temp = null;
        QueryResult result = null;
        Query query = new Query(keyword);

        int i = 2;
        do {
            result =;
            temp = result.getTweets();

        } while ((query = result.nextQuery()) != null && i < 5);
        return tweets;

    // Convert the Tweet Object to JSON Array
    public String getTweets(String[] keywordArray) throws Exception {
        StringBuffer buff = new StringBuffer();
        int start = 0;
        for (String keyword : keywordArray) {

            List<Status> actTemp = getTweets(keyword);

            // Remove the duplicate tweets
            Set<Status> temp = new HashSet<Status>();
            for (Status tweet : actTemp) {

            for (Status tweet : temp) {
                String text = tweet.getText();

                if (start != 0) {

                JSONObject obj = new JSONObject();
                obj.put("CreatedAt", tweet.getCreatedAt().toString());
                obj.put("User", tweet.getUser().getScreenName());
                obj.put("Text", text);
                obj.put("RetweetCount", tweet.getRetweetCount());
                obj.put("FriendsCount", tweet.getUser().getFriendsCount());
                obj.put("FollowersCount", tweet.getUser().getFollowersCount());
                obj.put("Id", new Long(tweet.getId()));
                        (tweet.getUser().getFollowersCount() / tweet.getUser()
                String rec = obj.toString();



        return buff.toString();

    // Write the file to HDFS
    private void writeOutput(Configuration conf, String hdfsPath, String content)
            throws Exception {
        FileSystem fs = FileSystem.get(conf);

        FSDataOutputStream fos = fs.create(new Path(hdfsPath));
        OutputStream out = fos.getWrappedStream();


    public static void main(String[] args) throws Exception {
        Configuration conf = new Configuration();
        String[] keywordArray = args[0].split("\n");
        String hdfsPath = args[1];
        String consumerKey = args[2];
        String consumerSecret = args[3];
        String accessToken = args[4];
        String accessTokenSecret = args[5];
        String proxyIP = args[6];
        String proxyPort = args[7];

        TweetReader reader = new TweetReader();
        reader.getConnection(consumerKey, consumerSecret, accessToken,
                accessTokenSecret, proxyIP, Integer.parseInt(proxyPort));
        String content = reader.getTweets(keywordArray);
        reader.writeOutput(conf, hdfsPath, content);



d) Edit the /TwitterSearchApp/BIApp/application/application.xml

<application-template xmlns="" xmlns:xsi="">
        <property isInputPath="false" isOutputPath="false" isRequired="true" label="Search" name="keywordArray" paramtype="STRING" uitype="textfield"/>
        <property isInputPath="false" isOutputPath="false" isRequired="true" label="Output Directory" name="hdfsPath" paramtype="STRING" uitype="textfield"/>
        <property isInputPath="false" isOutputPath="false" isRequired="true" label="consumerKey" name="consumerKey" paramtype="STRING" uitype="textfield"/>
        <property isInputPath="false" isOutputPath="false" isRequired="true" label="consumerSecret" name="consumerSecret" paramtype="STRING" uitype="textfield"/>
        <property isInputPath="false" isOutputPath="false" isRequired="true" label="accessToken" name="accessToken" paramtype="STRING" uitype="textfield"/>
        <property isInputPath="false" isOutputPath="false" isRequired="true" label="accessTokenSecret" name="accessTokenSecret" paramtype="STRING" uitype="textfield"/>
        <property isInputPath="false" isOutputPath="false" isRequired="false" label="proxyIP" name="proxyIP" paramtype="STRING" uitype="textfield"/>
        <property isInputPath="false" isOutputPath="false" isRequired="false" label="proxyPort" name="proxyPort" paramtype="STRING" uitype="textfield"/>
        <asset id="TwitterSearch" type="WORKFLOW"/>

e) Update the file /TwitterSearchApp/BIApp/workflow/config-default.xml with below details


f) Update the /TwitterSearchApp/BIApp/workflow/workflow.xml

<workflow-app name="wfapp" xmlns="uri:oozie:workflow:0.2">
    <start to="tweetSearch"/>
    <!-- add actions here -->
    <action name='tweetSearch'>
        <ok to="end"/>
        <error to="kill"/>
    <kill name="kill">
        <message>error message[${wf:errorMessage(wf:lastErrorNode())}]</message>
    <end name="end"/>

e) Publish the Applications to BigInsights Cluster. Refer to publish the Application.

Step 3: Running the Twitter Application from BigInsights Cluster

a) Deploy the published Application in the Cluster

The Search Parameter is used to define the keywords to be searched in Twitter.

Provide Consumer Key / Consumer Secret and Access Token / Access Token Secret.

Provide the ProxyIP & ProxyPort details to connect to Internet from your Cluster.

Click Run

After completing the Job, open the newly generated file  as sheets and File Reader as JSON Array to see the output.

The above example, we used the Twitter Rest API to get the tweets from Twitter. There are some API Rate Limitations per User or per Application - Refer If you want to avoid these limitations you can consider Gnip or DataSift.

No comments: