Публікації

Показано дописи з міткою "Real-time"

Real-time сохранение изменяющихся данных в Apache SOLR с помощью Kafka + Lily HBase Batch Indexer

   Apache Solr очень хороший инструмент как для поиска так и для аналитики. Но изменение данных является дорогой операцией и желательно это не делать вообще или делать это как можно реже.    Для этого был разработан подход на базе key-value базы HBase которая хорошо работает с обновлением данных. После сохранение обновлений в HBase, данные с помощью Lily HBase Batch Indexer считываются изменения с update log-a в HBase и заливают изменения в индекс Solr. Для начала создаём таблицу в HBase: hbase shell hbase(main):021:0> create 'table_for_index', {NAME => 'data', REPLICATION_SCOPE => 1} hbase(main):021:0> put 'table_for_index', 'row1', 'data', 'value' hbase(main):022:0> put 'table_for_index', 'row2', 'data', 'value2' Далее создаём индекс в Solr в который будут сохранятся изменения из HBase: solrctl instancedir --generate $HOME/hbase-index //Вносим необходимые поля из HBase таблички v...

Реализация real-time загрузки данных в Hive c помощью Kafka topic и Apache Flume

Для сохранения данных с Kafka topic напрямую в Hive можно использовать HiveSink: Если данные в Kafka топике у Вас сохранены у Вас  в Json-e, то для этого есть  .serializer = JSON. Также возможен вариант DELIMITED с последующим указанием разделителя для значений. Создаём таблицу в Hive для загрузки данных: CREATE TABLE `db_name.table_name`( `date` string,`cost` string,`cost_origin` string,`campaign` string, `currency` string) CLUSTERED BY(campaign) INTO 10 BUCKETS ROW FORMAT SERDE 'org.apache.hadoop.hive.ql.io.orc.OrcSerde' STORED AS INPUTFORMAT 'org.apache.hadoop.hive.ql.io.orc.OrcInputFormat' OUTPUTFORMAT 'org.apache.hadoop.hive.ql.io.orc.OrcOutputFormat' LOCATION '/user/flume/flumeingest/db.table_name'; Пример настройки Hive Sink: tier1.sinks.HiveSink5.channel = channel5 tier1.sinks.HiveSink5.type = hive tier1.sinks.HiveSink5.hive.metastore = thirft://127.0.0.1:9083 tier1.sinks.HiveSink5.hive.database = db_name tier1.sinks.HiveSink5.hive...

Real-time cохранение данных в SOLR с помощью Kafka + Flume

   Для заливки данных в колекцию SOLR можно использовать сочетание Kafka  + Flume MorphlineSolrSink. Это позволяет быстро настроить добавление новых данных в SOLR даже без необходимости писать код.     Предварительно сохранённая в формате json информация в Kafka topic читается, преобразовывается и сохраняется в SOLR с помощью MorphlineSolrSink. При этом все нюансы предподготовки данных перед сохранением в документ SOLR можно описать в конфигурационном файле для morphline.conf. Файл конфигураций для flume: tier1.sources = source1 tier1.channels = channel1 tier1.sinks = solrSink1 tier1.sources.source1.channels = channel1 tier1.sources.source1.type = org.apache.flume.source.kafka.KafkaSource tier1.sources.source1.zookeeperConnect = 127.0.0.1:2181 Zookeeper client port tier1.sources.source1.topic = test_topic tier1.sources.source1.kafka.auto.offset.reset = smallest tier1.sources.source1.groupId = flume_source_test_topic tier1.sources.source1.batch...

Оптимизация Apache Spark Streaming с помощью rdd.cache и настройки сериализации данных.

Зображення
   Одним из эффективных способов существенного ускорения вашего Spark Streaming приложения, являвляется применения кеширования RDD а также конфигурация способа кеширования ваших данных.    При использовании rdd.cache первая операция по rdd будет занимать тоже время что и без применения кешированя, но все последующие операции будут занимать значительно меньше времени. Это может быть удивительно, ведь Spark стараеться делать все операции с данными в памяти. Это сложно понять, но приминение rdd.cache имеет значимый эффект..    Тесты показывают, что если DStream или RDD используется несколько раз, то их кеширование значительно увеличивает скорость их обработки. Возможно это будет являться тем изменением, что принесёт вам наибольшее улучшнеием производительности. Пример применения rdd.cache: dstream.foreachRDD{rdd => rdd.cache() // cache the RDD before iterating! rdd.foreach{ key => rdd.filter(elem=> key(elem) == key).saveAsFooBar(....

Понимание параметров Apache Spark - пошаговое руководство для настройки и оптимизации Apache Spark

Всего несколько параметров позволять Вам существенно увеличить производительность Spark Streaming приложений, которые читают топики с Kafka и записывают их в хранилище(например в HDFS). Первым, что стоит оптимизировать, это получателя сообщений с Kafka. Таким получателем выступает Ваше Spark Streaming приложение, которое читает сообщения с топика Kafka порциями(блоками, партициями). Размер этих блоков нужно настраивать относительно предположительной скорости сохранения данных в хранилище. По умолчанию размер этих блоков не установлен и его нужно определить исходя из скорости обработки сообщений Вашим получателем. Вот формула по которой можно понять, сколько же сообщений за раз получит Ваш воркер: partitionSize = (1000 / blockInterval) * maxRate После нескольких запусков Вы уже приблизительно понимаете сколько записей в секунду обрабатывает Ваш получатель. К примеру мы знаем, что воркер справляется приблизительно с 1000 записей в секунду. Возьмём половину от этого и установим ...

Проблема с файлами маленьких размеров при сохранении данных в hdfs c помощью Аpache Kafka & Flume

Flume показался наиболее подходящим инструментом для сохранения потоковых данных в hdfs. Изначально был ностроен flume-agent для сохранения данных в hdfs следующим образом: tier1.sources = source1 tier1.channels = channel1 tier1.sinks = sink1 tier1.sources.source1.type = org.apache.flume.source.kafka.KafkaSource tier1.sources.source1.zookeeperConnect = 127.0.0.1:2181 tier1.sources.source1.topic = users tier1.sources.source1.groupId = 67 tier1.sources.source1.channels = channel1 tier1.sources.source1.interceptors = i1 tier1.sources.source1.interceptors.i1.type = timestamp tier1.sources.source1.kafka.consumer.timeout.ms = 100 tier1.channels.channel1.type = memory tier1.channels.channel1.capacity = 10000 tier1.channels.channel1.transactionCapacity = 1000 tier1.sinks.sink1.type = hdfs tier1.sinks.sink1.hdfs.path = /user/hdfs/tables/%{topic}/%Y-%m-%d # rollover file based on max time of 2 min tier1.sinks.sink1.hdfs.rollInterval = 120 # rollover file based on maximum size of 0 MB ti...

Apache Kafka - краткое описание

Apache Kafka - это распределённая, легко маштабируемая система обмена сообщениями c высокой пропускной способностью(Kafka быстрее  RabbitMQ раз в  5-20 ). Пишет сразу на диск сообщения и хранит там их указанный период времени. Kafka является единственным проектом, который на уровне архитектуры решает вопрос импорта большого объёма данных. Проблемы: Теоретически может выдерживать любые объёмы данных, но на практике показатели сильно преувеличены. От версии к версии интерфейс может полностью изменится, что очень мешает. Не работает ряд функций: группы потребителей, сдвиги для пользователей.  Простой рецепт, это запускать по одному потребителю на партицию очереди (topic, в терминологии Kafka) и вручную контролировать сдвиги.  Приемущества: Легко маштабируется Высокая производительность как для паблишеров так и для подписчиков Отказоустойчивая и автоматическо балансируется в случае отказа Полезные ссылки: Документация Next Generation Distributed Me...

Apache Storm - краткое описание

Apache Storm - распределенная система для обработки больших обьемов данных в реальном времени. Гарантирует отказоустойчивость благодаря механизму отслеживания успешной обработки данных. Лучше использовать вместе с Apache Kafka . Storm быстрее чем аналог Spark Streaming ;) Полезные ссылки: Документация развертывания кластера Описание от IBM Цыкл статей на habre Аналог от Spark - Spark Streaming Hortonworks: real-time events Udacity course: Real-Time Analytics with Apache Storm

Spark Streaming - краткое описание

Зображення
Spark Streaming - расширение ApacheSpark в виде маштабируемого отакзоустойчевого обработчика потоков данных в реальном времени. Позволяет писать и запускать приложения в потоковом режиме. При этом, приложение сможет одновременно работать как с потоками данных, так и осуществлять пакетную обработку – без существенных изменений в коде. Плюс ко всему, фреймворк способен автоматически восстанавливать данные после ошибочных действий со стороны системы – от пользователя в этом случае не потребуется написать ни строчки кода. Spark Streaming умеет забирать данные за фиксированный промежуток времени (например, за 30 секунд) из Kafka , Flume, ZeroMQ, Kinesis, TCP сокета, и т.д. Обрабатывать их и сохранять в файловую систему, базу данных. Можно применять функции map, reduce, join, машинное обучение и т.д. Обработанные данные могут быть записаны в файловую систему, базу данных, realtime графики. Можно применить машинное обучение и алгоритмы для работы с графами для обработки потока дан...