Публікації

Пример подключения к Impala в RStudio с помощью jdbc драйвера

Устанавливаем R на ubuntu: sudo apt-get -y install r-base sudo R CMD javareconf Далее скачиваем и устанавливаем RStudio  . Также скачиваем и распаковываем архив jdbc драйвера для подключения к Impala  . После запуска RStudio указываем путь к jdbc драйверу и работаем с базой: install.packages("RImpala") library(RImpala) rimpala.init(libs="/tmp/impala/jars/") rimpala.connect("192.168.10.1","21050") rimpala.invalidate() rimpala.showdatabases() rimpala.usedatabase("yourdatabase") rimpala.showtables() rimpala.describe("yourtablename") Полезные ссылки: http://blog.cloudera.com/blog/2013/12/how-to-do-statistical-analysis-with-impala-and-r/ https://github.com/Mu-Sigma/RImpala

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...

Пример использования "Regular expression filtering" в Flume

   Основываясь на регулярном выражении можно фильтровать сообщения на основе контента сообщения. К примеру Вы можете передавать на сохранение только те сообщения, которые совпадают регулярному выражению и наоборот Вы можете отфильтровывать ненужные Вам сообщение регулярным выражением. За это отвечает флаг excludeEvents = true | false. Для регулярных выражений используется Java-style синтаксис . Пример использования: agent.sources.projectSource1.interceptors=filterErr agent.sources.projectSource1.interceptors.filterErr.type=regex_filter agent.sources.projectSource1.interceptors.filterErr.regex= ERROR [0-4]: #agent.sources.projectSource1.interceptors.filterErr.regex= ".*(\\_event\\_name\\\"\\:\\\"cost\\_importer\\_data).*" agent.sources.projectSource1.interceptors.filterErr.excludeEvents=false Полезные ссылки: http://flume.apache.org/FlumeUserGuide.html#regex-filtering-interceptor http://docs.oracle.com/javase/6/docs/api/java/util/regex/Pattern.html

Импорт данных с Vertica в Hive помощью Sqoop

sqoop import --m=1 --connection-manager="org.apache.sqoop.manager.GenericJdbcManager" --driver='com.vertica.jdbc.Driver' --connect "jdbc:vertica://127.0.0.1:5433/schema?searchpath=db" --username username --password-file="/user/sqoop/vertica_db_name_pwd.txt" --compression-codec=snappy --as-parquetfile --hive-import --hive-database db --hive-table table_name --check-column="user_id" --last-value="000" --verbose --query 'select * from db.table_name tableSel WHERE $CONDITIONS' --split-by="user_id" --target-dir="/user/sqoop/import_sqoop/sqoop_dump_table_name" --incremental=append --map-column-hive user_id=String,uid=String,user_uid=String,funnel_id=String --map-column-java user_id=String,uid=String,user_uid=String,funnel_id=String

Импорт данных с Vertica в HBase помощью Sqoop

sqoop import -Doraoop.timestamp.string=false --m=1 --connection-manager="org.apache.sqoop.manager.GenericJdbcManager" --driver='com.vertica.jdbc.Driver' --connect "jdbc:vertica://127.0.0.1:5433/schema?searchpath=db" --username username --password yourpassword --verbose --table table_name --hbase-create-table --hbase-table db.table_name --column-family data --columns "user_id,uid,user_uid,funnel_id" --hbase-row-key user_id

solrctl - создание колекции для SOLR

Создадим директорию со всеми необходимыми файлами настройки для будующей колекции документов (schema.xml, solrconf.xml) solrctl instancedir --generate /var/lib/solr/collection3 -schemaless Далее можно внести нужные нужные изменения в конфигурационные файлы. Например описать нужные поля в документе в schema.xml: <field name="dt" type="date" indexed="true" stored="true"/> <field name="hash" type="string" indexed="true" stored="true"/> <field name="funnel_id" type="string" indexed="false" stored="true"/> Или изменить фактор репликации для лога транзакций в solrconf.xml: <updateLog> <str name="dir">${solr.ulog.dir:}</str> <int name="tlogDfsReplication">2</int> </updateLog> Далее загружаем созданную директорию с настройками для колекции в SolrCloud: solrctl instancedir --create coll...