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 Kafka. Show all posts
Showing posts with label Kafka. 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>
Monday, November 24, 2014
Integration of Kafka-Trident-MySQL
This post shows a most basic example in which user can integrate Kafka, Trident (on top of Storm) and MySQL. The example uses a Kafka producer which randomly produce messages to Kafka brokers (a random list of country names), a TransactionalTridentKafkaSpout is used pull data from Kafka messaging system and emits the tuples (containing the field "str" which is the country names from the Kafka producer) to a Trident operation that serialize the received messages into the mysql database.
Some ZooKeeper and Kafka settings need to be explained before we proceed. the source codes developed here assumes that the ZooKeepers runs on the following nodes:
192.168.2.2:2181
192.168.2.4:2181
and also assumes that the Kafka brokers runs at the following hostname:port:
192.168.2.2:9092
192.168.2.4:9092
The Kafka producer can be downloaded from the following link:
https://dl.dropboxusercontent.com/u/113201788/storm/kakfa-producer-for-trident.tar.gz
Basically the Kafka producer emits a random list of country names as messages sent to the Kafka brokers. The tutorial on how to implement Kafka producer can be found at:
Next we need to create Maven project (e.g. with groupId="com.memeanalytics" and artifactId="kafka-trident-consumer") which will consumes the Kafka message in a Trident topology and serialize it to mysql database. The source codes of the project can be downloaded from:
https://dl.dropboxusercontent.com/u/113201788/storm/kafka-trident-consumer.zip
Below we will explain how to prepare the pom.xml file, implement the storm operation for mysql serialization, as well as configuration of Trident topology which can consume messages from Kafka brokers.
Next add in the storm-kafka-0.8-plus dependency in the dependencies section (for TransactionalTridentKafkaSpout):
Next add in the storm-core dependency in the dependencies section (for storm):
Next add in the mysql-connector-java dependency in the dependencies section (for mysql):
Next add in the maven-assembly-plugin in the build/plugins section (for packaging the Maven project as jar for submitting to Storm cluster):
Next add in the exec-maven-plugin in the build/plugins section (for executing the Maven project):
This completes the pom.xml configuration. Next we will implements the Trident operation which serializes the TridentTuple data into mysql database.
In the above implementation, the mysqlSerializer member variable is responsible for actually storing the data in mysql database. it opens the mysql connection (in its constructor) in prepare() method and closes the mysql connection in cleanup() method. the data serialization happens in the execute() method. Below is the implementation of the class of the variable:
For demo purpose, its uses LocalCluster to submit and run the Trident topology and does not implement a DRPCStream which can be used to queried result. The main change is the TractionalTridentKafkaSpout which takes in a TridentKafkaConfig object as parameter in its constructor. The TransactionalTridentKafkaSpout emits tuples which is serialized by the TridentUtil.MySqlPersist that is the mysql serialization Trident operation.
Once done, we can compile and exec the Maven project by navigating to its root folder and run the following command:
> mvn compile exec:java
Some ZooKeeper and Kafka settings need to be explained before we proceed. the source codes developed here assumes that the ZooKeepers runs on the following nodes:
192.168.2.2:2181
192.168.2.4:2181
and also assumes that the Kafka brokers runs at the following hostname:port:
192.168.2.2:9092
192.168.2.4:9092
The Kafka producer can be downloaded from the following link:
https://dl.dropboxusercontent.com/u/113201788/storm/kakfa-producer-for-trident.tar.gz
Basically the Kafka producer emits a random list of country names as messages sent to the Kafka brokers. The tutorial on how to implement Kafka producer can be found at:
Next we need to create Maven project (e.g. with groupId="com.memeanalytics" and artifactId="kafka-trident-consumer") which will consumes the Kafka message in a Trident topology and serialize it to mysql database. The source codes of the project can be downloaded from:
https://dl.dropboxusercontent.com/u/113201788/storm/kafka-trident-consumer.zip
Below we will explain how to prepare the pom.xml file, implement the storm operation for mysql serialization, as well as configuration of Trident topology which can consume messages from Kafka brokers.
Prepare pom.xml
In the pom.xml, first add in the clojars repository in the repositories section:<repositories> <repository> <id>clojars</id> <url>http://clojars.org/repo</url> </repository> </repositories>
Next add in the storm-kafka-0.8-plus dependency in the dependencies section (for TransactionalTridentKafkaSpout):
<dependency> <groupId>net.wurstmeister.storm</groupId> <artifactId>storm-kafka-0.8-plus</artifactId> <version>0.4.0</version> </dependency>
Next add in the storm-core dependency in the dependencies section (for storm):
<dependency> <groupId>storm</groupId> <artifactId>storm-core</artifactId> <version>0.9.0.1</version> </dependency>
Next add in the mysql-connector-java dependency in the dependencies section (for mysql):
<dependency> <groupId>mysql</groupId> <artifactId>mysql-connector-java</artifactId> <version>5.1.6</version> </dependency>
Next add in the maven-assembly-plugin in the build/plugins section (for packaging the Maven project as 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>
Next add in the exec-maven-plugin in the build/plugins section (for executing 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.kafka_trident_consumer.App</mainClass> </configuration> </plugin>
This completes the pom.xml configuration. Next we will implements the Trident operation which serializes the TridentTuple data into mysql database.
Trident operation for mysql serialization
The Trident operation which serializes the TridentTuple data is a BaseFilter object from Trident that has the following implementation:package com.memeanalytics.kafka_trident_consumer;
import java.util.Map;
import storm.trident.operation.BaseFilter;
import storm.trident.operation.TridentOperationContext;
import storm.trident.tuple.TridentTuple;
public class TridentUtils {
public static class MySqlPersist extends BaseFilter{
private static final long serialVersionUID = 1L;
private MySqlDump mysqlSerializer=null;
public boolean isKeep(TridentTuple tuple) {
String country=tuple.getString(0);
mysqlSerializer.store(country);
System.out.println("Country: "+country);
return true;
}
@Override
public void prepare(Map conf, TridentOperationContext context) {
mysqlSerializer=new MySqlDump("localhost","mylog", "root", "[username]");
}
@Override
public void cleanup() {
mysqlSerializer.close();
}
}
}
In the above implementation, the mysqlSerializer member variable is responsible for actually storing the data in mysql database. it opens the mysql connection (in its constructor) in prepare() method and closes the mysql connection in cleanup() method. the data serialization happens in the execute() method. Below is the implementation of the class of the variable:
package com.memeanalytics.kafka_trident_consumer;
import java.sql.Connection;
import java.sql.SQLException;
import java.sql.PreparedStatement;
public class MySqlDump {
private String ip;
private String database;
private String username;
private String password;
private Connection conn;
public MySqlDump(String ip, String database, String username, String password)
{
this.ip = ip;
this.database=database;
this.username=username;
this.password=password;
conn=MySqlConnectionGenerator.open(ip, database, username, password);
}
public void store(String dataitem)
{
if(conn==null) return;
PreparedStatement statement = null;
try{
statement = conn.prepareStatement("insert into mylogtable (id, dataitem) values (default, ?)");
statement.setString(1, dataitem);
statement.executeUpdate();
}catch(Exception ex)
{
ex.printStackTrace();
}finally
{
if(statement != null)
{
try {
statement.close();
} catch (SQLException e) {
// TODO Auto-generated catch block
e.printStackTrace();
}
}
}
}
public void close()
{
if(conn==null) return;
try {
conn.close();
} catch (SQLException e) {
// TODO Auto-generated catch block
e.printStackTrace();
}
}
}
package com.memeanalytics.kafka_trident_consumer;
import java.sql.Connection;
import java.sql.DriverManager;
public class MySqlConnectionGenerator {
public static Connection open(String ip, String database, String username, String password)
{
Connection conn = null;
try{
Class.forName("com.mysql.jdbc.Driver");
conn=DriverManager.getConnection("jdbc:mysql://"+ip+"/"+database+"?user="+username+"&password="+password);
}catch(Exception ex)
{
ex.printStackTrace();
}
return conn;
}
}
Trident topology implementation
Once this is completed. we are ready to implement the Trident topology which consumes data from Kafka and saves it to mysql. this is implemented in the main class:package com.memeanalytics.kafka_trident_consumer;
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.tuple.Fields;
import storm.kafka.SpoutConfig;
import storm.kafka.StringScheme;
import storm.kafka.ZkHosts;
import storm.kafka.trident.TransactionalTridentKafkaSpout;
import storm.kafka.trident.TridentKafkaConfig;
import storm.trident.TridentTopology;
public class App
{
public static void main( String[] args ) throws AlreadyAliveException, InvalidTopologyException
{
ZkHosts zkHosts=new ZkHosts("192.168.2.4:2181");
String topic="country-topic";
String consumer_group_id="storm";
TridentKafkaConfig kafkaConfig=new TridentKafkaConfig(zkHosts, topic, consumer_group_id);
kafkaConfig.scheme=new SchemeAsMultiScheme(new StringScheme());
kafkaConfig.forceFromStart=true;
TransactionalTridentKafkaSpout spout=new TransactionalTridentKafkaSpout(kafkaConfig);
TridentTopology topology=new TridentTopology();
topology.newStream("spout", spout).shuffle().each(new Fields("str"), new TridentUtils.MySqlPersist());
LocalCluster cluster=new LocalCluster();
Config config=new Config();
cluster.submitTopology("KafkaTridentMysqlDemo", config, topology.build());
try{
Thread.sleep(10000);
}catch(InterruptedException ex)
{
ex.printStackTrace();
}
cluster.killTopology("KafkaTridentMysqlDemo");
cluster.shutdown();
}
}
For demo purpose, its uses LocalCluster to submit and run the Trident topology and does not implement a DRPCStream which can be used to queried result. The main change is the TractionalTridentKafkaSpout which takes in a TridentKafkaConfig object as parameter in its constructor. The TransactionalTridentKafkaSpout emits tuples which is serialized by the TridentUtil.MySqlPersist that is the mysql serialization Trident operation.
Once done, we can compile and exec the Maven project by navigating to its root folder and run the following command:
> 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.
Saturday, November 22, 2014
Integrate Kafka with Storm
Create a Maven project (for example, with groupId = "com.memeanalytics" and artifactId = "kafka-consumer-storm"). and then modify the pom.xml file to include the following repository:
<repositories>
...
<repository>
<id>github-releases</id>
<url>http://oss.sonatype.org/content/repositories/github-releases/</url>
</repository>
<repository>
<id>clojars</id>
<url>http://clojars.org/repo</url>
</repository>
</repositories>
In the dependencies section of the pom.xml, include the storm-kafka-0.8-plus (which contains the KafkaSpout which is a spout in the storm cluster that act as consumer to Kafka) and storm-core:
<dependencies>
...
<dependency>
<groupId>net.wurstmeister.storm</groupId>
<artifactId>storm-kafka-0.8-plus</artifactId>
<version>0.4.0</version>
</dependency>
<dependency>
<groupId>storm</groupId>
<artifactId>storm-core<artifactId>
<version>0.9.0.1</version>
</dependency>
<!--Utility dependencies-->
<dependency>
<groupId>commons-collections</groupId>
<artifactId>commons-collections</artifactId>
<version>3.2.1</version>
</dependency>
<dependency>
<groupId>com.google.guava</groupId>
<artifactId>guava</artifactId>
<version>15.0</version>
<dependency>
</dependencies>
In the build/plugins section of the pom.xml, include the maven-assembly-plugin [version:2.2.1] (contains maven plugin for jar packaging of the storm topology project, which can be submitted into storm cluster, instruction can be found at this link: http://czcodezone.blogspot.sg/2014/11/maven-pom-configuration-for-maven-build.html) and exec-maven-plugin [version:1.2.1] (contains maven plugin to execute java program, instruction can be found at this link: http://czcodezone.blogspot.sg/2014/11/maven-add-plugin-to-pom-configuration.html).
Now we are ready to create storm topology which pull data from Kafka messaging system. The topology will be very simple, it will use a KafkaSpout which reads message from Kafka broker and emit a tuple consisting the message to a very simple bolt which prints the content of the tuple out. The simple bolt, PrinterBolt, has the following implementation:
Below is the implementation of the main class com.memeanalytics.kafka_consumer_storm.App, which creates a KafkaSpout which emits tuple to the PrinterBolt above (LocalCluster is used but will be change to StormSubmitter in production environment):
As can be seen, the main component consists of create a SpoutConfig which is configured for the KafkaSpout object. The SpoutConfig in this case specifies the following information:
--zookeeper 192.168.2.4:2181
--topic test-topic
--consumer-group-id: id7
--from-beginning
--string-scheme
The complete source codes can be downloaded from the link below:
https://dl.dropboxusercontent.com/u/113201788/storm/kafka-consumer-storm.zip
Now before executing the project, we must have the kafka cluster running (follow instructions from this link: http://czcodezone.blogspot.sg/2014/11/setup-kafka-in-cluster.html) and a Kafka producer running (following instructions from this link: http://czcodezone.blogspot.sg/2014/11/write-and-test-simple-kafka-producer.html). Now compile and run the kafka-consumer-storm project by running the following command from the project's root directory in the terminal :
> mvn clean compile exec:java -Dmain.class=com.memeanalytics.kafka_consumer_storm.App
You will now see that words produced by the Kafka producer being printed by the storm's PrinterBolt.
<repositories>
...
<repository>
<id>github-releases</id>
<url>http://oss.sonatype.org/content/repositories/github-releases/</url>
</repository>
<repository>
<id>clojars</id>
<url>http://clojars.org/repo</url>
</repository>
</repositories>
In the dependencies section of the pom.xml, include the storm-kafka-0.8-plus (which contains the KafkaSpout which is a spout in the storm cluster that act as consumer to Kafka) and storm-core:
<dependencies>
...
<dependency>
<groupId>net.wurstmeister.storm</groupId>
<artifactId>storm-kafka-0.8-plus</artifactId>
<version>0.4.0</version>
</dependency>
<dependency>
<groupId>storm</groupId>
<artifactId>storm-core<artifactId>
<version>0.9.0.1</version>
</dependency>
<!--Utility dependencies-->
<dependency>
<groupId>commons-collections</groupId>
<artifactId>commons-collections</artifactId>
<version>3.2.1</version>
</dependency>
<dependency>
<groupId>com.google.guava</groupId>
<artifactId>guava</artifactId>
<version>15.0</version>
<dependency>
</dependencies>
In the build/plugins section of the pom.xml, include the maven-assembly-plugin [version:2.2.1] (contains maven plugin for jar packaging of the storm topology project, which can be submitted into storm cluster, instruction can be found at this link: http://czcodezone.blogspot.sg/2014/11/maven-pom-configuration-for-maven-build.html) and exec-maven-plugin [version:1.2.1] (contains maven plugin to execute java program, instruction can be found at this link: http://czcodezone.blogspot.sg/2014/11/maven-add-plugin-to-pom-configuration.html).
Now we are ready to create storm topology which pull data from Kafka messaging system. The topology will be very simple, it will use a KafkaSpout which reads message from Kafka broker and emit a tuple consisting the message to a very simple bolt which prints the content of the tuple out. The simple bolt, PrinterBolt, has the following implementation:
package com.memeanalytics.kafka_consumer_storm;
import org.apache.commons.lang.StringUtils;
import backtype.storm.topology.BasicOutputCollector;
import backtype.storm.topology.OutputFieldsDeclarer;
import backtype.storm.topology.base.BaseBasicBolt;
import backtype.storm.tuple.Tuple;
public class PrinterBolt extends BaseBasicBolt {
private static final long serialVersionUID = 1L;
public void execute(Tuple input, BasicOutputCollector collector) {
// TODO Auto-generated method stub
String word=input.getString(0);
if(StringUtils.isBlank(word))
{
return;
}
System.out.println("Word: "+word);
}
public void declareOutputFields(OutputFieldsDeclarer declarer) {
// TODO Auto-generated method stub
}
}
Below is the implementation of the main class com.memeanalytics.kafka_consumer_storm.App, which creates a KafkaSpout which emits tuple to the PrinterBolt above (LocalCluster is used but will be change to StormSubmitter in production environment):
package com.memeanalytics.kafka_consumer_storm;
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_name="test-topic";
String consumer_group_id="id7";
String zookeeper_root="";
SpoutConfig kafkaConfig=new SpoutConfig(zkHosts, topic_name, zookeeper_root, consumer_group_id);
kafkaConfig.scheme=new SchemeAsMultiScheme(new StringScheme());
kafkaConfig.forceFromStart=true;
TopologyBuilder builder=new TopologyBuilder();
builder.setSpout("KafkaSpout", new KafkaSpout(kafkaConfig), 1);
builder.setBolt("PrinterBolt", new PrinterBolt()).globalGrouping("KafkaSpout");
Config config=new Config();
LocalCluster cluster=new LocalCluster();
cluster.submitTopology("KafkaConsumerTopology", config, builder.createTopology());
try{
Thread.sleep(60000);
}catch(InterruptedException ex)
{
ex.printStackTrace();
}
cluster.killTopology("KafkaConsumerTopology");
cluster.shutdown();
}
}
As can be seen, the main component consists of create a SpoutConfig which is configured for the KafkaSpout object. The SpoutConfig in this case specifies the following information:
--zookeeper 192.168.2.4:2181
--topic test-topic
--consumer-group-id: id7
--from-beginning
--string-scheme
The complete source codes can be downloaded from the link below:
https://dl.dropboxusercontent.com/u/113201788/storm/kafka-consumer-storm.zip
Now before executing the project, we must have the kafka cluster running (follow instructions from this link: http://czcodezone.blogspot.sg/2014/11/setup-kafka-in-cluster.html) and a Kafka producer running (following instructions from this link: http://czcodezone.blogspot.sg/2014/11/write-and-test-simple-kafka-producer.html). Now compile and run the kafka-consumer-storm project by running the following command from the project's root directory in the terminal :
> mvn clean compile exec:java -Dmain.class=com.memeanalytics.kafka_consumer_storm.App
You will now see that words produced by the Kafka producer being printed by the storm's PrinterBolt.
Write and test a simple Kafka producer
First we would need to start a zookeeper cluster
Now create a Maven project in Eclipse or STS (e.g. groupId=com.memeanalytics artifactId=kafka-producer), and change the pom.xml to include the following dependencies and plugins:
<dependencies>
...
<dependency>
<groupId>org.apache.kafka</groupId>
<artifactId>kafka_2.8.0</artifactId>
<version>0.8.1.1</version>
<exclusions>
<exclusion>
<groupId>javax.jms</groupId>
<artifactId>jms</artifactId>
</exclusion>
<exclusion>
<groupId>com.sun.jdmk</groupId>
<artifactId>jmxtools</artifactId>
</exclusion>
<exclusion>
<groupId>com.sun.jmx</groupId>
<artifactId>jmxri</artifactId>
</exclusion>
</exclusions>
</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>com.memeanalytics.kafka_producer.App</mainClass>
</configuration>
</plugin>
</plugins>
</build>
The dependency setting include the kafka pom and the plugin setting include the maven exec plugin.
The implementation of a kafka producer in java is very simple a straight forward, we first create a KafkaConfig object which the following initialization:
metadata.broker.list: 192.168.2.4:9092
serializer.class: kafka.serializer.StringEncoder
request.required.acks: 1
The broker list we only need to specify one broker and the rest will be automatically discovered. Since we are to produce string data to the kafka broker, the StringEncoder is used as the data serializer, We also specify that we like to have acknowleges for request sent.
Next is to create a Producer object that is configured by the KafkaConfig object above:
Producer<String, String> kafkaProducer=new Producer<String, String>(config);
The actual sending of data is performed by the following line:
KeyedMessage<String, String> data=new KeyedMessage<String, String>("test-topic", "MyDataTextMessage");
kafkaProducer.send(data);
The "test-topic" is the name of the topic and the "MyDataTextMessage" is the actual data sent to kafka brokers. After all the data has been sent, the kafkaProducer needs to be closed:
kafkaProducer.close();
You can download the complete java project codes from the link below:
https://dl.dropboxusercontent.com/u/113201788/storm/kafka-producer.zip
To run, navigate the root folder of the Maven project, and run the following command in a terminal:
> mvn compile exec:java
To test, make sure that the kafka cluster is setup and running (follow the instructions given in this link: http://czcodezone.blogspot.sg/2014/11/setup-kafka-in-cluster.html) open another terminal and run a kafka console consumer using the following command:
Now create a Maven project in Eclipse or STS (e.g. groupId=com.memeanalytics artifactId=kafka-producer), and change the pom.xml to include the following dependencies and plugins:
<dependencies>
...
<dependency>
<groupId>org.apache.kafka</groupId>
<artifactId>kafka_2.8.0</artifactId>
<version>0.8.1.1</version>
<exclusions>
<exclusion>
<groupId>javax.jms</groupId>
<artifactId>jms</artifactId>
</exclusion>
<exclusion>
<groupId>com.sun.jdmk</groupId>
<artifactId>jmxtools</artifactId>
</exclusion>
<exclusion>
<groupId>com.sun.jmx</groupId>
<artifactId>jmxri</artifactId>
</exclusion>
</exclusions>
</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>com.memeanalytics.kafka_producer.App</mainClass>
</configuration>
</plugin>
</plugins>
</build>
The dependency setting include the kafka pom and the plugin setting include the maven exec plugin.
The implementation of a kafka producer in java is very simple a straight forward, we first create a KafkaConfig object which the following initialization:
metadata.broker.list: 192.168.2.4:9092
serializer.class: kafka.serializer.StringEncoder
request.required.acks: 1
The broker list we only need to specify one broker and the rest will be automatically discovered. Since we are to produce string data to the kafka broker, the StringEncoder is used as the data serializer, We also specify that we like to have acknowleges for request sent.
Next is to create a Producer object that is configured by the KafkaConfig object above:
Producer<String, String> kafkaProducer=new Producer<String, String>(config);
The actual sending of data is performed by the following line:
KeyedMessage<String, String> data=new KeyedMessage<String, String>("test-topic", "MyDataTextMessage");
kafkaProducer.send(data);
The "test-topic" is the name of the topic and the "MyDataTextMessage" is the actual data sent to kafka brokers. After all the data has been sent, the kafkaProducer needs to be closed:
kafkaProducer.close();
You can download the complete java project codes from the link below:
https://dl.dropboxusercontent.com/u/113201788/storm/kafka-producer.zip
To run, navigate the root folder of the Maven project, and run the following command in a terminal:
> mvn compile exec:java
To test, make sure that the kafka cluster is setup and running (follow the instructions given in this link: http://czcodezone.blogspot.sg/2014/11/setup-kafka-in-cluster.html) open another terminal and run a kafka console consumer using the following command:
> $KAFKA_HOME/bin/kafka-console-consumer.sh --zookeeper 192.168.2.2:2181 --topic test-topic --from-beginning
Setup Kafka in a cluster
To setup Kafka in a cluster, first we must have the zookeeper cluster setup and running (follow this link: http://czcodezone.blogspot.sg/2014/11/setup-zookeeper-in-cluster.html), suppose that the zookeeper cluster consists of the zookeeper servers running at the following hostname:ports:
192.168.2.2:2181
192.168.2.4:2181
As I have only two computers, therefore i will use the same computers (but at different ports) to host the kafka cluster. For this case, the Kafka servers/brokers will be running at the following hostname:ports
192.168.2.2:9092
192.168.2.4:9092
To do this, follow this link (http://czcodezone.blogspot.sg/2014/11/setup-kafka-in-single-machine-running.html) to setup the kafka server. Now navigate to the kafka root folder oif each computer and modify the server.properties in "config" sub folder:
In the server.properties file, search the line "zookeeper.connect" and change it to:
zookeeper.connect=192.168.2.2:2181,192.168.2.4:2181
Then search the line "broker.id" (unique id for each broker node) and change it to "broke.id=1" on computer 192.168.2.2 and to "broker.id=2" on computer 192.168.2.4
Next search the line "host.name" and change it to "host.name=192.168.2.2" on computer 192.168.2.2 and to "host.name=192.168.2.4" on computer 192.168.2.4
Make sure that the line "port=9092" is there and uncommented in the server.properties
Save and close the server.properties. Now start the kafka server on each computer:
At this point, the kafka cluster is set up and running. We can test the cluster by creating a topic named "v-topic":
> bin/kakfa-topics.sh --create --zookeeper 192.168.2.4:2181 --partitions 2 --replication-factor 1 --topic v-topic
Now run the following commands to list the topics in the kafka brokers:
> bin/kafka-topics.sh --zookeeper 192.168.2.4:2181 --list
Now run the following commands to get a description how the topic "v-topic" is partitioned in each broker:
> bin/kafka-topics.sh --describe --zookeeper 192.168.2.4:2181 --topic v-topic
To test the producer and consumer interaction, let's start a consoler producer on the computer 192.168.2.4 by running the following command on that computer's terminal:
Now open a terminal of the other computer 192.168.2.4 and start a console consumer:
Begin to type something in the console producer on 192.168.2.2 terminal and press ENTER, you will see the output displayed in the console consumer on 192.168.2.4 terminal.
Note:
It is also ok to set up multiple Kafka brokers on the same computer. For example, if we want to have two Kafka brokers running at two different ports on computer 192.168.2.2, say:
192.168.2.2:9092
192.168.2.2:9093
Now all that we need to do is to duplicate the server.properties after it is updated, and rename it server1.properties in the same "config" folder (note that name is not important, can be anything that make sense). Now in the server1.properties, modify to have the following settings:
broker.id=3
log.dirs=/var/kafka1-logs
port=9093
Save and close server1.properties (remember to create the folder /var/kafka1-logs with write permission), open two terminal in 192.168.2.2 and run the following command in the first terminal to start a kafka broker at port 9092:
On the second terminal, run the following command to start a second kafka broker at port 9093:
Now you will have two kafka brokers running on 192.168.2.2 on two different ports. To include the second broker for the console producer, change its start command to:
192.168.2.2:2181
192.168.2.4:2181
As I have only two computers, therefore i will use the same computers (but at different ports) to host the kafka cluster. For this case, the Kafka servers/brokers will be running at the following hostname:ports
192.168.2.2:9092
192.168.2.4:9092
To do this, follow this link (http://czcodezone.blogspot.sg/2014/11/setup-kafka-in-single-machine-running.html) to setup the kafka server. Now navigate to the kafka root folder oif each computer and modify the server.properties in "config" sub folder:
> cd $KAFKA_HOME > gedit config/server.properties
In the server.properties file, search the line "zookeeper.connect" and change it to:
zookeeper.connect=192.168.2.2:2181,192.168.2.4:2181
Then search the line "broker.id" (unique id for each broker node) and change it to "broke.id=1" on computer 192.168.2.2 and to "broker.id=2" on computer 192.168.2.4
Next search the line "host.name" and change it to "host.name=192.168.2.2" on computer 192.168.2.2 and to "host.name=192.168.2.4" on computer 192.168.2.4
Make sure that the line "port=9092" is there and uncommented in the server.properties
Save and close the server.properties. Now start the kafka server on each computer:
> cd $KAFKA_HOME> bin/kafka-server-start.sh config/server.properties
At this point, the kafka cluster is set up and running. We can test the cluster by creating a topic named "v-topic":
> bin/kakfa-topics.sh --create --zookeeper 192.168.2.4:2181 --partitions 2 --replication-factor 1 --topic v-topic
Now run the following commands to list the topics in the kafka brokers:
> bin/kafka-topics.sh --zookeeper 192.168.2.4:2181 --list
Now run the following commands to get a description how the topic "v-topic" is partitioned in each broker:
> bin/kafka-topics.sh --describe --zookeeper 192.168.2.4:2181 --topic v-topic
To test the producer and consumer interaction, let's start a consoler producer on the computer 192.168.2.4 by running the following command on that computer's terminal:
> cd $KAFKA_HOME > bin/kafka-console-producer.sh --broker-list 192.168.2.2:9092,192.168.2.4:9092 --topic v-topic
Now open a terminal of the other computer 192.168.2.4 and start a console consumer:
> cd $KAFKA_HOME > bin/kafka-console-consumer.sh --zookeeper 192.168.2.4:2181 --topic v-topic --from-beginning
Begin to type something in the console producer on 192.168.2.2 terminal and press ENTER, you will see the output displayed in the console consumer on 192.168.2.4 terminal.
Note:
It is also ok to set up multiple Kafka brokers on the same computer. For example, if we want to have two Kafka brokers running at two different ports on computer 192.168.2.2, say:
192.168.2.2:9092
192.168.2.2:9093
Now all that we need to do is to duplicate the server.properties after it is updated, and rename it server1.properties in the same "config" folder (note that name is not important, can be anything that make sense). Now in the server1.properties, modify to have the following settings:
broker.id=3
log.dirs=/var/kafka1-logs
port=9093
Save and close server1.properties (remember to create the folder /var/kafka1-logs with write permission), open two terminal in 192.168.2.2 and run the following command in the first terminal to start a kafka broker at port 9092:
> $KAFKA_HOME/bin/kafka-server-start.sh $KAFKA_HOME/config/server.properties
On the second terminal, run the following command to start a second kafka broker at port 9093:
> $KAFKA_HOME/bin/kafka-server-start.sh $KAFKA_HOME/config/server1.properties
Now you will have two kafka brokers running on 192.168.2.2 on two different ports. To include the second broker for the console producer, change its start command to:
> $KAKFA_HOME/bin/kafka-console-producer.sh --broker-list 192.168.2.2:9092,192.168.2.2:9093,192.168.2.4:9092 --topic v-topic
Setup Kafka in a single machine running Ubuntu 14.04 LTS
Kafka is a messaging system that can acts as a buffer and feeder for messages processed by Storm spouts. It can also be used as a output buffer for Storm bolts. This post shows how to setup and test Kafka on a single machine running Ubuntu.
Firstly download the kafka 0.8.1.1 from the link below:
https://www.apache.org/dyn/closer.cgi?path=/kafka/0.8.1.1/kafka_2.8.0-0.8.1.1.tgz
Next "tar -xvzf" the kafka_2.8.0-0.8.1.1.tgz file and move it to a destination folder (say, /Documents/Works/Kafka folder under the user root directory):
Now go back to the user root folder and open the .bashrc file for editing:
In the .bashrc file, add the following line to the end:
export KAFKA_HOME=$HOME/Documents/Works/Kakfa/kafka_2.8.0-0.8.1.1
Save and close the .bashrc and run "source .bashrc" to update the environment variables. Now navigate to the kafka home folder and edit the server.properties in its sub-directory "config":
In the server.properties file, search the line "zookeeper.connect" and change it to the following:
zookeeper.connect=192.168.2.2:2181,192.168.2.4:2181
search the line "log.dirs" and change it to the following:
log.dirs=/var/kafka-logs
Save and close the server.properties file (192.168.2.2 and 192.168.2.4 are the zookeeper nodes). Next we go and create the folder /var/kafka-logs (which will store the topics and partitions data for kafka) with write permissions:
> sudo mkdir /var/kafka-logs
> sudo chmod -R 777 /var/kafka-logs
Now set up and run the zookeeper cluster by following instructions in the link http://czcodezone.blogspot.sg/2014/11/setup-zookeeper-in-cluster.html. Once this is done, we are ready to start the kafka messaging system by running the following commands:
To start testing kafka setup, Ctrl+Alt+T to open a new terminal and run the following command to create a topic "verification-topic" (a topic is a named entity in kafka which contain one or more partitions which are message queues that can run in parallel and serialize to individual folder in /var/kafka-log folder):
The above command creates a topic named "verification-topic" which contains 1 partition (and with no replication)
Now we can check the list of topics in kafka by running the following command:
> bin/kafka-topics.sh --zookeeper 192.168.2.2:2181 --list
To test the producer and consumer interaction in kafka, fire up the console producer by running
> bin/kafka-console-producer.sh --broker-list localhost:9092 --topic verification-topic
9092 is the default port for a kafka broker node (which is localhost at the moment). Now the terminal enter interaction mode. Let's open another terminal and run the console consumer:
> bin/kafka-console-consumer.sh --zookeeper 192.168.2.2:2181 --topic verification-topic
Now enter some data in the console producer terminal and you should see the data immediately display in the console consumer terminal.
Firstly download the kafka 0.8.1.1 from the link below:
https://www.apache.org/dyn/closer.cgi?path=/kafka/0.8.1.1/kafka_2.8.0-0.8.1.1.tgz
Next "tar -xvzf" the kafka_2.8.0-0.8.1.1.tgz file and move it to a destination folder (say, /Documents/Works/Kafka folder under the user root directory):
> tar -xvzf kafka_2.8.0-0.8.1.1.tgz > mkdir $HOME/Documents/Works/Kafka > mv kafka_2.8.0-0.8.1.1 $HOME/Documents/Works/Kafka
Now go back to the user root folder and open the .bashrc file for editing:
> cd $HOME > gedit .bashrc
In the .bashrc file, add the following line to the end:
export KAFKA_HOME=$HOME/Documents/Works/Kakfa/kafka_2.8.0-0.8.1.1
Save and close the .bashrc and run "source .bashrc" to update the environment variables. Now navigate to the kafka home folder and edit the server.properties in its sub-directory "config":
> cd $KAFKA_HOME/config > gedit server.properties
In the server.properties file, search the line "zookeeper.connect" and change it to the following:
zookeeper.connect=192.168.2.2:2181,192.168.2.4:2181
search the line "log.dirs" and change it to the following:
log.dirs=/var/kafka-logs
Save and close the server.properties file (192.168.2.2 and 192.168.2.4 are the zookeeper nodes). Next we go and create the folder /var/kafka-logs (which will store the topics and partitions data for kafka) with write permissions:
> sudo mkdir /var/kafka-logs
> sudo chmod -R 777 /var/kafka-logs
Now set up and run the zookeeper cluster by following instructions in the link http://czcodezone.blogspot.sg/2014/11/setup-zookeeper-in-cluster.html. Once this is done, we are ready to start the kafka messaging system by running the following commands:
> cd $KAFKA_HOME > bin/kafka-server-start.sh config/server.properties
To start testing kafka setup, Ctrl+Alt+T to open a new terminal and run the following command to create a topic "verification-topic" (a topic is a named entity in kafka which contain one or more partitions which are message queues that can run in parallel and serialize to individual folder in /var/kafka-log folder):
> cd $KAKFA_HOME > bin/kafka-topics.sh --create --zookeeper 192.168.2.2:2181 --topic verification-topic --partitions 1 --replication-factor 1
The above command creates a topic named "verification-topic" which contains 1 partition (and with no replication)
Now we can check the list of topics in kafka by running the following command:
> bin/kafka-topics.sh --zookeeper 192.168.2.2:2181 --list
To test the producer and consumer interaction in kafka, fire up the console producer by running
> bin/kafka-console-producer.sh --broker-list localhost:9092 --topic verification-topic
9092 is the default port for a kafka broker node (which is localhost at the moment). Now the terminal enter interaction mode. Let's open another terminal and run the console consumer:
> bin/kafka-console-consumer.sh --zookeeper 192.168.2.2:2181 --topic verification-topic
Now enter some data in the console producer terminal and you should see the data immediately display in the console consumer terminal.
Subscribe to:
Posts (Atom)