Skip to content

Latest commit



312 lines (195 loc) · 12 KB

File metadata and controls

312 lines (195 loc) · 12 KB

Table of content

Set up a complete Kafka environment

Here, we have provisioned two kafka nodes with two zookeeper nodes.


  • Create 4 linux based VM/machines. (Ubuntu 17.10, CPU 4 Core, Memory 8GB, Disk 250GB)

Step #1: Install JAVA on all machines

Followed below steps to install Java on unbutu:

sudo add-apt-repository ppa:openjdk-r/ppa
sudo apt-get update 
sudo apt-get install openjdk-8-jdk 

Set up the appropriate environment variables for Java.

sudo vi /etc/profile
export JAVA_HOME=/usr/lib/jvm/jre-1.8.0-openjdk
export JRE_HOME=/usr/lib/jvm/jre

Step #2: Install Kafka

Installing kafka is very simple. One can check the latest release here.

  • First, download the kafka tar on all machine
  • Untar it
tar -xvzf kafka_2.12-1.0.1.tgz
  • Explore the extracted folder. You will see bin and config folders. These are important folders as you will be working with those in next steps.

Step #3: Configure Kafka nodes

After downloading Kafka, the next step is to change/update some configuration on Kafka.

  • Edit /kafka_2.12-1.0.1/config/ file. Provide IPs of zookeeper server (don't worry, we'll be configuring zookeeper in later steps), log output path and broker identifier (remember, it should be unique id for each kafka server).
  • Make sure that log output path already exist on nodes with right execution permissions.
mkdir /tmp/kafka-logs-0
drwxr-xr-x  2 root root 4096 Mar 19 17:17 kafka-logs-0

Step #4: Configure Kafka Zookeeper nodes

In this step, we will be configuring zookeeper nodes. Zookeeper is a high avalibilty coordination service that Kafka uses for coordination among brokers.

I followed the information provided in the apache zookeeper doc.

  • Edit /kafka_2.12-1.0.1/config/ file. Provide data directory path. Make sure dataDir file exists on the nodes.
# the port at which the clients will connect
# disable the per-ip limit on the number of connections since this is a non-production config

To let every zookeeper know about every other nodes in ensemble, we need to provide IDs, IPs and ports as follows.


123 and 125 are unique identifier provided in dataDir/myid file on each nodes.

myid file consists of a single line containing id of node.

2888 this is the port used by followers to connect to leader and 3888 is for leader election. (Make sure these ports are open on nodes)

Step #5: Start Apache Kafka and Zookeeper services on nodes

  • Run Kafka as demon process on Kafka nodes
nohup bin/ config/ &
./bin/ -daemon config/


Get the list of active brokers

./bin/ <<< "ls /brokers/ids"

screen shot 2018-03-23 at 3 03 16 pm

Push a topic with replica

Push a topic within Kafka cluster (containing two Kafka servers) with two partitions and replicate partition over two Kafka broker.

./bin/ --create --zookeeper --topic test-topic --partitions 2 --replication-factor 2

Created topic "test-topic".

Note that under the Kafka node's logs path, two partition has been created test-topic-0 & test-topic-1 on both nodes.

screen shot 2018-03-22 at 6 52 12 pm

Push a topic without replica

Push a topic within Kafka cluster (containing two Kafka servers) with two partitions and no replication.

./bin/ --create --zookeeper --topic non-repeat-topic --partitions 2 --replication-factor 1

Created topic "non-repeat-topic".

You will observe below.

On Kafka node #1:

screen shot 2018-03-22 at 7 02 30 pm

On Kafka node #2:

screen shot 2018-03-22 at 7 02 45 pm

Get the list of topics

Get the list of topics on Kafka cluster.

./bin/ --list --zookeeper


Describe topic

Describe test-topic and non-repeat-topic topics.

./bin/ --zookeeper --describe --topic test-topic

screen shot 2018-03-22 at 7 16 24 pm

screen shot 2018-03-22 at 7 19 10 pm

Topic: test-topic	Partition: 0	Leader: 1	Replicas: 1,0	Isr: 1

represents Kafka broker with id 0 acting as leader and handling all reads/writes for partition 0. Replicas: 1,0 meaning partition 0 is replicated on broker id 0 and 1.

Delete topic

Delete topic test-topic from Kafka cluster.

To enable deletion, we need to make sure to add delete.topic.enable=true on all Kafka brokers config file.

./bin/ --zookeeper --delete --topic test-topic

Run a producer

Run a producer process to publish data into test-topic Kafka topic within the Kafka broker.

bin/ --broker-list, --topic test-topic

Create two consumer groups

Edit /kafka_2.11-1.0.1/config/ with group id test-consumer-group-0 and test-consumer-group-1 by creating another consumer.properties1 config.

# consumer group id

Run a consumer process

Run a consumer process to process published data from a topic.

bin/ --zookeeper --from-beginning --topic topic-new

Run multiple consumer processes

First, run a producer process for pushing data to a topic-new topic.

bin/ --broker-list, --topic topic-new

Run two consumer processes in a consumer group.

bin/ --zookeeper --topic topic-new --consumer.config config/

Let's add messages in producer console as below:

screen shot 2018-03-23 at 5 23 59 pm

This is how consumer processes would process messages through Kafka cluster.

screen shot 2018-03-23 at 5 24 43 pm

screen shot 2018-03-23 at 5 24 55 pm


Each consumer in a consumer group is guaranteed to read a particular message by only one consumer in the group. In other words, data pushed to a Kafka topic is only processed once in a consumer group. Or we can say, the processing of data is distributed among consumer processes in a consumer group.

Run multiple consumer groups

First, run a producer process for pushing data to a cast-topic topic.

bin/ --broker-list, --topic cast-topic

Run two consumer processes in two seperate consumer groups. We need two consumer-properties config files as covered above.

bin/ --zookeeper --topic cast-topic --consumer.config config/
bin/ --zookeeper --topic cast-topic --consumer.config config/consumer.properties1

Now, put [one, two, three, four, five, six] messages in cast-topic through producer process. Following is how consumer processes would receive messages.

screen shot 2018-03-23 at 8 26 37 pm

screen shot 2018-03-23 at 8 26 30 pm


Here, consumers from different consumer group receive all the messages on a Kafka topic. Meaning, if multiple consumer groups subscribe to a topic, then Kafka would broadcast messages to each of them.


  • What happens if all the Kafka brokers died?

The consumer pulling data from topic gets Error while fetching metadata with correlation id 43 : {test-topic=LEADER_NOT_AVAILABLE} (org.apache.kafka.clients.NetworkClient) error.


Kafka is a publish-subscribe based messaging system that can be used to exchange data between processes, application and services. It has build-in partitioning, replication and fault-tolerance.

As we have already seen above, Kafka has capability to scale-processing (by distributing the data processing among consumer processes in a consumer group) and multi-subscriber (by broadcasting messages among consumer groups).