Skip to content

Instantly share code, notes, and snippets.

Last active December 22, 2022 09:39
What would you like to do?
Kafka Producer/Consumer Example in Scala
import java.util
import org.apache.kafka.clients.consumer.KafkaConsumer
import scala.collection.JavaConverters._
object ConsumerExample extends App {
import java.util.Properties
val TOPIC="test"
val props = new Properties()
props.put("bootstrap.servers", "localhost:9092")
props.put("key.deserializer", "org.apache.kafka.common.serialization.StringDeserializer")
props.put("value.deserializer", "org.apache.kafka.common.serialization.StringDeserializer")
props.put("", "something")
val consumer = new KafkaConsumer[String, String](props)
val records=consumer.poll(100)
for (record<-records.asScala){
object ProducerExample extends App {
import java.util.Properties
import org.apache.kafka.clients.producer._
val props = new Properties()
props.put("bootstrap.servers", "localhost:9092")
props.put("key.serializer", "org.apache.kafka.common.serialization.StringSerializer")
props.put("value.serializer", "org.apache.kafka.common.serialization.StringSerializer")
val producer = new KafkaProducer[String, String](props)
val TOPIC="test"
for(i<- 1 to 50){
val record = new ProducerRecord(TOPIC, "key", s"hello $i")
val record = new ProducerRecord(TOPIC, "key", "the end "+new java.util.Date)
Copy link

Above example works fine simply with below dependencies:-

"org.apache.kafka" % "kafka_2.11" % "1.1.1"
scala 2.11
spark 2.3.0

Ran it through spark shell with below commands

spark2-shell --packages org.apache.kafka:kafka_2.11:1.1.1

Copy link

anigos commented Feb 25, 2019

"org.apache.kafka" % "kafka_2.11" % "1.1.1" this worked fine for me.

name := "myApp"
version := "0.1"
scalaVersion := "2.11.12"
val sparkVersion = "2.3.0"

libraryDependencies ++= Seq(
"org.apache.bahir" %% "spark-streaming-twitter" % sparkVersion,
"org.apache.spark" %% "spark-core" % sparkVersion,
"org.apache.spark" % "spark-sql_2.11" % sparkVersion,
"org.apache.spark" %% "spark-streaming" % sparkVersion,
"org.apache.kafka" % "kafka_2.11" % "1.1.1"

Copy link

I'm using eclipse IDE getting below error. Any idea what is the issue.
log4j:WARN No appenders could be found for logger (org.apache.kafka.clients.producer.ProducerConfig).
log4j:WARN Please initialize the log4j system properly.
log4j:WARN See for more info.

Copy link

Hi, I am new to Kafka (using for only 3 days now) but I think
[error] Modules were resolved with conflicting cross-version suffixes in {file:/home/edureka/Desktop/code_practicals/}code_practicals:
[error] org.apache.kafka:kafka _2.11, _2.10

this means that you should use scala version 2.11 instead of 2.10
correct me if I am wrong...

Copy link

A full kafkasteams example here

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment