Today i encountered a problem in which i need to know what is the batch size for a TransactionalTridentKafkaSpout.
The default output field by KafkaSpout is "str", which can be of variable length. By some simple thought, it is reasonable that kafka spout will set the batch size dynamically based on the actual sizes of the "str" messages it retrieves in a batch. In other words, i suspect that the batch size may be determined by the following pseudo codes:
batchSize = 0;
totalBytesConsumed=0;
while(totalBytesConsumed < maxBytesConsumed)
{
totalBytesConsumed+= getMessageSize("str");
batchSize ++;
}
After quickly read through trident-kafka's source codes on TransactionalTridentKafkaSpout and trace down to the following classes and their methods:
TransactionalTridentKafkaSpout.getEmitter(...)
TridentKafkaEmitter.fetchMessages(...)
KafkaUtils.fetchMessages(...)
by the builder.addFetch(...) line in "KafkaUtils.fetchMessages(...)", it looks like the batch size is determine by a variable defined in KafkaConfig.fetchSizeBytes, which is defaulted to 1024 * 1024.
I also noticed another variable KafkaConfig.bufferSizeBytes, which is used by the SimpleConsumer class, this is also defaulted to 1024 * 1024.
Therefore i suspect that the batch size of the kafka spout depends on both Math.Min(KafkaConfig.bufferSizeBytes, kafkaConfig.fetchSizeBytes).
After some googling, I noticed for handling huge data in Kafka, the following settings have been used:
kafkaConfig.bufferSizeBytes = 1024 * 1024 * 4;
kafkaConfig.fetchSizeBytes = 1024 * 1024 * 4;
Showing posts with label Storm. Show all posts
Showing posts with label Storm. Show all posts
Wednesday, December 17, 2014
Friday, December 12, 2014
Storm: java.lang.NoClassDefFoundError: org/apache/curator/RetryPolicy
Recently I was working on a Trident topology with Trident-ML which uses nathan marz's storm-kafka which pushes the results from the Trident topology to be read by another storm topology. While the program works perfectly in local cluster testing. When it was deployed in a storm cluster, the following error showed up from nowhere:
java.lang.NoClassDefFoundError: org/apache/curator/RetryPolicy
The error got me stuck for more than half an hour, before i figured out a workable solution. It seems that i am using there is a version compatibility issues with the version of storm, trident-ml, and the storm-kafka. Originally i had been using the following dependencies in the pom file:
<dependency>
<groupId>org.apache.storm</groupId>
<artifactId>storm-core</artifactId>
<version>0.9.2-incubating</version>
<scope>provided</scope>
</dependency>
<dependency>
<groupId>net.wurstmeister.storm</groupId>
<artifactId>storm-kafka-plus-0.8</artifactId>
<version>0.4.0</version>
</dependency>
<dependency>
<groupId>com.github.pmerienne</groupId>
<artifactId>trident-ml</artifactId>
<version>0.0.4</version>
</dependency>
The error was gone, after I changed the storm-kafka dependency from storm-kafka-plus-0.8 to the following:
<dependency>
<groupId>org.apache.storm</groupId>
<artifactId>storm-kafka</artifactId>
<version>0.9.2-incubating</version>
</dependency>
I figured the updated dependency solves the curator framework dependency in the storm-kafka's pom file. Note that if you have a pom-assembly.xml, remember to include the following:
<include>org.apache.storm:*</include>
java.lang.NoClassDefFoundError: org/apache/curator/RetryPolicy
The error got me stuck for more than half an hour, before i figured out a workable solution. It seems that i am using there is a version compatibility issues with the version of storm, trident-ml, and the storm-kafka. Originally i had been using the following dependencies in the pom file:
<dependency>
<groupId>org.apache.storm</groupId>
<artifactId>storm-core</artifactId>
<version>0.9.2-incubating</version>
<scope>provided</scope>
</dependency>
<dependency>
<groupId>net.wurstmeister.storm</groupId>
<artifactId>storm-kafka-plus-0.8</artifactId>
<version>0.4.0</version>
</dependency>
<dependency>
<groupId>com.github.pmerienne</groupId>
<artifactId>trident-ml</artifactId>
<version>0.0.4</version>
</dependency>
The error was gone, after I changed the storm-kafka dependency from storm-kafka-plus-0.8 to the following:
<dependency>
<groupId>org.apache.storm</groupId>
<artifactId>storm-kafka</artifactId>
<version>0.9.2-incubating</version>
</dependency>
I figured the updated dependency solves the curator framework dependency in the storm-kafka's pom file. Note that if you have a pom-assembly.xml, remember to include the following:
<include>org.apache.storm:*</include>
Thursday, December 11, 2014
Storm: Debugging topology running in storm cluster
Usually debugging is much easier when running storm topology in a local cluster, especially when it comes to command line output (e.g. "System.out.println()"). However, once the storm topology is deployed in a storm cluster. the command line output will longer be available. In such case, the logs created by storm becomes very useful as your command line output will be stored there. To this end, suppose you configure to have the storm logs stored in /var/storm-logs directory. and you want to view and continuously print out of command line just like what "System.out.println()" did in a local cluster. And suppose the logs is stored in a file named worker-6666.log. The easiest way to do this is to run the following commands:
> cd /var/storm-log
> tail -f worker-6666.log
The "tail -f" command will print out the last ten lines of the worker-6666.log continuously,
> cd /var/storm-log
> tail -f worker-6666.log
The "tail -f" command will print out the last ten lines of the worker-6666.log continuously,
Wednesday, November 26, 2014
Trident-ML: Indexing Sentiment Classification results with ElasticSearch
In Storm, we may have scenarios in which we like to index results obtained from real-time processing or machine learning into a search and analytics engine. For example, we may have some text streaming in from Kafka messaging system which will go through a TwitterSentimentClassifier (which is available in Trident-ML). After that, we may wish to save the text together with the classified sentiment label as an indexed document in ElasticSearch. This post shows one way to realize such an implementation.
First create a Maven project (e.g. with groupId="com.memeanalytics" and artifactId="es-create-index"), the complete source code of the project can be downloaded from the link below:
https://dl.dropboxusercontent.com/u/113201788/storm/es-create-index.tar.gz
The above code requires the following dependency in pom.xml:
<dependency>
<groupId>org.elasticsearch</groupId>
<artifactId>elasticsearch</artifactId>
<version>1.4.0</version>
</dependency>
However, as the pom and coding of this library has dependency on lucene-core [version=3.6.0] that is an older version that is not compatible with lucene-analyzers [version=3.6.2] which is currently one of Trident-ML's dependency (The TwitterTokenizer in TwitterSentimentClassifier uses this library). As a result, the elasticsearch library above cannot be used if the TwitterSentimentClassifier in Trident-ML is to be used in this project.
Since the above java code and elastic library cannot be used in this project, the project uses httpclient [version=4.3] from org.apache.httpcomponents in its place to communicate with elasticsearch via RESTful api. The httpclient provides CloseableHttpClient and operators such as HttpGet, HttpPut, HttpDelete,
The dependencies section of the pom for this project looks like the following:
> mvn compile exec:java
First create a Maven project (e.g. with groupId="com.memeanalytics" and artifactId="es-create-index"), the complete source code of the project can be downloaded from the link below:
https://dl.dropboxusercontent.com/u/113201788/storm/es-create-index.tar.gz
Configure pom.xml and libraries to be used
Before we proceed, I would like to discuss how to write a elasticsearch client which is compatible with Trident-ML as we will be using both in this project. Traditionally an elasticsearch java client can be implemented using native code such as this:import static org.elasticsearch.node.NodeBuiler.*;
Node node=nodeBuilder().clusterName("elasticsearch").node();
Client client=node.getClient();
//TODO HERE: put document, delete document, etc using the client
node.close();
The above code requires the following dependency in pom.xml:
<dependency>
<groupId>org.elasticsearch</groupId>
<artifactId>elasticsearch</artifactId>
<version>1.4.0</version>
</dependency>
However, as the pom and coding of this library has dependency on lucene-core [version=3.6.0] that is an older version that is not compatible with lucene-analyzers [version=3.6.2] which is currently one of Trident-ML's dependency (The TwitterTokenizer in TwitterSentimentClassifier uses this library). As a result, the elasticsearch library above cannot be used if the TwitterSentimentClassifier in Trident-ML is to be used in this project.
Since the above java code and elastic library cannot be used in this project, the project uses httpclient [version=4.3] from org.apache.httpcomponents in its place to communicate with elasticsearch via RESTful api. The httpclient provides CloseableHttpClient and operators such as HttpGet, HttpPut, HttpDelete,
The dependencies section of the pom for this project looks like the following:
<dependency> <groupId>storm</groupId> <artifactId>storm</artifactId> <version>0.9.0.1</version> </dependency> <dependency> <groupId>com.github.pmerienne</groupId> <artifactId>trident-ml</artifactId> <version>0.0.4</version> <scope>provided</scope> </dependency> <dependency> <groupId>org.apache.httpcomponents</groupId> <artifactId>httpclient</artifactId> <version>4.3</version> </dependency>
Spout
Once the pom.xml is properly updated, we can move to implement the code for the Storm spout used in this project. The spout, named TweetCommentSpout, reads tweets from "src/test/resources/twitter-sentiment.csv" and emits them in batch to the Trident topology. the implementation of the spout is shown below:package com.memeanalytics.es_create_index;
import java.io.BufferedReader;
import java.io.FileInputStream;
import java.io.FileNotFoundException;
import java.io.IOException;
import java.io.InputStreamReader;
import java.util.ArrayList;
import java.util.List;
import java.util.Map;
import backtype.storm.task.TopologyContext;
import backtype.storm.tuple.Fields;
import backtype.storm.tuple.Values;
import com.github.pmerienne.trident.ml.core.TextInstance;
import com.github.pmerienne.trident.ml.preprocessing.EnglishTokenizer;
import com.github.pmerienne.trident.ml.preprocessing.TextTokenizer;
import storm.trident.operation.TridentCollector;
import storm.trident.spout.IBatchSpout;
public class TweetCommentSpout implements IBatchSpout {
private static final long serialVersionUID = 1L;
private static List<List<Object>> data=new ArrayList<List<Object>>();
private int batchIndex;
private int batchSize=10;
static{
BufferedReader br=null;
FileInputStream is=null;
String filePath="src/test/resources/twitter-sentiment.csv";
try {
is=new FileInputStream(filePath);
br=new BufferedReader(new InputStreamReader(is));
String line=null;
while((line=br.readLine())!=null)
{
String[] values = line.split(",");
Integer label=Integer.parseInt(values[0]);
String text=values[1];
// TextTokenizer tokenizer=new EnglishTokenizer();
// List<String> tokens = tokenizer.tokenize(text);
// TextInstance<Integer> instance=new TextInstance<Integer>(label, tokens);
data.add(new Values(text, label));
}
} catch (FileNotFoundException e) {
e.printStackTrace();
}catch(IOException ex)
{
ex.printStackTrace();
}
}
public void open(Map conf, TopologyContext context) {
// TODO Auto-generated method stub
}
public void emitBatch(long batchId, TridentCollector collector) {
// TODO Auto-generated method stub
int maxBatchCount = data.size() / batchSize;
if(maxBatchCount > 0 && batchIndex < maxBatchCount)
{
for(int i=(batchSize * batchIndex); i < data.size() && i < (batchIndex+1) * batchSize; ++i)
{
collector.emit(data.get(i));
}
batchIndex++;
}
}
public void ack(long batchId) {
// TODO Auto-generated method stub
}
public void close() {
// TODO Auto-generated method stub
}
public Map getComponentConfiguration() {
// TODO Auto-generated method stub
return null;
}
public Fields getOutputFields() {
// TODO Auto-generated method stub
return new Fields("text", "label");
}
}
The tuples emitted by the spout contains two fields: "text" and "label", the label is ignored, we are going to have the Trident-ML's TweetSentimentClassifier predict the sentiment label for us instead.Trident operation for ElasticSearch
Next we are going to implement a BaseFilter, named CreateESIndex, which is a Trident operation that create an indexed document in ElasticSearch from each tweet text and its predicted sentiment label. The implementation of the Trident operation is shown below:package com.memeanalytics.es_create_index;
import java.io.IOException;
import java.text.SimpleDateFormat;
import java.util.Calendar;
import java.util.Date;
import java.util.Map;
import org.apache.http.HttpEntity;
import org.apache.http.client.ClientProtocolException;
import org.apache.http.client.HttpRequestRetryHandler;
import org.apache.http.client.methods.CloseableHttpResponse;
import org.apache.http.client.methods.HttpDelete;
import org.apache.http.client.methods.HttpGet;
import org.apache.http.client.methods.HttpPost;
import org.apache.http.client.methods.HttpPut;
import org.apache.http.entity.StringEntity;
import org.apache.http.impl.client.CloseableHttpClient;
import org.apache.http.impl.client.HttpClients;
import org.apache.http.message.BasicHeader;
import org.apache.http.protocol.HTTP;
import org.apache.http.protocol.HttpContext;
import org.apache.http.util.EntityUtils;
import storm.trident.operation.BaseFilter;
import storm.trident.operation.TridentOperationContext;
import storm.trident.tuple.TridentTuple;
public class CreateESIndex extends BaseFilter{
private static final long serialVersionUID = 1L;
private int esIndex=1;
private String wsUrl="http://127.0.0.1:9200";
private String indexName="twittersentiment"; //must be lowercase
private String typeName="trident";
private CloseableHttpClient client;
private String lastIndexedDocumentIdQueryJson="{\"query\": { \"match_all\": {}}, \"size\": 1,"+
"\"sort\": ["+
"{"+
"\"_timestamp\": {"+
"\"order\": \"desc\""+
"}"+
"}"+
"]"+
"}";
public boolean isKeep(TridentTuple tuple) {
// TODO Auto-generated method stub
Boolean prediction =tuple.getBooleanByField("prediction");
String comment=tuple.getStringByField("text");
System.out.println(comment + " >> " + prediction);
if(client != null)
{
HttpPut method=new HttpPut(wsUrl+"/"+indexName+"/"+typeName+"/"+esIndex);
Date currentTime= new Date();
SimpleDateFormat format1 = new SimpleDateFormat("yyyy-MM-dd");
SimpleDateFormat format2 = new SimpleDateFormat("HH:mm:ss");
String dateString = format1.format(currentTime)+"T"+format2.format(currentTime);
CloseableHttpResponse response=null;
try{
String json = "{\"text\":\""+comment+"\", \"prediction\":\""+prediction+"\", \"postTime\":\""+dateString+"\"}";
System.out.println(json);
StringEntity params=new StringEntity(json);
params.setContentType(new BasicHeader(HTTP.CONTENT_TYPE, "application/json"));
method.setEntity(params);
method.addHeader("Accept", "application/json");
method.addHeader("Content-type", "application/json");
response = client.execute(method);
HttpEntity entity=response.getEntity();
String responseText=EntityUtils.toString(entity);
System.out.println(responseText);
}catch(IOException ex) {
ex.printStackTrace();
}finally {
method.releaseConnection();
}
esIndex++;
}
return true;
}
@Override
public void prepare(Map conf, TridentOperationContext context) {
client=HttpClients.custom().setRetryHandler(new MyRetryHandler()).build();
CloseableHttpResponse response=null;
HttpDelete method=new HttpDelete(wsUrl+"/"+indexName);
try{
response = client.execute(method);
HttpEntity entity=response.getEntity();
String responseBody=EntityUtils.toString(entity);
System.out.println(responseBody);
}catch(IOException ex)
{
ex.printStackTrace();
}
}
private class MyRetryHandler implements HttpRequestRetryHandler {
public boolean retryRequest(IOException arg0, int arg1, HttpContext arg2) {
// TODO Auto-generated method stub
return false;
}
}
@Override
public void cleanup() {
try {
client.close();
} catch (IOException e) {
// TODO Auto-generated catch block
e.printStackTrace();
}
}
}
In the prepare() method of the CreateESIndex, a RESTful DELETE call is performed to delete all indexed documents under twittersentiment/trident in ElasticSearch. This is to ensure that no data will be under twittersentiment/trident when the bolt is run. Now in its isKeep() method, the tweet text and its associated predicted sentiment label is serialized to a json and sent to elasticsearch via a http PUT call. The CloseableHttpClient object is closed in its cleanup() method.Trident topology
Now we have the neccessary spout and trident operation, we can define a simple Trident topology which stream tweets-> classified by TwitterSentimentClassifier -> indexed by ElasticSearch. Below is the implementation in the main class:package com.memeanalytics.es_create_index;
import com.github.pmerienne.trident.ml.nlp.TwitterSentimentClassifier;
import storm.trident.TridentTopology;
import backtype.storm.Config;
import backtype.storm.LocalCluster;
import backtype.storm.generated.StormTopology;
import backtype.storm.tuple.Fields;
public class App
{
public static void main( String[] args )
{
LocalCluster cluster=new LocalCluster();
Config config=new Config();
cluster.submitTopology("TridentWriteToESDemo", config, buildTopology());
try{
Thread.sleep(10000);
}catch(InterruptedException ex)
{
}
cluster.killTopology("TridentWriteToESDemo");
cluster.shutdown();
}
private static StormTopology buildTopology()
{
TridentTopology topology=new TridentTopology();
TweetCommentSpout spout=new TweetCommentSpout();
topology.newStream("classifyAndIndex", spout).each(new Fields("text"), new TwitterSentimentClassifier(), new Fields("prediction")).each(new Fields("text", "prediction"), new CreateESIndex());
return topology.build();
}
}
Once it is completed, run the following command in the project root folder:> mvn compile exec:java
Tuesday, November 25, 2014
Trident-ML: Regression using Passive-Aggressive algorithm
This post shows some very basic example of how to use the Passive-Aggressive algorithm as regression algorithm in Trident-ML to process data from Storm Spout.
Firstly create a Maven project (e.g. with groupId="com.memeanalytics" artifactId="trident-regression-pa"). The complete source codes of the project can be downloaded from the link:
https://dl.dropboxusercontent.com/u/113201788/storm/trident-regression-pa.tar.gz
For the start we need to configure the pom.xml file in the project.
Configure pom.xml:
Firstly we need to add the clojars repository to the repositories section:
Next we need to add the storm dependency to the dependencies section (for storm):
Next we need to add the strident-ml dependency to the dependencies section (for PA regression):
Next we need to add the exec-maven-plugin to the build/plugins section (for execute the Maven project):
Next we need to add the maven-assembly-plugin to the build/plugins section (for packacging the Maven project to jar for submitting to Storm cluster):
As can be seen above, the BirthDataSpout is derived from IBatchSpout, and emits a batch of 10 tuples at one time, each tuple is a training record containing the fields ("instance"). The "instance" field contains a datatype Instance<Double> which contains a double array as features and a double value as label. The training records are obtained from a births.csv file residing in src/test/resources.
As can be seen above, the Trident topology has the BirthDataSpout emits Instance<Double> training data can be consumed by RegressionUpdater. The RegressionUpdater object from Trident-ML updates the underlying regressionModel via PA algorithm.
The DRPCStream allows user to pass in a new testing instance to the regressionModel which will then return a "predict" field, that contains the predicted output of the testing instance. The DRPCArgsToInstance is a BaseFunction operation which converts the arguments passed into the LocalDRPC.execute() into an Instance<Double> which can be passed into the RegressionQuery which then uses PARegressor and regressionModel to determine the predicted output value.
Once the coding is completed, we can run the project by navigating to the project root folder and run the following commands:
> .mvn compile exec:java
Firstly create a Maven project (e.g. with groupId="com.memeanalytics" artifactId="trident-regression-pa"). The complete source codes of the project can be downloaded from the link:
https://dl.dropboxusercontent.com/u/113201788/storm/trident-regression-pa.tar.gz
For the start we need to configure the pom.xml file in the project.
Configure pom.xml:
Firstly we need to add the clojars repository to the repositories section:
<repositories> <repository> <id>clojars</id> <url>http://clojars.org/repo</url> </repository> </repositories>
Next we need to add the storm dependency to the dependencies section (for storm):
<dependency> <groupId>storm</groupId> <artifactId>storm</artifactId> <version>0.9.0.1</version> <scope>provided</scope> </dependency>
Next we need to add the strident-ml dependency to the dependencies section (for PA regression):
<dependency> <groupId>com.github.pmerienne</groupId> <artifactId>trident-ml</artifactId> <version>0.0.4</version> </dependency>
Next we need to add the exec-maven-plugin to the build/plugins section (for execute the Maven project):
<plugin> <groupId>org.codehaus.mojo</groupId> <artifactId>exec-maven-plugin</artifactId> <version>1.2.1</version> <executions> <execution> <goals> <goal>exec</goal> </goals> </execution> </executions> <configuration> <includeProjectDependencies>true</includeProjectDependencies> <includePluginDependencies>false</includePluginDependencies> <executable>java</executable> <classpathScope>compile</classpathScope> <mainClass>com.memeanalytics.trident_regression_pa.App</mainClass> </configuration> </plugin>
Next we need to add the maven-assembly-plugin to the build/plugins section (for packacging the Maven project to jar for submitting to Storm cluster):
<plugin> <artifactId>maven-assembly-plugin</artifactId> <version>2.2.1</version> <configuration> <descriptorRefs> <descriptorRef>jar-with-dependencies</descriptorRef> </descriptorRefs> <archive> <manifest> <mainClass></mainClass> </manifest> </archive> </configuration> <executions> <execution> <id>make-assembly</id> <phase>package</phase> <goals> <goal>single</goal> </goals> </execution> </executions> </plugin>
Implement Spout for training data
Once the pom.xml update is completed, we can move to implement the BirthDataSpout which is the Storm spout that emits batches of training data to the Trident topology:package com.memeanalytics.trident_regression_pa;
import java.io.BufferedReader;
import java.io.FileInputStream;
import java.io.FileNotFoundException;
import java.io.IOException;
import java.io.InputStreamReader;
import java.util.ArrayList;
import java.util.Collections;
import java.util.List;
import java.util.Map;
import com.github.pmerienne.trident.ml.core.Instance;
import com.github.pmerienne.trident.ml.testing.data.Datasets;
import backtype.storm.task.TopologyContext;
import backtype.storm.tuple.Fields;
import backtype.storm.tuple.Values;
import storm.trident.operation.TridentCollector;
import storm.trident.spout.IBatchSpout;
public class BirthDataSpout implements IBatchSpout {
private static final long serialVersionUID = 1L;
private int batchSize=10;
private int batchIndex=0;
private static List<Instance<Double>> sample_data=new ArrayList<Instance<Double>>();
private static List<Instance<Double>> testing_data=new ArrayList<Instance<Double>>();
public static List<String> getDRPCArgsList()
{
List<String> drpc_args_list =new ArrayList<String>();
for(Instance<Double> instance : testing_data)
{
double[] features = instance.getFeatures();
String drpc_args="";
for(int i=0; i < features.length; ++i)
{
if(i==0)
{
drpc_args+=features[i];
}
else
{
drpc_args+=(","+features[i]);
}
}
drpc_args+=(","+instance.label);
drpc_args_list.add(drpc_args);
}
return drpc_args_list;
}
static{
FileInputStream is=null;
BufferedReader br=null;
try{
String filePath="src/test/resources/births.csv";
is=new FileInputStream(filePath);
br=new BufferedReader(new InputStreamReader(is));
List<Instance<Double>> temp=new ArrayList<Instance<Double>>();
String line=null;
while((line=br.readLine())!=null)
{
String[] values = line.split(";");
double label= Double.parseDouble(values[values.length-1]);
double[] features=new double[values.length-1];
for(int i=0; i < values.length-1; ++i)
{
features[i]=Double.parseDouble(values[i]);
}
Instance<Double> instance=new Instance<Double>(label, features);
temp.add(instance);
}
Collections.shuffle(temp);
for(Instance<Double> instance : temp)
{
if(testing_data.size() < 10)
{
testing_data.add(instance);
}
else
{
sample_data.add(instance);
}
}
}catch(FileNotFoundException ex)
{
ex.printStackTrace();
}catch(IOException ex)
{
ex.printStackTrace();
}finally
{
try {
if(is!=null) is.close();
if(br !=null) br.close();
} catch (IOException e) {
e.printStackTrace();
}
}
}
public BirthDataSpout()
{
}
public void open(Map conf, TopologyContext context) {
// TODO Auto-generated method stub
}
public void emitBatch(long batchId, TridentCollector collector) {
// TODO Auto-generated method stub
int maxBatchCount=sample_data.size() / batchSize;
if(maxBatchCount > 0 && batchIndex < maxBatchCount)
{
for(int i=batchIndex * batchSize; i < sample_data.size() && i < (batchIndex+1) * batchSize; ++i)
{
Instance<Double> instance = sample_data.get(i);
collector.emit(new Values(instance));
}
batchIndex=(batchIndex+1) % maxBatchCount;
}
}
public void ack(long batchId) {
// TODO Auto-generated method stub
}
public void close() {
// TODO Auto-generated method stub
}
public Map getComponentConfiguration() {
// TODO Auto-generated method stub
return null;
}
public Fields getOutputFields() {
// TODO Auto-generated method stub
return new Fields("instance");
}
}
As can be seen above, the BirthDataSpout is derived from IBatchSpout, and emits a batch of 10 tuples at one time, each tuple is a training record containing the fields ("instance"). The "instance" field contains a datatype Instance<Double> which contains a double array as features and a double value as label. The training records are obtained from a births.csv file residing in src/test/resources.
PA Regression in Trident topology using Trident-ML implementation
Once we have the training data spout, we can build a Trident topology which uses the training data to create a predicted output value for each of the data record using PA regression algorithm in Trident-ML. This is implemented in the main class shown below:package com.memeanalytics.trident_regression_pa;
import java.util.List;
import com.github.pmerienne.trident.ml.regression.PARegressor;
import com.github.pmerienne.trident.ml.regression.RegressionQuery;
import com.github.pmerienne.trident.ml.regression.RegressionUpdater;
import storm.trident.TridentState;
import storm.trident.TridentTopology;
import storm.trident.testing.MemoryMapState;
import backtype.storm.Config;
import backtype.storm.LocalCluster;
import backtype.storm.LocalDRPC;
import backtype.storm.generated.AlreadyAliveException;
import backtype.storm.generated.InvalidTopologyException;
import backtype.storm.generated.StormTopology;
import backtype.storm.tuple.Fields;
/**
* Hello world!
*
*/
public class App
{
public static void main( String[] args ) throws AlreadyAliveException, InvalidTopologyException
{
LocalDRPC drpc=new LocalDRPC();
LocalCluster cluster=new LocalCluster();
Config config=new Config();
cluster.submitTopology("RegressionDemo", config, buildTopology(drpc));
try{
Thread.sleep(10000);
}catch(InterruptedException ex)
{
ex.printStackTrace();
}
List<String> drpc_args_list=BirthDataSpout.getDRPCArgsList();
for(String drpc_args : drpc_args_list)
{
System.out.println(drpc.execute("predict", drpc_args));
}
cluster.killTopology("RegressionDemo");
cluster.shutdown();
drpc.shutdown();
}
private static StormTopology buildTopology(LocalDRPC drpc)
{
TridentTopology topology=new TridentTopology();
BirthDataSpout spout=new BirthDataSpout();
TridentState regressionModel = topology.newStream("training", spout).partitionPersist(new MemoryMapState.Factory(), new Fields("instance"), new RegressionUpdater("regression", new PARegressor()));
topology.newDRPCStream("predict", drpc).each(new Fields("args"), new DRPCArgsToInstance(), new Fields("instance")).stateQuery(regressionModel, new Fields("instance"), new RegressionQuery("regression"), new Fields("prediction")).project(new Fields("args", "prediction"));
return topology.build();
}
}
package com.memeanalytics.trident_regression_pa;
import backtype.storm.tuple.Values;
import com.github.pmerienne.trident.ml.core.Instance;
import storm.trident.operation.BaseFunction;
import storm.trident.operation.TridentCollector;
import storm.trident.tuple.TridentTuple;
public class DRPCArgsToInstance extends BaseFunction {
private static final long serialVersionUID = 1L;
public void execute(TridentTuple tuple, TridentCollector collector) {
String drpc_args = tuple.getString(0);
String[] args=drpc_args.split(",");
Double label=Double.parseDouble(args[args.length-1]);
double[] features=new double[args.length-1];
for(int i=0; i < args.length-1; ++i)
{
features[i]=Double.parseDouble(args[i]);
}
Instance<Double> instance=new Instance<Double>(label, features);
collector.emit(new Values(instance));
}
}
As can be seen above, the Trident topology has the BirthDataSpout emits Instance<Double> training data can be consumed by RegressionUpdater. The RegressionUpdater object from Trident-ML updates the underlying regressionModel via PA algorithm.
The DRPCStream allows user to pass in a new testing instance to the regressionModel which will then return a "predict" field, that contains the predicted output of the testing instance. The DRPCArgsToInstance is a BaseFunction operation which converts the arguments passed into the LocalDRPC.execute() into an Instance<Double> which can be passed into the RegressionQuery which then uses PARegressor and regressionModel to determine the predicted output value.
Once the coding is completed, we can run the project by navigating to the project root folder and run the following commands:
> .mvn compile exec:java
Labels:
BigData,
Java,
Machine Learning,
Storm,
Trident,
Trident-ML
Trident-ML: Sentiment Analysis Classifier
Trident-ML comes with a pre-trained twitter sentiment classifier, this post shows how to use this classifier to perform sentiment analysis in Storm.
This post shows some very basic example of how to use the pre-trained twitter sentiment classifier in Trident-ML to classifier sentiment of text which will return true (positive) or false (negative).
Firstly create a Maven project (e.g. with groupId="com.memeanalytics" artifactId="trident-sentiment-classifier"). The complete source codes of the project can be downloaded from the link:
https://dl.dropboxusercontent.com/u/113201788/storm/trident-sentiment-classifier.tar.gz
For the start we need to configure the pom.xml file in the project.
Next we need to add the storm dependency to the dependencies section (for storm):
Next we need to add the strident-ml dependency to the dependencies section (for text classification):
Next we need to add the exec-maven-plugin to the build/plugins section (for execute the Maven project):
Next we need to add the maven-assembly-plugin to the build/plugins section (for packacging the Maven project to jar for submitting to Storm cluster):
The DRPCStream allows user to pass in a text string to the TwitterSentimentClassifier which will then return a "sentiment" field, that contains the predicted label (true for positive; false for negative) of the testing text.
Next copy the following two files into the "main/resources" folder under the project root folder:
twitter-sentiment-classifier-classifier.json:
https://github.com/pmerienne/trident-ml/blob/master/src/main/resources/twitter-sentiment-classifier-classifier.json
twitter-sentiment-classifier-extractor.json:
https://github.com/pmerienne/trident-ml/blob/master/src/main/resources/twitter-sentiment-classifier-extractor.json
The above step can be important, otherwise you may get a FileNotFoundException during runtime.
Once the coding is completed, we can run the project by navigating to the project root folder and run the following commands:
> .mvn compile exec:java
This post shows some very basic example of how to use the pre-trained twitter sentiment classifier in Trident-ML to classifier sentiment of text which will return true (positive) or false (negative).
Firstly create a Maven project (e.g. with groupId="com.memeanalytics" artifactId="trident-sentiment-classifier"). The complete source codes of the project can be downloaded from the link:
https://dl.dropboxusercontent.com/u/113201788/storm/trident-sentiment-classifier.tar.gz
For the start we need to configure the pom.xml file in the project.
Configure pom.xml:
Firstly we need to add the clojars repository to the repositories section:<repositories> <repository> <id>clojars</id> <url>http://clojars.org/repo</url> </repository> </repositories>
Next we need to add the storm dependency to the dependencies section (for storm):
<dependency> <groupId>storm</groupId> <artifactId>storm</artifactId> <version>0.9.0.1</version> <scope>provided</scope> </dependency>
Next we need to add the strident-ml dependency to the dependencies section (for text classification):
<dependency> <groupId>com.github.pmerienne</groupId> <artifactId>trident-ml</artifactId> <version>0.0.4</version> </dependency>
Next we need to add the exec-maven-plugin to the build/plugins section (for execute the Maven project):
<plugin> <groupId>org.codehaus.mojo</groupId> <artifactId>exec-maven-plugin</artifactId> <version>1.2.1</version> <executions> <execution> <goals> <goal>exec</goal> </goals> </execution> </executions> <configuration> <includeProjectDependencies>true</includeProjectDependencies> <includePluginDependencies>false</includePluginDependencies> <executable>java</executable> <classpathScope>compile</classpathScope> <mainClass>com.memeanalytics.trident_sentiment_classifier.App</mainClass> </configuration> </plugin>
Next we need to add the maven-assembly-plugin to the build/plugins section (for packacging the Maven project to jar for submitting to Storm cluster):
<plugin> <artifactId>maven-assembly-plugin</artifactId> <version>2.2.1</version> <configuration> <descriptorRefs> <descriptorRef>jar-with-dependencies</descriptorRef> </descriptorRefs> <archive> <manifest> <mainClass></mainClass> </manifest> </archive> </configuration> <executions> <execution> <id>make-assembly</id> <phase>package</phase> <goals> <goal>single</goal> </goals> </execution> </executions> </plugin>
Sentiment Classification in Trident topology using Trident-ML implementation
Once the pom.xml update is completed, we can build a Trident topology which uses TwitterSentimentClassifier in a DRPCStream to classify text sentiment in Trident-ML. This is implemented in the main class shown below:package com.memeanalytics.trident_sentiment_classifier;
import com.github.pmerienne.trident.ml.nlp.TwitterSentimentClassifier;
import storm.trident.TridentTopology;
import backtype.storm.Config;
import backtype.storm.LocalCluster;
import backtype.storm.LocalDRPC;
import backtype.storm.generated.StormTopology;
import backtype.storm.tuple.Fields;
public class App
{
public static void main( String[] args )
{
LocalDRPC drpc=new LocalDRPC();
LocalCluster cluster=new LocalCluster();
Config config=new Config();
cluster.submitTopology("SentimentClassifierDemo", config, buildTopology(drpc));
try{
Thread.sleep(2000);
}catch(InterruptedException ex)
{
ex.printStackTrace();
}
System.out.println(drpc.execute("classify", "Have a nice day!"));
System.out.println(drpc.execute("classify", "I feel really bad!"));
System.out.println(drpc.execute("classify", "Whatever, i don't really care"));
System.out.println(drpc.execute("classify", "feel sleepy zzzz...."));
cluster.killTopology("SentimentClassifierDemo");
cluster.shutdown();
drpc.shutdown();
}
private static StormTopology buildTopology(LocalDRPC drpc)
{
TridentTopology topology=new TridentTopology();
topology.newDRPCStream("classify", drpc).each(new Fields("args"), new TwitterSentimentClassifier(), new Fields("sentiment"));
return topology.build();
}
}
The DRPCStream allows user to pass in a text string to the TwitterSentimentClassifier which will then return a "sentiment" field, that contains the predicted label (true for positive; false for negative) of the testing text.
Next copy the following two files into the "main/resources" folder under the project root folder:
twitter-sentiment-classifier-classifier.json:
https://github.com/pmerienne/trident-ml/blob/master/src/main/resources/twitter-sentiment-classifier-classifier.json
twitter-sentiment-classifier-extractor.json:
https://github.com/pmerienne/trident-ml/blob/master/src/main/resources/twitter-sentiment-classifier-extractor.json
The above step can be important, otherwise you may get a FileNotFoundException during runtime.
Once the coding is completed, we can run the project by navigating to the project root folder and run the following commands:
> .mvn compile exec:java
Trident-ML: Text Classification using KLD
This post shows some very basic example of how to use the Kullback-Leibler Distance text classification algorithm in Trident-ML to process data from Storm Spout.
Firstly create a Maven project (e.g. with groupId="com.memeanalytics" artifactId="trident-text-classifier-kld"). The complete source codes of the project can be downloaded from the link:
https://dl.dropboxusercontent.com/u/113201788/storm/trident-text-classifier-kld.tar.gz
For the start we need to configure the pom.xml file in the project.
Next we need to add the storm dependency to the dependencies section (for storm):
Next we need to add the strident-ml dependency to the dependencies section (for text classification):
Next we need to add the exec-maven-plugin to the build/plugins section (for execute the Maven project):
Next we need to add the maven-assembly-plugin to the build/plugins section (for packacging the Maven project to jar for submitting to Storm cluster):
As can be seen above, the ReuterNewsSpout is derived from IBatchSpout, and emits a batch of 10 tuples at one time, each tuple is a new article containing the fields ("label", "text"). The "label" field is integer value (represents the topic of the news article), while "text" field is a string which is text of the news article. the training records are obtained in such a way that the correct prediction learned from the text classification should be predicting the topic of a news article given the text of the news article.
As can be seen above, the Trident topology has a TextInstanceCreator<Integer> trident operation which convert raw ("label", "text") tuple into an TextInstance<Integer> object which can be consumed by TextClassifierUpdater. The TextClassifierUpdater object from Trident-ML updates the underlying classifierModel via KLDClassifier training algorithm.
The DRPCStream allows user to pass in a new testing instance to the classifierModel which will then return a "predict" field, that contains the predicted label of the testing instance. The DRPCArgsToInstance is a BaseFunction operation which converts the arguments passed into the LocalDRPC.execute() into an TextInstance<Integer> (Note you can set the label to null in DRPCArgsToInstance.execute() method as the label will be predicted instead) which can be passed into the ClassifyTextQuery which then uses KLD and classifierModel to determine the predicted label.
Once the coding is completed, we can run the project by navigating to the project root folder and run the following commands:
> .mvn compile exec:java
Firstly create a Maven project (e.g. with groupId="com.memeanalytics" artifactId="trident-text-classifier-kld"). The complete source codes of the project can be downloaded from the link:
https://dl.dropboxusercontent.com/u/113201788/storm/trident-text-classifier-kld.tar.gz
For the start we need to configure the pom.xml file in the project.
Configure pom.xml:
Firstly we need to add the clojars repository to the repositories section:<repositories> <repository> <id>clojars</id> <url>http://clojars.org/repo</url> </repository> </repositories>
Next we need to add the storm dependency to the dependencies section (for storm):
<dependency> <groupId>storm</groupId> <artifactId>storm</artifactId> <version>0.9.0.1</version> <scope>provided</scope> </dependency>
Next we need to add the strident-ml dependency to the dependencies section (for text classification):
<dependency> <groupId>com.github.pmerienne</groupId> <artifactId>trident-ml</artifactId> <version>0.0.4</version> </dependency>
Next we need to add the exec-maven-plugin to the build/plugins section (for execute the Maven project):
<plugin> <groupId>org.codehaus.mojo</groupId> <artifactId>exec-maven-plugin</artifactId> <version>1.2.1</version> <executions> <execution> <goals> <goal>exec</goal> </goals> </execution> </executions> <configuration> <includeProjectDependencies>true</includeProjectDependencies> <includePluginDependencies>false</includePluginDependencies> <executable>java</executable> <classpathScope>compile</classpathScope> <mainClass>com.memeanalytics.trident_text_classifier_kld.App</mainClass> </configuration> </plugin>
Next we need to add the maven-assembly-plugin to the build/plugins section (for packacging the Maven project to jar for submitting to Storm cluster):
<plugin> <artifactId>maven-assembly-plugin</artifactId> <version>2.2.1</version> <configuration> <descriptorRefs> <descriptorRef>jar-with-dependencies</descriptorRef> </descriptorRefs> <archive> <manifest> <mainClass></mainClass> </manifest> </archive> </configuration> <executions> <execution> <id>make-assembly</id> <phase>package</phase> <goals> <goal>single</goal> </goals> </execution> </executions> </plugin>
Implement Spout for training data
Once the pom.xml update is completed, we can move to implement the ReuterNewsSpout which is the Storm spout that emits batches of training data to the Trident topology:package com.memeanalytics.trident_text_classifier_kld;
import java.io.BufferedReader;
import java.io.FileInputStream;
import java.io.FileNotFoundException;
import java.io.IOException;
import java.io.InputStreamReader;
import java.util.ArrayList;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
import backtype.storm.task.TopologyContext;
import backtype.storm.tuple.Fields;
import backtype.storm.tuple.Values;
import storm.trident.operation.TridentCollector;
import storm.trident.spout.IBatchSpout;
public class ReuterNewsSpout implements IBatchSpout {
private static final long serialVersionUID = 1L;
private List<List<Object>> trainingData=new ArrayList<List<Object>>();
private static Map<Integer, List<Object>> testingData=new HashMap<Integer, List<Object>>();
private int batchSize=10;
private int batchIndex=0;
public ReuterNewsSpout()
{
try{
loadReuterNews();
}catch(FileNotFoundException ex)
{
ex.printStackTrace();
}catch(IOException ex)
{
ex.printStackTrace();
}
}
public static List<List<Object>> getTestingData()
{
List<List<Object>> result=new ArrayList<List<Object>>();
for(Integer topic_index : testingData.keySet())
{
result.add(testingData.get(topic_index));
}
return result;
}
private void loadReuterNews() throws FileNotFoundException, IOException
{
Map<String, Integer> topics=new HashMap<String, Integer>();
String filePath="src/test/resources/reuters.csv";
FileInputStream inputStream=new FileInputStream(filePath);
BufferedReader reader= new BufferedReader(new InputStreamReader(inputStream));
String line;
while((line = reader.readLine())!=null)
{
String topic = line.split(",")[0];
if(!topics.containsKey(topic))
{
topics.put(topic, topics.size());
}
Integer topic_index=topics.get(topic);
int index = line.indexOf(" - ");
if(index==-1) continue;
String text=line.substring(index, line.length()-1);
if(testingData.containsKey(topic_index))
{
List<Object> values=new ArrayList<Object>();
values.add(topic_index);
values.add(text);
trainingData.add(values);
}
else
{
testingData.put(topic_index, new Values(topic_index, text));
}
}
reader.close();
}
public void open(Map conf, TopologyContext context) {
// TODO Auto-generated method stub
}
public void emitBatch(long batchId, TridentCollector collector) {
// TODO Auto-generated method stub
int maxBatchIndex = (trainingData.size() / batchSize);
if(trainingData.size() > batchSize && batchIndex < maxBatchIndex)
{
for(int i=batchIndex * batchSize; i < trainingData.size() && i < (batchIndex+1) * batchSize; ++i)
{
collector.emit(trainingData.get(i));
}
batchIndex++;
//System.out.println("Progress: "+batchIndex +" / "+maxBatchIndex);
}
}
public void ack(long batchId) {
// TODO Auto-generated method stub
}
public void close() {
// TODO Auto-generated method stub
}
public Map getComponentConfiguration() {
// TODO Auto-generated method stub
return null;
}
public Fields getOutputFields() {
// TODO Auto-generated method stub
return new Fields("label", "text");
}
}
As can be seen above, the ReuterNewsSpout is derived from IBatchSpout, and emits a batch of 10 tuples at one time, each tuple is a new article containing the fields ("label", "text"). The "label" field is integer value (represents the topic of the news article), while "text" field is a string which is text of the news article. the training records are obtained in such a way that the correct prediction learned from the text classification should be predicting the topic of a news article given the text of the news article.
KLD Text Classification in Trident topology using Trident-ML implementation
Once we have the training data spout, we can build a Trident topology which uses the training data to create a class label for each of the data record using KLD classifier algorithm in Trident-ML. This is implemented in the main class shown below:package com.memeanalytics.trident_text_classifier_kld;
import java.util.List;
import com.github.pmerienne.trident.ml.nlp.ClassifyTextQuery;
import com.github.pmerienne.trident.ml.nlp.KLDClassifier;
import com.github.pmerienne.trident.ml.nlp.TextClassifierUpdater;
import com.github.pmerienne.trident.ml.preprocessing.TextInstanceCreator;
import storm.trident.TridentState;
import storm.trident.TridentTopology;
import storm.trident.testing.MemoryMapState;
import backtype.storm.Config;
import backtype.storm.LocalCluster;
import backtype.storm.LocalDRPC;
import backtype.storm.generated.AlreadyAliveException;
import backtype.storm.generated.InvalidTopologyException;
import backtype.storm.generated.StormTopology;
import backtype.storm.tuple.Fields;
public class App
{
public static void main( String[] args ) throws AlreadyAliveException, InvalidTopologyException
{
LocalDRPC drpc=new LocalDRPC();
LocalCluster cluster=new LocalCluster();
Config config=new Config();
cluster.submitTopology("KLDDemo", config, buildTopology(drpc));
try{
Thread.sleep(20000);
}catch(InterruptedException ex)
{
ex.printStackTrace();
}
List<List<Object>> testingData = ReuterNewsSpout.getTestingData();
for(int i=0; i < testingData.size(); ++i)
{
List<Object> testingDataRecord=testingData.get(i);
String drpc_args="";
for(Object val : testingDataRecord){
if(drpc_args.equals(""))
{
drpc_args+=val;
}
else
{
drpc_args+=(","+val);
}
}
System.out.println(drpc.execute("predict", drpc_args));
}
cluster.killTopology("KLDDemo");
cluster.shutdown();
drpc.shutdown();
}
private static StormTopology buildTopology(LocalDRPC drpc)
{
ReuterNewsSpout spout=new ReuterNewsSpout();
TridentTopology topology=new TridentTopology();
TridentState classifierModel = topology.newStream("training", spout).each(new Fields("label", "text"), new TextInstanceCreator<Integer>(), new Fields("instance")).partitionPersist(new MemoryMapState.Factory(), new Fields("instance"), new TextClassifierUpdater("newsClassifier", new KLDClassifier(9)));
topology.newDRPCStream("predict", drpc).each(new Fields("args"), new DRPCArgsToInstance(), new Fields("instance")).stateQuery(classifierModel, new Fields("instance"), new ClassifyTextQuery("newsClassifier"), new Fields("prediction"));
return topology.build();
}
}
package com.memeanalytics.trident_text_classifier_kld;
import java.util.ArrayList;
import java.util.List;
import backtype.storm.tuple.Values;
import com.github.pmerienne.trident.ml.core.TextInstance;
import com.github.pmerienne.trident.ml.preprocessing.EnglishTokenizer;
import com.github.pmerienne.trident.ml.preprocessing.TextTokenizer;
import storm.trident.operation.BaseFunction;
import storm.trident.operation.TridentCollector;
import storm.trident.tuple.TridentTuple;
public class DRPCArgsToInstance extends BaseFunction{
private static final long serialVersionUID = 1L;
public void execute(TridentTuple tuple, TridentCollector collector) {
// TODO Auto-generated method stub
String drpc_args=tuple.getString(0);
String[] args=drpc_args.split(",");
Integer label=Integer.parseInt(args[0]);
String text=args[1];
TextTokenizer textAnalyzer=new EnglishTokenizer();
List<String> tokens=textAnalyzer.tokenize(text);
TextInstance<Integer> instance=new TextInstance<Integer>(label, tokens);
collector.emit(new Values(instance));
}
}
As can be seen above, the Trident topology has a TextInstanceCreator<Integer> trident operation which convert raw ("label", "text") tuple into an TextInstance<Integer> object which can be consumed by TextClassifierUpdater. The TextClassifierUpdater object from Trident-ML updates the underlying classifierModel via KLDClassifier training algorithm.
The DRPCStream allows user to pass in a new testing instance to the classifierModel which will then return a "predict" field, that contains the predicted label of the testing instance. The DRPCArgsToInstance is a BaseFunction operation which converts the arguments passed into the LocalDRPC.execute() into an TextInstance<Integer> (Note you can set the label to null in DRPCArgsToInstance.execute() method as the label will be predicted instead) which can be passed into the ClassifyTextQuery which then uses KLD and classifierModel to determine the predicted label.
Once the coding is completed, we can run the project by navigating to the project root folder and run the following commands:
> .mvn compile exec:java
Labels:
BigData,
Java,
Machine Learning,
Storm,
Trident,
Trident-ML
Trident-ML: Classification using Perceptron
This post shows some very basic example of how to use the perceptron classification algorithm in Trident-ML to process data from Storm Spout.
Firstly create a Maven project (e.g. with groupId="com.memeanalytics" artifactId="trident-classifier-perceptron"). The complete source codes of the project can be downloaded from the link:
https://dl.dropboxusercontent.com/u/113201788/storm/trident-classifier-perceptron.tar.gz
For the start we need to configure the pom.xml file in the project.
Configure pom.xml:
Firstly we need to add the clojars repository to the repositories section:
Next we need to add the storm dependency to the dependencies section (for storm):
Next we need to add the strident-ml dependency to the dependencies section (for perceptron classification):
Next we need to add the exec-maven-plugin to the build/plugins section (for execute the Maven project):
Next we need to add the maven-assembly-plugin to the build/plugins section (for packacging the Maven project to jar for submitting to Storm cluster):
As can be seen above, the NANDSpout is derived from IBatchSpout, and emits a batch of 10 tuples at one time, each tuple is a training record containing the fields ("label", "x0", "x1"). The label is boolean value, while x0, x1 are double values which are either 1 (true) or 0 (false). the training records are obtained in such a way that the correct prediction should be a NAND gate from the classification.
As can be seen above, the Trident topology has a InstanceCreator<Boolean> trident operation which convert raw ("label", "x0", "x1") tuple into an Instance<Boolean> object which can be consumed by ClassifierUpdater. The ClassifierUpdater object from Trident-ML updates the underlying classifierModel via perceptron training algorithm.
The DRPCStream allows user to pass in a new testing instance to the classifierModel which will then return a "predict" field, that contains the predicted label of the testing instance. The DRPCArgsToInstance is a BaseFunction operation which converts the arguments passed into the LocalDRPC.execute() into an Instance<Boolean> which can be passed into the ClassifyQuery which then uses perceptron and classifierModel to determine the predicted label.
Once the coding is completed, we can run the project by navigating to the project root folder and run the following commands:
> .mvn compile exec:java
Firstly create a Maven project (e.g. with groupId="com.memeanalytics" artifactId="trident-classifier-perceptron"). The complete source codes of the project can be downloaded from the link:
https://dl.dropboxusercontent.com/u/113201788/storm/trident-classifier-perceptron.tar.gz
For the start we need to configure the pom.xml file in the project.
Configure pom.xml:
Firstly we need to add the clojars repository to the repositories section:
<repositories> <repository> <id>clojars</id> <url>http://clojars.org/repo</url> </repository> </repositories>
Next we need to add the storm dependency to the dependencies section (for storm):
<dependency> <groupId>storm</groupId> <artifactId>storm</artifactId> <version>0.9.0.1</version> <scope>provided</scope> </dependency>
Next we need to add the strident-ml dependency to the dependencies section (for perceptron classification):
<dependency> <groupId>com.github.pmerienne</groupId> <artifactId>trident-ml</artifactId> <version>0.0.4</version> </dependency>
Next we need to add the exec-maven-plugin to the build/plugins section (for execute the Maven project):
<plugin> <groupId>org.codehaus.mojo</groupId> <artifactId>exec-maven-plugin</artifactId> <version>1.2.1</version> <executions> <execution> <goals> <goal>exec</goal> </goals> </execution> </executions> <configuration> <includeProjectDependencies>true</includeProjectDependencies> <includePluginDependencies>false</includePluginDependencies> <executable>java</executable> <classpathScope>compile</classpathScope> <mainClass>com.memeanalytics.trident_classifier_perceptron.App</mainClass> </configuration> </plugin>
Next we need to add the maven-assembly-plugin to the build/plugins section (for packacging the Maven project to jar for submitting to Storm cluster):
<plugin> <artifactId>maven-assembly-plugin</artifactId> <version>2.2.1</version> <configuration> <descriptorRefs> <descriptorRef>jar-with-dependencies</descriptorRef> </descriptorRefs> <archive> <manifest> <mainClass></mainClass> </manifest> </archive> </configuration> <executions> <execution> <id>make-assembly</id> <phase>package</phase> <goals> <goal>single</goal> </goals> </execution> </executions> </plugin>
Implement Spout for training data
Once the pom.xml update is completed, we can move to implement the NANDSpout which is the Storm spout that emits batches of training data to the Trident topology:package com.memeanalytics.trident_classifier_perceptron;
import java.util.ArrayList;
import java.util.List;
import java.util.Map;
import java.util.Random;
import backtype.storm.task.TopologyContext;
import backtype.storm.tuple.Fields;
import storm.trident.operation.TridentCollector;
import storm.trident.spout.IBatchSpout;
public class NANDSpout implements IBatchSpout {
private int batchSize=10;
public void open(Map conf, TopologyContext context) {
// TODO Auto-generated method stub
}
public void emitBatch(long batchId, TridentCollector collector) {
// TODO Auto-generated method stub
final Random rand=new Random();
for(int i=0; i < batchSize; ++i)
{
boolean x0=rand.nextBoolean();
boolean x1=rand.nextBoolean();
boolean label = !(x0 && x1);
List<Object> values=new ArrayList<Object>();
values.add(label);
values.add(x0 ? 1.0 : 0.0);
values.add(x1 ? 1.0 : 0.0);
//values.add(x0 ? 1.0 + noise(rand) : 0.0 + noise(rand));
//values.add(x1 ? 1.0 + noise(rand) : 0.0 + noise(rand));
collector.emit(values);
}
}
public static double noise(Random rand)
{
return rand.nextDouble()* 0.0001 - 0.00005;
}
public void ack(long batchId) {
// TODO Auto-generated method stub
}
public void close() {
// TODO Auto-generated method stub
}
public Map getComponentConfiguration() {
// TODO Auto-generated method stub
return null;
}
public Fields getOutputFields() {
// TODO Auto-generated method stub
return new Fields("label", "x0", "x1");
}
}
As can be seen above, the NANDSpout is derived from IBatchSpout, and emits a batch of 10 tuples at one time, each tuple is a training record containing the fields ("label", "x0", "x1"). The label is boolean value, while x0, x1 are double values which are either 1 (true) or 0 (false). the training records are obtained in such a way that the correct prediction should be a NAND gate from the classification.
Perceptron Classification in Trident topology using Trident-ML implementation
Once we have the training data spout, we can build a Trident topology which uses the training data to create a class label for each of the data record using perceptron classifier algorithm in Trident-ML. This is implemented in the main class shown below:package com.memeanalytics.trident_classifier_perceptron;
import java.util.Random;
import com.github.pmerienne.trident.ml.classification.ClassifierUpdater;
import com.github.pmerienne.trident.ml.classification.ClassifyQuery;
import com.github.pmerienne.trident.ml.classification.PerceptronClassifier;
import com.github.pmerienne.trident.ml.preprocessing.InstanceCreator;
import storm.trident.TridentState;
import storm.trident.TridentTopology;
import storm.trident.testing.MemoryMapState;
import backtype.storm.Config;
import backtype.storm.LocalCluster;
import backtype.storm.LocalDRPC;
import backtype.storm.generated.AlreadyAliveException;
import backtype.storm.generated.InvalidTopologyException;
import backtype.storm.generated.StormTopology;
import backtype.storm.tuple.Fields;
/**
* Hello world!
*
*/
public class App
{
public static void main( String[] args ) throws AlreadyAliveException, InvalidTopologyException
{
LocalDRPC drpc=new LocalDRPC();
LocalCluster cluster=new LocalCluster();
Config config=new Config();
cluster.submitTopology("PerceptronDemo", config, buildTopology(drpc));
try{
Thread.sleep(10000);
}catch(InterruptedException ex)
{
ex.printStackTrace();
}
for(int i=0; i < 10; ++i)
{
String drpc_args=createDRPCTestingSample();
System.out.println(drpc.execute("predict", drpc_args));
try{
Thread.sleep(1000);
}catch(InterruptedException ex)
{
ex.printStackTrace();
}
}
cluster.killTopology("PerceptronDemo");
cluster.shutdown();
drpc.shutdown();
}
private static String createDRPCTestingSample()
{
String drpc_args="";
final Random rand=new Random();
boolean bit_x0=rand.nextBoolean();
boolean bit_x1=rand.nextBoolean();
boolean label = !(bit_x0 && bit_x1);
double x0=bit_x0 ? 1.0 + NANDSpout.noise(rand) : 0.0 + NANDSpout.noise(rand);
double x1=bit_x1 ? 1.0 + NANDSpout.noise(rand) : 0.0 + NANDSpout.noise(rand);
drpc_args+=label;
drpc_args+=(","+x0);
drpc_args+=(","+x1);
return drpc_args;
}
private static StormTopology buildTopology(LocalDRPC drpc)
{
TridentTopology topology=new TridentTopology();
NANDSpout spout=new NANDSpout();
TridentState classifierModel = topology.newStream("training", spout).shuffle().each(new Fields("label", "x0", "x1"), new InstanceCreator<Boolean>(), new Fields("instance")).partitionPersist(new MemoryMapState.Factory(), new Fields("instance"), new ClassifierUpdater<Boolean>("perceptron", new PerceptronClassifier()));
topology.newDRPCStream("predict", drpc).each(new Fields("args"), new DRPCArgsToInstance(), new Fields("instance")).stateQuery(classifierModel, new Fields("instance"), new ClassifyQuery<Boolean>("perceptron"), new Fields("predict"));
return topology.build();
}
}
package com.memeanalytics.trident_classifier_perceptron;
import backtype.storm.tuple.Values;
import com.github.pmerienne.trident.ml.core.Instance;
import storm.trident.operation.BaseFunction;
import storm.trident.operation.TridentCollector;
import storm.trident.tuple.TridentTuple;
public class DRPCArgsToInstance extends BaseFunction {
private static final long serialVersionUID = 1L;
public void execute(TridentTuple tuple, TridentCollector collector) {
// TODO Auto-generated method stub
String drpc_args=tuple.getString(0);
String[] args=drpc_args.split(",");
boolean label=Boolean.parseBoolean(args[0]);
double[] features=new double[args.length-1];
for(int i=1; i < args.length; ++i)
{
features[i-1]=Double.parseDouble(args[i]);
}
Instance<Boolean> instance=new Instance<Boolean>(label, features);
collector.emit(new Values(instance));
}
}
As can be seen above, the Trident topology has a InstanceCreator<Boolean> trident operation which convert raw ("label", "x0", "x1") tuple into an Instance<Boolean> object which can be consumed by ClassifierUpdater. The ClassifierUpdater object from Trident-ML updates the underlying classifierModel via perceptron training algorithm.
The DRPCStream allows user to pass in a new testing instance to the classifierModel which will then return a "predict" field, that contains the predicted label of the testing instance. The DRPCArgsToInstance is a BaseFunction operation which converts the arguments passed into the LocalDRPC.execute() into an Instance<Boolean> which can be passed into the ClassifyQuery which then uses perceptron and classifierModel to determine the predicted label.
Once the coding is completed, we can run the project by navigating to the project root folder and run the following commands:
> .mvn compile exec:java
Labels:
BigData,
Java,
Machine Learning,
Storm,
Trident,
Trident-ML
Monday, November 24, 2014
Trident-ML: Clustering using K-Means
This post shows some very basic example of how to use the k means clustering algorithm in Trident-ML to process data from Storm Spout.
Firstly create a Maven project (e.g. with groupId="com.memeanalytics" artifactId="trident-k-means"). The complete source codes of the project can be downloaded from the link:
https://dl.dropboxusercontent.com/u/113201788/storm/trident-k-means.tar.gz
For the start we need to configure the pom.xml file in the project.
Configure pom.xml:
Firstly we need to add the clojars repository to the repositories section:
Next we need to add the storm dependency to the dependencies section (for storm):
Next we need to add the strident-ml dependency to the dependencies section (for k-means clustering):
Next we need to add the exec-maven-plugin to the build/plugins section (for execute the Maven project):
Next we need to add the maven-assembly-plugin to the build/plugins section (for packacging the Maven project to jar for submitting to Storm cluster):
As can be seen above, the RandomFeatureSpout is derived from IBatchSpout, and emits a batch of 10 tuples at one time, each tuple is a training record containing the fields ("label", "x0", "x1", "x2"). The label is integer, while x0, x1, x2 are double values. the training records are obtained from Trident-ML's DataSets.generateDataForMultiLabelClassification() method.
As can be seen above, the Trident topology has a InstanceCreator<Integer> trident operation which convert raw ("label", "x0", "x1", "x2") tuple into an Instance<Integer> object which can be consumed by ClusterUpdator. The ClusterUpdate object from Trident-ML updates the underlying clusterModel via k-Means algorithm.
The DRPCStream allows user to pass in a new testing instance to the clusterModel which will then return a "predict" field, that contains the predicted label of the testing instance. The DRPCArgsToInstance is a BaseFunction operation which converts the arguments passed into the LocalDRPC.execute() into an Instance<Integer> which can be passed into the ClusterQuery which then uses kmeans and clusterModel to determine the predicted label.
Once the coding is completed, we can run the project by navigating to the project root folder and run the following commands:
> .mvn compile exec:java
Firstly create a Maven project (e.g. with groupId="com.memeanalytics" artifactId="trident-k-means"). The complete source codes of the project can be downloaded from the link:
https://dl.dropboxusercontent.com/u/113201788/storm/trident-k-means.tar.gz
For the start we need to configure the pom.xml file in the project.
Configure pom.xml:
Firstly we need to add the clojars repository to the repositories section:
<repositories> <repository> <id>clojars</id> <url>http://clojars.org/repo</url> </repository> </repositories>
Next we need to add the storm dependency to the dependencies section (for storm):
<dependency> <groupId>storm</groupId> <artifactId>storm</artifactId> <version>0.9.0.1</version> <scope>provided</scope> </dependency>
Next we need to add the strident-ml dependency to the dependencies section (for k-means clustering):
<dependency> <groupId>com.github.pmerienne</groupId> <artifactId>trident-ml</artifactId> <version>0.0.4</version> </dependency>
Next we need to add the exec-maven-plugin to the build/plugins section (for execute the Maven project):
<plugin> <groupId>org.codehaus.mojo</groupId> <artifactId>exec-maven-plugin</artifactId> <version>1.2.1</version> <executions> <execution> <goals> <goal>exec</goal> </goals> </execution> </executions> <configuration> <includeProjectDependencies>true</includeProjectDependencies> <includePluginDependencies>false</includePluginDependencies> <executable>java</executable> <classpathScope>compile</classpathScope> <mainClass>com.memeanalytics.trident_k_means.App</mainClass> </configuration> </plugin>
Next we need to add the maven-assembly-plugin to the build/plugins section (for packacging the Maven project to jar for submitting to Storm cluster):
<plugin> <artifactId>maven-assembly-plugin</artifactId> <version>2.2.1</version> <configuration> <descriptorRefs> <descriptorRef>jar-with-dependencies</descriptorRef> </descriptorRefs> <archive> <manifest> <mainClass></mainClass> </manifest> </archive> </configuration> <executions> <execution> <id>make-assembly</id> <phase>package</phase> <goals> <goal>single</goal> </goals> </execution> </executions> </plugin>
Implement Spout for training data
Once the pom.xml update is completed, we can move to implement the RandomFeatureSpout which is the Storm spout that emits batches of training data to the Trident topology:package com.memeanalytics.trident_k_means;
import java.util.ArrayList;
import java.util.List;
import java.util.Map;
import com.github.pmerienne.trident.ml.core.Instance;
import com.github.pmerienne.trident.ml.testing.data.Datasets;
import backtype.storm.task.TopologyContext;
import backtype.storm.tuple.Fields;
import storm.trident.operation.TridentCollector;
import storm.trident.spout.IBatchSpout;
public class RandomFeatureSpout implements IBatchSpout{
private int batchSize=10;
private int numFeatures=3;
private int numClasses=3;
public void open(Map conf, TopologyContext context) {
// TODO Auto-generated method stub
}
public void emitBatch(long batchId, TridentCollector collector) {
// TODO Auto-generated method stub
List<Instance<Integer>> data = Datasets.generateDataForMultiLabelClassification(batchSize, numFeatures, numClasses);
for(Instance<Integer> instance : data)
{
List<Object> values=new ArrayList<Object>();
values.add(instance.label);
for(double feature : instance.getFeatures())
{
values.add(feature);
}
collector.emit(values);
}
}
public void ack(long batchId) {
// TODO Auto-generated method stub
}
public void close() {
// TODO Auto-generated method stub
}
public Map getComponentConfiguration() {
// TODO Auto-generated method stub
return null;
}
public Fields getOutputFields() {
// TODO Auto-generated method stub
return new Fields("label", "x0", "x1", "x2");
}
}
As can be seen above, the RandomFeatureSpout is derived from IBatchSpout, and emits a batch of 10 tuples at one time, each tuple is a training record containing the fields ("label", "x0", "x1", "x2"). The label is integer, while x0, x1, x2 are double values. the training records are obtained from Trident-ML's DataSets.generateDataForMultiLabelClassification() method.
K-means in Trident topology using Trident-ML implementation
Once we have the training data spout, we can build a Trident topology which uses the training data to create a class label for each of the data record using k-means algorithm in Trident-ML. This is implemented in the main class shown below:package com.memeanalytics.trident_k_means;
import java.util.ArrayList;
import java.util.List;
import java.util.Random;
import com.github.pmerienne.trident.ml.clustering.ClusterQuery;
import com.github.pmerienne.trident.ml.clustering.ClusterUpdater;
import com.github.pmerienne.trident.ml.clustering.KMeans;
import com.github.pmerienne.trident.ml.core.Instance;
import com.github.pmerienne.trident.ml.preprocessing.InstanceCreator;
import com.github.pmerienne.trident.ml.testing.data.Datasets;
import storm.trident.TridentState;
import storm.trident.TridentTopology;
import storm.trident.testing.MemoryMapState;
import backtype.storm.Config;
import backtype.storm.LocalCluster;
import backtype.storm.LocalDRPC;
import backtype.storm.generated.AlreadyAliveException;
import backtype.storm.generated.InvalidTopologyException;
import backtype.storm.generated.StormTopology;
import backtype.storm.tuple.Fields;
public class App
{
public static void main( String[] args ) throws AlreadyAliveException, InvalidTopologyException
{
LocalDRPC drpc=new LocalDRPC();
Config config=new Config();
LocalCluster cluster=new LocalCluster();
cluster.submitTopology("KMeansDemo", config, buildTopology(drpc));
try{
Thread.sleep(10000);
}catch(InterruptedException ex)
{
ex.printStackTrace();
}
for(int i=0; i < 10; ++i)
{
String drpc_args=generateRandomTestingArgs();
System.out.println(drpc.execute("predict", drpc_args));
try{
Thread.sleep(1000);
}catch(InterruptedException ex)
{
ex.printStackTrace();
}
}
cluster.killTopology("KMeansDemo");
cluster.shutdown();
drpc.shutdown();
}
private static String generateRandomTestingArgs()
{
int batchSize=10;
int numFeatures=3;
int numClasses=3;
final Random rand=new Random();
List<Instance<Integer>> data = Datasets.generateDataForMultiLabelClassification(batchSize, numFeatures, numClasses);
String args="";
Instance<Integer> instance = data.get(rand.nextInt(data.size()));
args+=instance.label;
for(double feature : instance.getFeatures())
{
args+=(","+feature);
}
return args;
}
private static StormTopology buildTopology(LocalDRPC drpc)
{
TridentTopology topology=new TridentTopology();
RandomFeatureSpout spout=new RandomFeatureSpout();
TridentState clusterModel = topology.newStream("training", spout).each(new Fields("label", "x0", "x1", "x2"), new InstanceCreator<Integer>(), new Fields("instance")).partitionPersist(new MemoryMapState.Factory(), new Fields("instance"), new ClusterUpdater("kmeans", new KMeans(3)));
topology.newDRPCStream("predict", drpc).each(new Fields("args"), new DRPCArgsToInstance(), new Fields("instance")).stateQuery(clusterModel, new Fields("instance"), new ClusterQuery("kmeans"), new Fields("predict"));
return topology.build();
}
}
package com.memeanalytics.trident_k_means;
import java.util.ArrayList;
import java.util.List;
import backtype.storm.tuple.Values;
import com.github.pmerienne.trident.ml.core.Instance;
import storm.trident.operation.BaseFunction;
import storm.trident.operation.TridentCollector;
import storm.trident.tuple.TridentTuple;
public class DRPCArgsToInstance extends BaseFunction{
private static final long serialVersionUID = 1L;
public void execute(TridentTuple tuple, TridentCollector collector) {
// TODO Auto-generated method stub
String drpc_args = tuple.getString(0);
String[] args = drpc_args.split(",");
Integer label=Integer.parseInt(args[0]);
double[] features=new double[args.length-1];
for(int i=1; i < args.length; ++i)
{
double feature=Double.parseDouble(args[i]);
features[i-1] = feature;
}
Instance<Integer> instance=new Instance<Integer>(label, features);
collector.emit(new Values(instance));
}
}
As can be seen above, the Trident topology has a InstanceCreator<Integer> trident operation which convert raw ("label", "x0", "x1", "x2") tuple into an Instance<Integer> object which can be consumed by ClusterUpdator. The ClusterUpdate object from Trident-ML updates the underlying clusterModel via k-Means algorithm.
The DRPCStream allows user to pass in a new testing instance to the clusterModel which will then return a "predict" field, that contains the predicted label of the testing instance. The DRPCArgsToInstance is a BaseFunction operation which converts the arguments passed into the LocalDRPC.execute() into an Instance<Integer> which can be passed into the ClusterQuery which then uses kmeans and clusterModel to determine the predicted label.
Once the coding is completed, we can run the project by navigating to the project root folder and run the following commands:
> .mvn compile exec:java
Integration of Kafka-Storm-MySQL
This post will discuss how to create a minimum basic storm topology which integrate Kafka, Storm and MySQL. The scenerios is: a Kafka producer will push some dummy data to the Kafka brokers, and a KafkaSpout (with consumer group id = "id7" and zookeeper = "192.168.2.4:2181") from storm cluster consumes data from the Kafka brokers. The KafkaSpout will then emits Kafka messages as tuples to a BaseBasicBolt which will then persists the data to the MySQL server.
The post assumes user already have a zookeeper cluster set up on two hostname:port:
192.168.2.2:2181
192.168.2.4:2181
The post also assumes that the Kafka brokers runs at the following hostname:port:
192.168.2.2:9092
192.168.2.4:9092
Details of how to set up ZooKeeper and Kafka cluster can be found at the following links:
http://czcodezone.blogspot.sg/2014/11/setup-zookeeper-in-cluster.html
http://czcodezone.blogspot.sg/2014/11/setup-kafka-in-cluster.html
The Kafka producer source codes can be downloaded from the link: https://dl.dropboxusercontent.com/u/113201788/kafka-producer.zip. The details about how to implement a simple Kafka producer in java can be found at this link: http://czcodezone.blogspot.sg/2014/11/write-and-test-simple-kafka-producer.html).
To create the Storm topology which consumes messages from Kafka and persists them to MySQL database, create a Maven project (e.g. with groupId="com.memeanalytics" and artifactId="storm-kafka-mysql"). The complete source code of the project can be downloaded from the following link:
https://dl.dropboxusercontent.com/u/113201788/storm/storm-kafka-mysql.tar.gz
Next we add the mysql dependency into the dependencies section:
Next we add the storm-kafka-0.8-plus dependency into the dependencies section (for KafkaSpout which consumers Kafka messages in Storm):
Next we add the storm-core dependency into the dependencies section (for Storm):
Next we add the exec-maven-plugin into the build/plugins section (for compile and execute java project):
Next we add the maven-assembly-plugin into the build/plugins section (for packaging the java project as jar to submit to the Storm cluster):
As illustrated above, the MySqlConnection manages the opening and closing of the connection.
Next we will create a class which saves the tuples emitted from the KafkaSpout into the MySQL server.
The above classes opens the mysql connection in its constructor, and will saves the tuple via its persist() method, it also has a close() method which can be invoked to close the mysql connection.
Before we proceed further, we need to create the necessary database and datatable in mysql so that data can be inserted into. For demo purpose, we create a very basic datatable named "mylogtable" in a database named "mylog". To do this, access the mysql server by running the command:
> mysql -u [username] -p
Once logged into the mysql, run the following mysql queries to create the mylogtable:
> create database mylog;
> user mylog;
> create table mylogtable( id INT NOT NULL AUTO_INCREMENT, dataitem VARCHAR(255) NOT NULL, PRIMARY KEY(id));
As can be seen above, the bolt opens the mysql database connection in its prepare() method and close the mysql database connection in its cleanup() method. it also persists the tuple received into mysql datatable in its execute() method.
> mvn exec:java -Dmain.class=com.memeanalytics.storm_kafka_mysql.App
Note that if you change the ${main.class} to com.memeanalytics.storm_kafka_mysql.App, that you don't need to add the -Dmain.class arguments in the above command.
The post assumes user already have a zookeeper cluster set up on two hostname:port:
192.168.2.2:2181
192.168.2.4:2181
The post also assumes that the Kafka brokers runs at the following hostname:port:
192.168.2.2:9092
192.168.2.4:9092
Details of how to set up ZooKeeper and Kafka cluster can be found at the following links:
http://czcodezone.blogspot.sg/2014/11/setup-zookeeper-in-cluster.html
http://czcodezone.blogspot.sg/2014/11/setup-kafka-in-cluster.html
The Kafka producer source codes can be downloaded from the link: https://dl.dropboxusercontent.com/u/113201788/kafka-producer.zip. The details about how to implement a simple Kafka producer in java can be found at this link: http://czcodezone.blogspot.sg/2014/11/write-and-test-simple-kafka-producer.html).
To create the Storm topology which consumes messages from Kafka and persists them to MySQL database, create a Maven project (e.g. with groupId="com.memeanalytics" and artifactId="storm-kafka-mysql"). The complete source code of the project can be downloaded from the following link:
https://dl.dropboxusercontent.com/u/113201788/storm/storm-kafka-mysql.tar.gz
Maven setup: pom.xml
We will start by explaining setup in the pom.xml file. Firstly, we need to put in the clojars in the repositories:<repositories> <repository> <id>clojars</id> <url>http://clojars.org/repo</url> </repository> </repositories>
Next we add the mysql dependency into the dependencies section:
<dependency> <groupId>mysql</groupId> <artifactId>mysql-connector-java</artifactId> <version>5.1.6</version> </dependency>
Next we add the storm-kafka-0.8-plus dependency into the dependencies section (for KafkaSpout which consumers Kafka messages in Storm):
<dependency> <groupId>net.wurstmeister.storm</groupId> <artifactId>storm-kafka-0.8-plus</artifactId> <version>0.4.0</version> </dependency>
Next we add the storm-core dependency into the dependencies section (for Storm):
<dependency> <groupId>storm</groupId> <artifactId>storm-core</artifactId> <version>0.9.0.1</version> <scope>provided</scope> </dependency>
Next we add the exec-maven-plugin into the build/plugins section (for compile and execute java project):
<plugin>
<groupId>org.codehaus.mojo</groupId>
<artifactId>exec-maven-plugin</artifactId>
<version>1.2.1</version>
<executions>
<execution>
<goals>
<goal>exec</goal>
</goals>
</execution>
</executions>
<configuration>
<includeProjectDependencies>true</includeProjectDependencies>
<includePluginDependencies>false</includePluginDependencies>
<executable>java</executable>
<classpathScope>compile</classpathScope>
<mainClass>${main.class}</mainClass>
</configuration>
</plugin>
Next we add the maven-assembly-plugin into the build/plugins section (for packaging the java project as jar to submit to the Storm cluster):
<plugin> <artifactId>maven-assembly-plugin</artifactId> <version>2.2.1</version> <configuration> <descriptorRefs> <descriptorRef>jar-with-dependencies</descrptorRef> </descriptorRefs> <archive> <manifest> <mainClass></mainClass> </manifest> </archive> </configuration> <executions> <execution> <id>make-assembly</id> <phase>package</phase> <goals> <goal>single</goal> </goals> </execution> </executions> </plugin>
MySQL Connection and Management Classes
Next we will create a class that manages MySQL connection, the source codes of which are show below:package com.memeanalytics.storm_kafka_mysql;
import java.sql.Connection;
import java.sql.DriverManager;
public class MySqlConnection {
private String ip;
private String database;
private String username;
private String password;
private Connection conn;
public MySqlConnection(String ip, String database, String username, String password)
{
this.ip=ip;
this.database=database;
this.username=username;
this.password=password;
}
public Connection getConnection()
{
return conn;
}
public boolean open()
{
boolean successful=true;
try{
Class.forName("com.mysql.jdbc.Driver");
conn = DriverManager.getConnection("jdbc:mysql://"+ip+"/"+database+"?"+"user="+username+"&password="+password);
}catch(Exception ex)
{
successful=false;
ex.printStackTrace();
}
return successful;
}
public boolean close()
{
if(conn==null)
{
return false;
}
boolean successful=true;
try{
conn.close();
}catch(Exception ex)
{
successful=false;
ex.printStackTrace();
}
return successful;
}
}
As illustrated above, the MySqlConnection manages the opening and closing of the connection.
Next we will create a class which saves the tuples emitted from the KafkaSpout into the MySQL server.
package com.memeanalytics.storm_kafka_mysql;
import java.sql.PreparedStatement;
import backtype.storm.tuple.Tuple;
public class MySqlDump {
private MySqlConnection conn;
public MySqlDump(String ip, String database, String username, String password)
{
conn = new MySqlConnection(ip, database, username, password);
conn.open();
}
public void persist(Tuple tuple)
{
PreparedStatement statement=null;
try{
statement = conn.getConnection().prepareStatement("insert into mylogtable (id, dataitem) values (default, ?)");
statement.setString(1, tuple.getString(0));
statement.executeUpdate();
}catch(Exception ex)
{
ex.printStackTrace();
}finally
{
if(statement != null)
{
try{
statement.close();
}catch(Exception ex)
{
ex.printStackTrace();
}
}
}
}
public void close()
{
conn.close();
}
}
The above classes opens the mysql connection in its constructor, and will saves the tuple via its persist() method, it also has a close() method which can be invoked to close the mysql connection.
Before we proceed further, we need to create the necessary database and datatable in mysql so that data can be inserted into. For demo purpose, we create a very basic datatable named "mylogtable" in a database named "mylog". To do this, access the mysql server by running the command:
> mysql -u [username] -p
Once logged into the mysql, run the following mysql queries to create the mylogtable:
> create database mylog;
> user mylog;
> create table mylogtable( id INT NOT NULL AUTO_INCREMENT, dataitem VARCHAR(255) NOT NULL, PRIMARY KEY(id));
Storm Bolt to persists data
Once this is done. We are ready to create a Storm bolt which route the received tuple from KafkaSpout to the MySqlDump which saves it to the database:package com.memeanalytics.storm_kafka_mysql;
import java.util.Map;
import backtype.storm.task.TopologyContext;
import backtype.storm.topology.BasicOutputCollector;
import backtype.storm.topology.OutputFieldsDeclarer;
import backtype.storm.topology.base.BaseBasicBolt;
import backtype.storm.tuple.Tuple;
public class MySqlDumpBolt extends BaseBasicBolt{
private static final long serialVersionUID = 1L;
private MySqlDump mySqlDump;
@Override
public void prepare(Map stormConf, TopologyContext context)
{
mySqlDump=new MySqlDump("localhost", "mylog","root","[username]");
}
public void execute(Tuple input, BasicOutputCollector collector) {
// TODO Auto-generated method stub
mySqlDump.persist(input);
//System.out.println(input);
}
public void declareOutputFields(OutputFieldsDeclarer declarer) {
// TODO Auto-generated method stub
}
@Override
public void cleanup() {
mySqlDump.close();
}
}
As can be seen above, the bolt opens the mysql database connection in its prepare() method and close the mysql database connection in its cleanup() method. it also persists the tuple received into mysql datatable in its execute() method.
Submit Topology in main()
Now we can implement the main() method to submit a topology. For simplicity, we only uses local cluster:package com.memeanalytics.storm_kafka_mysql;
import backtype.storm.Config;
import backtype.storm.LocalCluster;
import backtype.storm.generated.AlreadyAliveException;
import backtype.storm.generated.InvalidTopologyException;
import backtype.storm.spout.SchemeAsMultiScheme;
import backtype.storm.topology.TopologyBuilder;
import storm.kafka.KafkaSpout;
import storm.kafka.SpoutConfig;
import storm.kafka.StringScheme;
import storm.kafka.ZkHosts;
public class App
{
public static void main( String[] args ) throws AlreadyAliveException, InvalidTopologyException
{
ZkHosts zkHosts=new ZkHosts("192.168.2.4:2181");
String topic="test-topic";
String consumer_group_id="id7";
SpoutConfig kafkaConfig=new SpoutConfig(zkHosts, topic, "", consumer_group_id);
kafkaConfig.forceFromStart=true;
kafkaConfig.scheme=new SchemeAsMultiScheme(new StringScheme());
KafkaSpout kafkaSpout=new KafkaSpout(kafkaConfig);
TopologyBuilder builder=new TopologyBuilder();
builder.setSpout("KafkaSpout", kafkaSpout);
builder.setBolt("MySqlBolt", new MySqlDumpBolt()).globalGrouping("KafkaSpout");
LocalCluster cluster=new LocalCluster();
Config config=new Config();
cluster.submitTopology("MySqlDemoTopology", config, builder.createTopology());
try{
Thread.sleep(10000);
}catch(InterruptedException ex)
{
ex.printStackTrace();
}
cluster.killTopology("MySqlDemoTopology");
cluster.shutdown();
}
}
Execute Maven project
Now navigate to the project root folder and run the following command:> mvn exec:java -Dmain.class=com.memeanalytics.storm_kafka_mysql.App
Note that if you change the ${main.class} to com.memeanalytics.storm_kafka_mysql.App, that you don't need to add the -Dmain.class arguments in the above command.
Write and test a simple Distributed RPC with Storm Trident
Distributed RPC can be used to query results from a Trident topology running in a storm cluster in real-time. This post shows how to use DRPC to query the accumulated country count in real-time on a storm Trident topology implemented in http://czcodezone.blogspot.sg/2014/11/write-and-test-trident-non.html).
The source codes of the project can be downloaded from the following link:
https://dl.dropboxusercontent.com/u/113201788/storm/trident-drpc-test.tar.gz
The main difference between the source codes in this post and that in http://czcodezone.blogspot.sg/2014/11/write-and-test-trident-non.html is in the main class, for which the source codes are show below:
As shown in the source codes,the Trident topology now has a second stream which is for DRPC, it queries the TridentState object which store accumulated count (from the start of the program) of country frequency in a MemoryMapState object (which is a memory-based map object). The DRPCStream has a transaction id "Count", which can be used by the DRPClient.execute() or LocalDRPC.execute(), to query the country accumulated via the "args" (which in this case is "China,Russia,USA", the TridentComps.CountrySplit is a BaseFunction trident operation object which split the "args" tuple ["China,Russia,USA"] into 3 tuples:
["China.Russia.USA", "China"]
["China.Russia.USA", "Russia"]
["China.Russia.USA", "USA"]
The stateQuery() method query the MemoryMapState object for the accumulated count of "Country" field stored.
By the way, for remote DRPC, the user needs to add the following lines to the storm.yaml in the STORM_HOME/conf folder:
drpc.servers:
- "192.168.2.4"
where 192.168.2.4 is the drpc server's hostname (may be the same machine as the master node in Storm). and then in the terminal, run the following command to start the drpc server:
> $STORM_HOME/bin/storm drpc
The source codes of the project can be downloaded from the following link:
https://dl.dropboxusercontent.com/u/113201788/storm/trident-drpc-test.tar.gz
The main difference between the source codes in this post and that in http://czcodezone.blogspot.sg/2014/11/write-and-test-trident-non.html is in the main class, for which the source codes are show below:
package com.memeanalytics.trident_drpc_test;
import storm.trident.TridentState;
import storm.trident.TridentTopology;
import storm.trident.operation.builtin.Count;
import storm.trident.operation.builtin.FilterNull;
import storm.trident.operation.builtin.MapGet;
import storm.trident.testing.MemoryMapState;
import backtype.storm.Config;
import backtype.storm.LocalCluster;
import backtype.storm.LocalDRPC;
import backtype.storm.StormSubmitter;
import backtype.storm.generated.AlreadyAliveException;
import backtype.storm.generated.InvalidTopologyException;
import backtype.storm.generated.StormTopology;
import backtype.storm.tuple.Fields;
import backtype.storm.utils.DRPCClient;
public class App
{
public static void main( String[] args ) throws Exception
{
Config config=new Config();
config.setMaxSpoutPending(20);
if(args.length==0)
{
LocalDRPC drpc=new LocalDRPC();
LocalCluster cluster=new LocalCluster();
cluster.submitTopology("DRPCTridentDemo", config, buildTopology(drpc));
try{
Thread.sleep(2000);
}catch(InterruptedException ex)
{
ex.printStackTrace();
}
for(int i=0; i < 10; ++i)
{
System.out.println(drpc.execute("Count", "China,Russia,USA"));
try{
Thread.sleep(1000);
}catch(InterruptedException ex)
{
ex.printStackTrace();
}
}
cluster.killTopology("DRPCTridentDemo");
cluster.shutdown();
}
else
{
config.setNumWorkers(3);
try{
StormSubmitter.submitTopology(args[0], config, buildTopology(null));
}catch(AlreadyAliveException ex)
{
ex.printStackTrace();
}catch(InvalidTopologyException ex)
{
ex.printStackTrace();
}
try{
Thread.sleep(2000);
}catch(InterruptedException ex)
{
ex.printStackTrace();
}
DRPCClient client=new DRPCClient("RRPC-Server",1234);
System.out.print(client.execute("Count", "China,Russia,USA"));
}
}
private static StormTopology buildTopology(LocalDRPC drpc)
{
TridentTopology topology=new TridentTopology();
RandomWordSpout spout=new RandomWordSpout(10);
TridentState countryCount = topology.newStream("spout1", spout).shuffle().each(new Fields("Country", "Rank"), new TridentComps.CountryFilter()).groupBy(new Fields("Country")).persistentAggregate(new MemoryMapState.Factory(), new Fields("Country"), new Count(), new Fields("count")).parallelismHint(2);
try{
Thread.sleep(2000);
}catch(InterruptedException ex)
{
ex.printStackTrace();
}
topology.newDRPCStream("Count", drpc).each(new Fields("args"), new TridentComps.CountrySplit(), new Fields("Country")).stateQuery(countryCount, new Fields("Country"), new MapGet(), new Fields("count")).each(new Fields("count"), new FilterNull());
return topology.build();
}
}
As shown in the source codes,the Trident topology now has a second stream which is for DRPC, it queries the TridentState object which store accumulated count (from the start of the program) of country frequency in a MemoryMapState object (which is a memory-based map object). The DRPCStream has a transaction id "Count", which can be used by the DRPClient.execute() or LocalDRPC.execute(), to query the country accumulated via the "args" (which in this case is "China,Russia,USA", the TridentComps.CountrySplit is a BaseFunction trident operation object which split the "args" tuple ["China,Russia,USA"] into 3 tuples:
["China.Russia.USA", "China"]
["China.Russia.USA", "Russia"]
["China.Russia.USA", "USA"]
The stateQuery() method query the MemoryMapState object for the accumulated count of "Country" field stored.
By the way, for remote DRPC, the user needs to add the following lines to the storm.yaml in the STORM_HOME/conf folder:
drpc.servers:
- "192.168.2.4"
where 192.168.2.4 is the drpc server's hostname (may be the same machine as the master node in Storm). and then in the terminal, run the following command to start the drpc server:
> $STORM_HOME/bin/storm drpc
Sunday, November 23, 2014
Write and test a Trident non transactional topology in Storm
Trident provides high-level abstraction data model for Storm with the concepts of base function, filter, projection, aggregate, grouping, etc. Though it adds overhead to storm but it makes Storm easier to implement as well as provides support such as at least oncely processing or exactly oncely processing.
The trident non transactional topology to be implemented is extremely simple, a dummy spout (derived from IBatchSpout) emits batch (size:10) of tuples having the form of ["{CountryName}", "{Rank}"]. The tuples emitted contains illegal country names which needs to be filtered away. The tuples having the same country in partitions belonging to the same batch is groupped together. Then the frequency of a particular country appearing in the same batch is counted and printed out.
The source codes of the project can be downloaded from the link below:
https://dl.dropboxusercontent.com/u/113201788/storm/trident-test.tar.gz
To start, create a Maven project (in my case, with groupId="memeanalytics", artifactId="trident-test"), and modify the pom.xml file as shown below:
The pom specifies where to download storm as well as maven plugins for building, executing (exec-maven-plugin), and packaging (maven-assembly-plugin) the java project. Now lets create the spout which emits tuples in batch:
Now we will create a set of Trident operations including CountryFilter (which filters away tuples containing non-country name), Print (which prints values of the count of tuples containing particular country in a batch). The code is as shown below:
Now we are ready to implement the main class:
The main() method is quite straightforward, if there is arguments in the command line, then the project should be packaged into jar and submitted into a storm cluster, otherwise run a local storm cluster and submit the topology there to run. To run locally, navigate to the project root folder and run the following command:
> mvn compile exec:java -Dmain.class=com.memeanalytics.trident_test.App
To run in a storm cluster, make sure the zookeeper cluster and storm cluster is running (following instructions at this link: http://czcodezone.blogspot.sg/2014/11/setup-storm-in-cluster.html), run the following command:
> mvn clean install
After that a trident-test-0.0.1-SNAPSHOT-jar-with-dependencies.jar will be created in the "target" folder under the project root folder.
Now upload the jar by running the following command:
> $STORM_HOME/bin/storm jar [projectRootFolder]/target/trident-test-0.0.1-SNAPSHOT-jar-with-dependencies.jar com.memeanalytics.trident.trident_test.App TridentDemo
The trident non transactional topology to be implemented is extremely simple, a dummy spout (derived from IBatchSpout) emits batch (size:10) of tuples having the form of ["{CountryName}", "{Rank}"]. The tuples emitted contains illegal country names which needs to be filtered away. The tuples having the same country in partitions belonging to the same batch is groupped together. Then the frequency of a particular country appearing in the same batch is counted and printed out.
The source codes of the project can be downloaded from the link below:
https://dl.dropboxusercontent.com/u/113201788/storm/trident-test.tar.gz
To start, create a Maven project (in my case, with groupId="memeanalytics", artifactId="trident-test"), and modify the pom.xml file as shown below:
<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>com.memeanalytics</groupId>
<artifactId>trident-test</artifactId>
<version>0.0.1-SNAPSHOT</version>
<packaging>jar</packaging>
<name>trident-test</name>
<url>http://maven.apache.org</url>
<properties>
<project.build.sourceEncoding>UTF-8</project.build.sourceEncoding>
</properties>
<repositories>
<repository>
<id>clojars</id>
<url>http://clojars.org/repo</url>
</repository>
</repositories>
<dependencies>
<dependency>
<groupId>junit</groupId>
<artifactId>junit</artifactId>
<version>3.8.1</version>
<scope>test</scope>
</dependency>
<dependency>
<groupId>storm</groupId>
<artifactId>storm</artifactId>
<version>0.9.0.1</version>
<scope>provided</scope>
</dependency>
</dependencies>
<build>
<plugins>
<plugin>
<groupId>org.codehaus.mojo</groupId>
<artifactId>exec-maven-plugin</artifactId>
<version>1.2.1</version>
<executions>
<execution>
<goals>
<goal>exec</goal>
</goals>
</execution>
</executions>
<configuration>
<includeProjectDependencies>true</includeProjectDependencies>
<includePluginDependencies>false</includePluginDependencies>
<executable>java</executable>
<classpathScope>compile</classpathScope>
<mainClass>${main.class}</mainClass>
</configuration>
</plugin>
<plugin>
<artifactId>maven-assembly-plugin</artifactId>
<version>2.2.1</version>
<configuration>
<descriptorRefs>
<descriptorRef>jar-with-dependencies</descriptorRef>
</descriptorRefs>
<archive>
<manifest>
<mainClass></mainClass>
</manifest>
</archive>
</configuration>
<executions>
<execution>
<id>make-assembly</id>
<phase>package</phase>
<goals>
<goal>single</goal>
</goals>
</execution>
</executions>
</plugin>
</plugins>
</build>
</project>
The pom specifies where to download storm as well as maven plugins for building, executing (exec-maven-plugin), and packaging (maven-assembly-plugin) the java project. Now lets create the spout which emits tuples in batch:
package com.memeanalytics.trident_test;
import java.util.ArrayList;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
import java.util.Random;
import backtype.storm.task.TopologyContext;
import backtype.storm.tuple.Fields;
import backtype.storm.tuple.Values;
import storm.trident.operation.TridentCollector;
import storm.trident.spout.IBatchSpout;
public class RandomWordSpout implements IBatchSpout {
private static final long serialVersionUID = 1L;
private static String[] countries=new String[]{
"China",
"USA",
"Ruassia",
"UK",
"France",
"Rubbish",
"Garbage"
};
private static Integer[] ranks=new Integer[]{
1,
2,
3,
4,
5
};
private Map<Long, List<List<Object>>> dataStore=new HashMap<Long, List<List<Object>>>();
private int batchSize;
public RandomWordSpout(int batchSize)
{
this.batchSize = batchSize;
}
public void open(Map conf, TopologyContext context) {
// TODO Auto-generated method stub
}
public void emitBatch(long batchId, TridentCollector collector) {
// TODO Auto-generated method stub
List<List<Object>> batch=dataStore.get(batchId);
if(batch == null)
{
final Random rand=new Random();
batch=new ArrayList<List<Object>>();
for(int i=0; i < batchSize; ++i)
{
batch.add(new Values(
countries[rand.nextInt(countries.length)],
ranks[rand.nextInt(ranks.length)]
));
}
dataStore.put(batchId, batch);
}
for(List<Object> tuple : batch)
{
collector.emit(tuple);
}
}
public void ack(long batchId) {
// TODO Auto-generated method stub
dataStore.remove(batchId);
}
public void close() {
// TODO Auto-generated method stub
}
public Map getComponentConfiguration() {
return null;
}
public Fields getOutputFields() {
return new Fields("Country","Rank");
}
}
The spout basically create batch based on the batchId and emits the tuples in that batch. When acknowledgement is received, the acknowledge batch having the batchId is then removed from the spout. Note that the spout will emit tuples containing non-country name such as "Rubbish" and "Garbage"Now we will create a set of Trident operations including CountryFilter (which filters away tuples containing non-country name), Print (which prints values of the count of tuples containing particular country in a batch). The code is as shown below:
package com.memeanalytics.trident_test;
import backtype.storm.tuple.Values;
import storm.trident.operation.BaseFilter;
import storm.trident.operation.BaseFunction;
import storm.trident.operation.TridentCollector;
import storm.trident.tuple.TridentTuple;
public class TridentComps {
public static class CountryFilter extends BaseFilter{
private static final long serialVersionUID = 1L;
public boolean isKeep(TridentTuple tuple) {
// TODO Auto-generated method stub
String country_candidate = tuple.getString(0);
return !country_candidate.equals("Garbage") && !country_candidate.equals("Rubbish");
}
}
public static class CountrySplit extends BaseFunction{
private static final long serialVersionUID = 1L;
public void execute(TridentTuple tuple, TridentCollector collector) {
// TODO Auto-generated method stub
String country_comps=tuple.getString(0);
for(String country_candidate : country_comps.split("\\s"))
{
collector.emit(new Values(country_candidate.trim()));
}
}
}
public static class Print extends BaseFilter{
private static final long serialVersionUID = 1L;
public boolean isKeep(TridentTuple tuple) {
// TODO Auto-generated method stub
System.out.println(tuple);
return true;
}
}
}
Now we are ready to implement the main class:
package com.memeanalytics.trident_test;
import storm.trident.TridentTopology;
import storm.trident.operation.builtin.Count;
import backtype.storm.Config;
import backtype.storm.LocalCluster;
import backtype.storm.StormSubmitter;
import backtype.storm.generated.AlreadyAliveException;
import backtype.storm.generated.InvalidTopologyException;
import backtype.storm.generated.StormTopology;
import backtype.storm.tuple.Fields;
public class App
{
public static void main( String[] args ) throws Exception
{
Config config=new Config();
config.setMaxPendingSpouts(20);
if(args.length==0)
{
LocalCluster cluster=new LocalCluster();
cluster.submitTopology("TridentDemo", config, buildTopology());
try{
Thread.sleep(10000);
}catch(InterruptedException ex)
{
ex.printStackTrace();
}
cluster.killTopology("TridentDemo");
cluster.shutdown();
}
else
{
config.setNumWorkers(3);
try{
StormSubmitter.submitTopology(args[0], config, buildTopology());
}catch(AlreadyAliveException ex)
{
ex.printStackTrace();
}catch(InvalidTopologyException ex)
{
ex.printStackTrace();
}
}
}
private static StormTopology buildTopology()
{
RandomWordSpout spout=new RandomWordSpout(10);
TridentTopology topology=new TridentTopology();
topology.newStream("TridentTxId", spout).shuffle().each(new Fields("Country"), new TridentComps.CountryFilter()).groupBy(new Fields("Country")).aggregate(new Fields("Country"), new Count(), new Fields("Count")).each(new Fields("Count"), new TridentComps.Print()).parallelismHint(2);
return topology.build();
}
}
The static method buildTopology() creates a Trident non transactional topology, which uses the spout created as data source, the tuples are then filtered by the CountryFilter, and the groupped by the "Country" field value within each batch, a frequency count is then generated via the aggregate method. Finally it is then printed out into the console. (Note that the "aggregate(new Fields("Country"), new Count(), new Fields("Count")" will lead the TridentComps.Print to print out the "Count" value, if you want to print the country as well, then change it to "aggregate(new Fields("Country"), new Count(), new Fields("Country", "Count")")The main() method is quite straightforward, if there is arguments in the command line, then the project should be packaged into jar and submitted into a storm cluster, otherwise run a local storm cluster and submit the topology there to run. To run locally, navigate to the project root folder and run the following command:
> mvn compile exec:java -Dmain.class=com.memeanalytics.trident_test.App
To run in a storm cluster, make sure the zookeeper cluster and storm cluster is running (following instructions at this link: http://czcodezone.blogspot.sg/2014/11/setup-storm-in-cluster.html), run the following command:
> mvn clean install
After that a trident-test-0.0.1-SNAPSHOT-jar-with-dependencies.jar will be created in the "target" folder under the project root folder.
Now upload the jar by running the following command:
> $STORM_HOME/bin/storm jar [projectRootFolder]/target/trident-test-0.0.1-SNAPSHOT-jar-with-dependencies.jar com.memeanalytics.trident.trident_test.App TridentDemo
Subscribe to:
Posts (Atom)