Интеграция Apache Kafka и Spark
Пару слов о Spark
Spark Streaming API обеспечивает масштабируемую, высокопроизводительную и отказоустойчивую обработку потоков данных в режиме реального времени. Данные могут поступать из различных источников, таких как Kafka, Flume, Twitter и т. д.. Их обработка происоходит с помощью сложных алгоритмов, таких как высокоуровневые функции map, reduce, join и window. Обработанные данные могут быть переданы в файловые системы, базы данных и дашборды. Устойчивые распределенные наборы данных (Resilient Distributed Datasets, RDD) - это основах Spark, представляющая собой неизменяемую распределенную коллекцию объектов. Каждый набор данных в RDD делится на логические разделы, каждый из которых может вычисляться на разных узлах кластера.
Интеграция с Spark
Kafka - это платформа обмена сообщениями, которая может быть интегрирована со = Spark. Kafka выступает в качестве центрального узла для потоков данных в режиме реального времени, которые обрабатываются с помощью сложных алгоритмов в Spark Streaming. После обработки данных Spark Streaming может публиковать результаты в топик Kafka или хранить их в HDFS, базах данных или дашбордах.
На диаграмме, представленной ниже, показан схематический поток данных.
Теперь давайте подробно рассмотрим интеграцию API Kafka-Spark.
SparkConf API
Представляет собой конфигурацию для приложения Spark, которая применяется для определения его ключевых параметров.
Параметры:
- set(string key, string value) – устанавливает переменную конфигурации;
- remove(string key) – удаляет ключ из конфигурации;
- setAppName(string name) − задает имя приложения;
- get(string key) – получает ключ
StreamingContext API
Это основная точка входа в систему Spark и доступа к ее функционалу. SparkContext представляет собой соединение с кластером Spark и может использоваться для создания RDD и широковещательных переменных на кластере.
public StreamingContext(String master, String appName, Duration batchDuration, String sparkHome, scala.collection.Seq<String> jars, scala.collection.Map<String,String> environment)
- master − URL-адрес кластера, к которому необходимо подключиться (e.g. mesos://host:port, spark://host:port, local[4]).
- appName – название Вашего задания для отображения в веб-интерфейсе кластера
- batchDuration − интервал времени, через который потоковые данные будут делиться на партии
Создайте StreamingContext, предоставив конфигурацию, необходимую для новой SparkContext.
- conf – параметры Spark
- batchDuration − интервал времени, через который потоковые данные будут делиться на партии
KafkaUtils API
API KafkaUtils используется для подключения кластера Kafka к потоковой передаче данных Spark. В этом API есть важная метод createStream:
public static ReceiverInputDStream<scala.Tuple2<String,String>> createStream( StreamingContext ssc, String zkQuorum, String groupId, scala.collection.immutable.Map<String,Object> topics, StorageLevel storageLevel)
Метод, описанный выше, используется для создания входного потока, извлекающего сообщения из брокеров Kafka.
- ssc – объект StreamingContext
- zkQuorum – кворумZookeeper
- groupId − группа id для консьюмера.
- topics – возвращает список топиков для консьюмера.
- storageLevel − Уровень хранения, используемый для хранения полученных объектов.
В API KafkaUtils есть еще один метод под названием createDirectStream, который используется для создания входного потока, напрямую извлекающего сообщения из брокеров Kafka без использования приемника. Этот метод гарантирует то, что каждое сообщение из Kafka будет обработано лишь один раз.
Пример приложения, приведенный ниже, выполнен на языке Scala. Для компиляции приложения загрузите и установите sbt и инструмент сборки Scala (аналогичный maven). Основной код приложения представлен ниже.
import java.util.HashMap
import org.apache.kafka.clients.producer.{KafkaProducer, ProducerConfig, Produc-erRecord}
import org.apache.spark.SparkConf
import org.apache.spark.streaming._
import org.apache.spark.streaming.kafka._
object KafkaWordCount {
def main(args: Array[String]) {
if (args.length < 4) {
System.err.println("Usage: KafkaWordCount <zkQuorum><group> <topics> <numThreads>")
System.exit(1)
}
val Array(zkQuorum, group, topics, numThreads) = args
val sparkConf = new SparkConf().setAppName("KafkaWordCount")
val ssc = new StreamingContext(sparkConf, Seconds(2))
ssc.checkpoint("checkpoint")
val topicMap = topics.split(",").map((_, numThreads.toInt)).toMap
val lines = KafkaUtils.createStream(ssc, zkQuorum, group, topicMap).map(_._2)
val words = lines.flatMap(_.split(" "))
val wordCounts = words.map(x => (x, 1L))
.reduceByKeyAndWindow(_ + _, _ - _, Minutes(10), Seconds(2), 2)
wordCounts.print()
ssc.start()
ssc.awaitTermination()
}
}
Составление скрипта
Интеграция spark-kafka зависит от jar интеграции spark, spark streaming и spark Kafka. Создайте новый файл build.sbt и укажите в нем данные о приложении и его зависимости. В процессе компиляции и упаковки приложения sbt загрузит необходимый jar-файл.
name := "Spark Kafka Project" version := "1.0" scalaVersion := "2.10.5" libraryDependencies += "org.apache.spark" %% "spark-core" % "1.6.0" libraryDependencies += "org.apache.spark" %% "spark-streaming" % "1.6.0" libraryDependencies += "org.apache.spark" %% "spark-streaming-kafka" % "1.6.0"
Компиляция / Упаковка
Для компиляции и упаковки jar-файла приложения выполните следующую команду. Для запуска приложения отправьте jar-файл в консоль spark
sbt package
Отправка приложения в Spark
Запустите Kafka Producer CLI (смотрите предыдущую статью), создайте новый топик под названием my-first-topic и предоставьте несколько сообщений:
Another spark test message
Для отправки приложения в консоль Spark используйте следующую команду:
/usr/local/spark/bin/spark-submit --packages org.apache.spark:spark-streaming -kafka_2.10:1.6.0 --class "KafkaWordCount" --master local[4] target/scala-2.10/spark -kafka-project_2.10-1.0.jar localhost:2181 <group name> <topic name> <number of threads>
Пример вывода приложения:
spark console messages .. (Test,1) (spark,1) (another,1) (message,1) spark console message ..




