Показаны сообщения с ярлыком Java. Показать все сообщения
Показаны сообщения с ярлыком Java. Показать все сообщения

среда, 1 мая 2013 г.

Apache Thrift RPC Server. Дружим C++ и Java

Thrift - технология для организации межпроцессного взаимодействия между компонентами системы. Была разработана где то в недрах Facebook. Посути это кросс-языковой фреймворк для создания RPC сервисов, на бинарном протоколе. С помощью этого решения можно "подружить" компоненты написанные на разных языках  C#, C++, Delphi, Erlang, Go, Java, PHP, Python, Ruby, итд. Описание сигнатур сервисов и данных осуществляется с помощью специального IDL - языка. Технология, по своей сути, похожа на COM, но без всей этой обвязки с регистрацией компонент. Так же не будем забывать, что COM это технология только для Windows, в то время как Thrift - кросплатформенна.

Вобщем решил поэкспериментировать, попробовать вынести часть нагруженной-вычислительной логики из Java в С++, в надежде что нативный С++ код будет производительней, за одно опробовать Thrift RPC, в надежде что это быстрее чем REST.
Как и положено, без бубнов и граблей не обошлось!

И так, для начала надо всё это поставить:
1. Ставим поддержку для Boost, ибо всё завязано на нём

$ sudo apt-get install libboost-dev libboost-test-dev libboost-program-options-dev libevent-dev automake libtool flex bison pkg-config g++ libssl-dev

2. качаем thrift tarball http://apache.softded.ru/thrift/0.9.0/thrift-0.9.0.tar.gz
распаковываем, запускаем configure, затем собираем

$ ./configure
$ make
$ sudo make install

Вроде бы всё... Можно даже попробовать сгенерировать код из учебника, который идет вместе с thrift tarball

$ thrift --gen cpp tutorial.thrift

команда thrift сгенерирует С++ обвертки, и бережно положит их в директорию gen-cpp. Тоже самое можно сделать для Java, PHP итд...

Пробуем скомпилировать и собрать нагенеренные исходники

$ g++ -Wall -I/usr/local/include/thrift *.cpp -L/usr/local/lib -lthrift -o something

Упс, получите:  error: ‘uint32_t’ does not name a type
Оказывается есть небольшой таракан в библиотеках thrift связанный с uint32_t.  Лечится  добавлением "#include <stdint.h>" в "Thrift.h", или (что лучше всего) специальными опциями компилятора -DHAVE_NETINET_IN_H -DHAVE_INTTYPES_H

Теперь это выглядит так:

$ g++ -DHAVE_INTTYPES_H -DHAVE_NETINET_IN_H -Wall -I/usr/local/include/thrift *.cpp -L/usr/local/lib -lthrift -o something

Вот и всё, появился исполнимый файл, под названием something.
Запускаем, и получаем: error while loading shared libraries: libthrift.so.0: cannot open shared object file: No such file or directory
Возможно есть какие-то элегантные методы решения этой проблемы, но я решил её в лоб, копированием thrift файлов из /usr/local/lib в /lib

Всё, пример запустился. Значит, все на местах, и всё работает.

Теперь можно писать свой RPC сервер.
Его задача, быть key-value хранилищем. Хранить длинные (в несколько сот тысяч) битовые маски. Складывать их (AND), и отдавать клиенту массив индексов, в которых получились еденицы. Да, почти тоже самое умеет Redis, но он мне не подходит.

Полный код лежит здесь: https://github.com/2anikulin/fast-hands.git

Описываем сигнатуры данных и сервисов, в thrift definition file:

namespace cpp fasthands
namespace java fasthands
namespace php fasthands
namespace perl fasthands

exception InvalidOperation {
  1: i32 what,
  2: string why
}

service FastHandsService {

 i32 put(1:i32 key, 2:binary value),
 
 binary get(1:i32 key),
 
 list <i32> bitAnd(1:list<i32> keys) throws (1:InvalidOperation ouch)
}

И генерируем обвертки.
Реализация на C++
Этот код, создает, и запускает RPC сервер

#define PORT 9090
#define THREAD_POOL_SIZE 15


int main() {

  printf("FastHands Server started\n");

  shared_ptr<TProtocolFactory> protocolFactory(new TBinaryProtocolFactory());
  shared_ptr<FastHandsHandler> handler(new FastHandsHandler());
  shared_ptr<TProcessor> processor(new FastHandsServiceProcessor(handler));

  shared_ptr<ThreadManager> threadManager = ThreadManager::newSimpleThreadManager(THREAD_POOL_SIZE);
  shared_ptr<PosixThreadFactory> threadFactory = shared_ptr<PosixThreadFactory>(new PosixThreadFactory());
  threadManager->threadFactory(threadFactory);
  threadManager->start();

  TNonblockingServer server(processor, protocolFactory, PORT, threadManager);
  server.serve();
  printf("done.\n");

  return 0;
}

В классе FastHandsHandler - имплементируется весь наш, прикладной функционал
Это заголовочный файл

class FastHandsHandler : virtual public FastHandsServiceIf {

 public:
  FastHandsHandler();
  int32_t put(const int32_t key, const std::string& value);
  void get(std::string& _return, const int32_t key);
  void bitAnd(std::vector<int32_t> & _return, const std::vector<int32_t> & keys);

 private:
  void appendBitPositions(std::vector<int32_t> & positions, unsigned char bits, int offset);

 private:
  std::map<int32_t, std::string> m_store;

};


Пробуем собрать, и получаем очередную ошибку: c++ undefined reference to apache::thrift::server::TNonblockingServer
Дело в том, что, в отличии от учебника, мой сервер - ассинхронный, и использует класс TNonblockingServer. Чтобы код собирался, надо добавить дополнительные библиотеки -lthriftnb -levent

т.е сборка сейчас будет выглядеть так:
g++ -DHAVE_INTTYPES_H -DHAVE_NETINET_IN_H -Wall -I/usr/local/include/thrift *.cpp -L/usr/local/lib -lthrift -lthriftnb -levent -o something

Теперь все хорошо. Код скомпилирован, слинкован, на выходе файл по имени something
Генерируем обвертки для java, и пишем вот такого клиента


import fasthands.FastHandsService;
import fasthands.InvalidOperation;
import org.apache.thrift.TException;
import org.apache.thrift.transport.TFramedTransport;
import org.apache.thrift.transport.TTransport;
import org.apache.thrift.transport.TSocket;
import org.apache.thrift.transport.TTransportException;
import org.apache.thrift.protocol.TBinaryProtocol;
import org.apache.thrift.protocol.TProtocol;

import java.nio.ByteBuffer;
import java.util.ArrayList;
import java.util.List;

public class JavaClient {

    public static void main(String [] args) {
        TTransport transport = new TFramedTransport(new TSocket("localhost", 9090));
        TProtocol protocol = new TBinaryProtocol(transport);
        final FastHandsService.Client client = new FastHandsService.Client(protocol);

        final List<Integer> filters = new ArrayList<Integer>();

        try {
            transport.open();

            int count = 12500;
            byte bt[] = new byte[count];
            for (int i =0; i < count; i++) {
                bt[i] = (byte)0xFF;
            }

            for (int i = 0; i < 50; i++) {
                client.put(i, ByteBuffer.wrap(bt)) ;
                filters.add(i);
            }

            List<Integer> list = client.bitAnd(filters);
            System.out.println(list.size());  

        } catch (TTransportException e) {
            e.printStackTrace();
        } catch (TException e) {
            e.printStackTrace();  
        }

        transport.close();
    }
}

Что в итоге.
Интересная технология, и не плохой способ прикрутить транспортный функционал к голому коду на С++. Но не скажу, что это намного быстрее чем REST, бинарные данные прекрасно передаются и по http. Что касается производительности кода, вынесенного из Java в С++, то чуда не произошло, он быстрее в 1.2 - 1.4 раза, ибо сериализация + расходы на транспортный уровень.


Полезные ссылки:
http://thrift.apache.org/
http://wiki.apache.org/thrift/ThriftUsageC%2B%2B
http://fundoonick.blogspot.ru/2010/06/sample-thrift-program-for-server-in.html

четверг, 28 февраля 2013 г.

Экспорт Hadoop - MapReduce в MongoDB

Понадобилось направить в MongoDB выход хадуповского M/R .
Пришлось повозиться.

Используем MongoDB+Hadoop Connector, не забываем про avro без него не будет работать
<dependency>
   <groupId>org.mongodb</groupId>
   <artifactId>mongo-hadoop-core_cdh4b1</artifactId>
   <version>1.0.0</version>
</dependency>
<dependency>
   <groupId>org.apache.hadoop</groupId>
   <artifactId>avro</artifactId>
   <version>1.1.0</version>
</dependency>
1. MongoConfigUtil.setOutputURI - надо ставить до вызова   Job job = new Job(configuration, JOB_NAME); Иначе получаем ошибку связанную с невозможностью подключиться к Монге, дело в том что setOutputURI ставит параметры в configuration, которые Job копирует при создании.

2. Надо обязательно ставить выходной тип ключа для редьюса MongoConfigUtil.setOutputKey(configuration, Text.class);
иначе по умолчанию MongoConfigUtil ставит BSON, и становится тяжело понять, почему валится job

3. Указываем Output Format  job.setOutputFormatClass(MongoOutputFormat.class);
public static Job createJob(Configuration configuration,  String... paths)   throws IOException, JSONException {
        MongoConfigUtil.setOutputURI(
            configuration, "mongodb://mongo-stg-2:27102/buffer.jobout"
        );           
        MongoConfigUtil.setOutputKey(configuration, Text.class);
        Job job = new Job(configuration, JOB_NAME);
        job.setJarByClass(MongoWriterJob.class);
        job.setReducerClass(MongoWriterReducer.class);
        job.setMapOutputKeyClass(ImmutableBytesWritable.class);
        job.setMapOutputValueClass(ResultJSONWritable.class);
        job.setOutputFormatClass(MongoOutputFormat.class);
        job.setOutputKeyClass(Text.class);
        job.setOutputValueClass(DBObject.class);
        job.setInputFormatClass(SequenceFileInputFormat.class);
        for (String path : paths) {
            FileInputFormat.addInputPath(job, new Path(path));
        }
        return job;
  }

Ну и наконец, reducer на выходе должен отдавать DBObject :
public static class MongoWriterReducer extends Reducer<
            ImmutableBytesWritable,
            ResultJSONWritable,
            Text,
            DBObject> {...


суббота, 16 февраля 2013 г.

Hadoop - Слоники большие и маленькие

Hadoop в последнее время шагает в массы. Вот и мне довелось с ним работать.
Отличная система, но грабельки большие и маленькие всё таки там раскиданы.
Вроде бы мелочи, но на решение каждой уходит по пол дня.

И так, во первых Hadoop бывает в нескольких дистрибутивных исполнениях:

1. Apache, как есть http://hadoop.apache.org/releases.html
Hadoop в чистом виде. Ставить на всех нодах надо руками, править конфиги. Учить его видеть HBase. В общем надо с головой уходить в администрирование. Свои плюсы в этом конечно есть, но где взять столько времени ?

2. Cloudera http://www.cloudera.com
Половина, из создателей Hadoop организовала свою кантору и свой дистрибутив.
Главная его изюминка - это ClouderaManager, с помощью которого можно очень легко управлять кластером, мониторить, добавлять ноды. Ставиться всё в два клика, система разворачивается на кластере, всё нужное в конфигах прописывается автоматически. В общем прелесть.
Это коммерческая система, но есть бесплатная версия, с ограничением в 50 нод на кластер.

Так выглядит Web консоль ClouderaManager


3. MapR http://www.mapr.com/
Другая половина создателей Hadoop решала сделать свой лунопарк, и выпустила свой дистрибутив. Всё как в Cloudera или в Cloudera всё как в MapR. Плюс комунити, форум, и отсутствие ограничений на размер кластера. Вобщем тоже красота.

MapR я не пробовал, а начал работать с Cloudera, по этому о всех изысках по порядку:

1. У Cloudera свой мавновский репозиторий. 
И необходимо использовать именно их зависимости в купе с версиями системы на вашем кластере
---------------------------------------------------------------------------
 <repositories>
     <repository>
         <id>apache</id>
         <url>https://repository.apache.org/content/repositories/releases/
         </url>
     </repository>
     <repository>
        <id>cloudera</id>
        <url>https://repository.cloudera.com/artifactory/cloudera-repos/
        </url>
    </repository>
</repositories>
---------------------------------------------------------------------------

<dependency>
     <groupId>org.apache.hadoop</groupId>
      <artifactId>hadoop-client</artifactId>
      <version>2.0.0-mr1-cdh4.1.2</version>
      <scope>provided</scope>
</dependency>
---------------------------------------------------------------------------

2. Минимизация Super Jar
Я привык собирать super jar. Это удобно, все нужные зависимости находятся уже внутри.
Таким способом можно избежать возможных проблем в продакшене.
Но как оказалось, для map - reduce такой подход опасен. Когда в jar включаются клоудеровские зависимости, в run time начинается какой то бардак. Такое впечатление что вместо одних методов начинают вызываются совсем другие. Поэтому из конечного архива пришлось исключить все клоудеровские зависимости. Сделать это можно просто, добавив ключ  <scope>provided</scope> к мавновской зависимости. После этого m/r работает как надо

Но если надо работать с HBase, то следует полдностью включить org.apache.hbase
Иначе хадуп просто не видит этих библиотек

3. Map-reduce значения по умолчанию
Очень важно непосредственно определять входные и выходные форматы Job-а, т.к по умолчанию они всегда текстовые.

job.setInputFormatClass(SequenceFileInputFormat.class);
job.setOutputFormatClass(SequenceFileOutputFormat.class);

Тоже самое касается форматов входа и выхода Mapper-ов и Reducer-ов. Особенно болезненно это проявляется в задачах со сцепленными Job-ми

job.setOutputKeyClass(ReportReduceWritable.class);   //Отдельно для Reduce
job.setOutputValueClass(MyKeyWritable.class);

job.setMapOutputKeyClass(Text.class);                           //Отдельно для Map
job.setMapOutputValueClass(MyReduceWritable.class);

4. Удаление директорий
Перед запуском m/r задачи надо удалять все выходные директории, которые были созданы после предыдущего запуска. Нначе Job будет ругаться и вываливаться.

5. Итератор Reducer-a нельзя использовать повторно
Интересная грабелька. Когда мы имплементируем org.apache.hadoop.mapreduce.Reducer
на вход приходит некий итератор Iterable<MyValue> values. Он позволяет считать данные только один раз. После этого перевести курсор в начало не возможно. Приходится выкручиваться сохранением элементов во временную коллекцию.


6. Запуск Job от имени другого пользователя и TableMapReduceUtil.initTableMapperJob
Убили на эти грабли день, прежде чем понять в чем дело.
Запускаем Job по крону от имени специального - системного пользователя
dmpservice. Валится !

java.io.IOException: java.lang.RuntimeException: java.io.IOException: Permission denied
at org.apache.hadoop.hbase.mapreduce.TableMapReduceUtil.findOrCreateJar(TableMapReduceUtil.java:521)
at org.apache.hadoop.hbase.mapreduce.TableMapReduceUtil.addDependencyJars(TableMapReduceUtil.java:472)
at org.apache.hadoop.hbase.mapreduce.TableMapReduceUtil.addDependencyJars(TableMapReduceUtil.java:438)
at org.apache.hadoop.hbase.mapreduce.TableMapReduceUtil.initTableMapperJob(TableMapReduceUtil.java:138)
at org.apache.hadoop.hbase.mapreduce.TableMapReduceUtil.initTableMapperJob(TableMapReduceUtil.java:215)
at org.apache.hadoop.hbase.mapreduce.TableMapReduceUtil.initTableMapperJob(TableMapReduceUtil.java:81)
at ru.crystaldata.analyticengine.mr.HBaseReaderJob.createJob(HBaseReaderJob.java:206)
at ru.crystaldata.analyticengine.JobsRunner.main(JobsRunner.java:52)
at sun.reflect.NativeMethodAccessorImpl.invoke0(Native Method)
at sun.reflect.NativeMethodAccessorImpl.invoke(NativeMethodAccessorImpl.java:39)
at sun.reflect.DelegatingMethodAccessorImpl.invoke(DelegatingMethodAccessorImpl.java:25)
at java.lang.reflect.Method.invoke(Method.java:597)
at org.apache.hadoop.util.RunJar.main(RunJar.java:208)
Caused by: java.lang.RuntimeException: java.io.IOException: Permission denied
at org.apache.hadoop.util.JarFinder.getJar(JarFinder.java:164)
at sun.reflect.NativeMethodAccessorImpl.invoke0(Native Method)
at sun.reflect.NativeMethodAccessorImpl.invoke(NativeMethodAccessorImpl.java:39)
at sun.reflect.DelegatingMethodAccessorImpl.invoke(DelegatingMethodAccessorImpl.java:25)
at java.lang.reflect.Method.invoke(Method.java:597)
at org.apache.hadoop.hbase.mapreduce.TableMapReduceUtil.findOrCreateJar(TableMapReduceUtil.java:518)
... 12 more
Caused by: java.io.IOException: Permission denied
at java.io.UnixFileSystem.createFileExclusively(Native Method)
at java.io.File.checkAndCreate(File.java:1704)
at java.io.File.createTempFile(File.java:1792)
at org.apache.hadoop.util.JarFinder.getJar(JarFinder.java:156)
... 17 more

Делаем всё тоже самое но от обычного пользователя - работает.
Убираем из джоба HBase - мепперы, запускаем от dmpservice - работает.
Вывод: что то темное творится в недрах TableMapReduceUtil.initTableMapperJob

Как оказалось TableMapReduceUtil.initTableMapperJob создает темповый файл для своих внутренних нужд, и пытается разместить его в HOME, а так как у пользователя dmpservice,
не заданна данная переменная, то запись производится в root /, на который естественно нет прав.

Лечится - при создании задачи крону, прописываем HOME=var/tmp


Пока это всё. Удачи в начинаниях :)








пятница, 5 октября 2012 г.

AeroSpike (Citrusleaf)

AeroSpike, он же Citrusleaf http://www.aerospike.com/
Распределенное хранилище данных, в виде ключ - значение.
Основной фичей является крайне быстрый поиск и передача записей.

Термины AeroSpike
Node - хост, на котором функционирует AeroSpike
Cluster - объединение нод в одном информационном пространстве
Name space - База данных (Data base) - в терминах SQL
Record - запись (row)
Bin - колонка

Есть бесплатная Community Edition версия, с ограничением размера кластера, не больше 2-х нод. Но и этого вполне хватает.

К AeroSpike прилагается умный клиент. Его легко встроить в своё приложение. Реализации SDK есть на Java, C#, Python, PHP.
Он отслеживает динамическое изменения кластера, управляет балансировкой нагрузки. Т.е делает всё сам.

Работать с данным SDK очень просто

//Создаем подключение
CitrusleafClient cc = new CitrusleafClient();
cc.addHost("10.1.1.60",3000);

//Добавляем запись
String namespace = "test";
ClResultCode rc = cc.set(namespace, "myset", "mykey",
                       "mybin", "myvalue", null, null);

//Считываем запись
ClResult cr = cc.get(namespace, "myset", "mykey", "name", null);
if (cr.result != null) {
    System.out.println("got name:" + cr.result);
}

Схема name space может динамически меняться. Для каждой записи, можно добавлять или удалять бины (колонки). Т.е по сути, никакой схемы нет...

Ноды сами занимаются репликацией данных, информация полностью дублируется на каждой ноде. Ноды могут динамически удаляться и добавляться в кластер. Потери данных происходить не должно.

Проблемы:
Пока не удалось запустить кластер из 2-х нод, на двух локальных виртуальных машинах (похоже проблемы с сетевыми мостами). Зато получилось это сделать в связке 1 нода - на виртуальной машине, 2-я на хосте. Всё работает. Кластер самоидентифицируется.

Производительность:
База, 100 000 записей, один запрос на чтение приходит за 0.6 ms
Шустро.

четверг, 4 октября 2012 г.

Java, RESTful Web Service.

В прошлом посте, для тестов использовался SOAP Web Service (javax.WS)
Время одного вызова в синхронном режиме занимало 43 ms.
В этот раз я протестировал RESTful / JSON Web Service (javax.RS), время обработки одного запроса, катастрофично упало до 2 ms.

SOAP всетаки тяжелая технология, подходит для тех, кому некуда спешить.

Как создать JSON Web Service
http://www.mkyong.com/webservices/jax-rs/restful-java-client-with-jersey-client/
http://www.mkyong.com/webservices/jax-rs/json-example-with-jersey-jackson/

среда, 3 октября 2012 г.

Java, Web Service и Асинхронность

      Совсем не давно возникла необходимость создать ассинхронный веб сервис. Я выяснил, что на самом деле асинхронность, это нечто чужеродное  для веб сервисов. Поэтому, если хотите асинхронности придется по извращаться. Ниже я опишу решения, которые удалось раскопать.

1. Асинхронный клиент для синхронного сервиса.
Данный подход, заключается в эмуляции асинхронности на стороне клиента.
Реализация хорошо описана здесь http://www.ibm.com/developerworks/ru/library/wes-0804_sedov/

От себя  добавлю, что необходимо, с помощью утилиты wsimport, создать обвертку из Java классов. Которую мы будем использовать в клиенте, для коммуникации с WS.

Вот так, можно сформировать обычную синхронную обвертку
wsimport -keep -verbose http://localhost:8080/AeroSpike?wsdl

Но нам нужны асинхронные методы. Для этого запускаем утилиту, с такими параметрами
wsimport -b binding.xml -keep -verbose http://localhost:8080/AsyncSpikeServiceImpl?wsdl

где binding.xml файл биндинга, следующего содержания

<bindings
    xmlns:xsd="http://www.w3.org/2001/XMLSchema"
    xmlns:wsdl="http://schemas.xmlsoap.org/wsdl/"
    wsdlLocation="http://localhost:8080/AsyncSpikeServiceImpl?wsdl"
    xmlns="http://java.sun.com/xml/ns/jaxws">
    <bindings node="wsdl:definitions">
         <enableAsyncMapping>true</enableAsyncMapping>
    </bindings>
</bindings>

В качестве wsdlLocation, может быть указан как URL, так и WSDL файл.
При этом, если у веб сервиса есть методы, помеченные аннотацией OneWay (т.е метод вызываемый в не блокирующем режиме), для таких методов асинхронные обвертки создаваться не будут.

Производительность
Сделал замеры производительности на клиенте. В синхронном режиме, на  один запрос уходит 43 ms. В асинхронном (с обратным вызовом AsyncHandler)  84 ms.



2. Дуплексный подход. Клиент - он же сервис обратного вызова.

Это самый правильный подход, базирующийся на спецификации WS-Addressing. С помощью него, можно организовать самую настоящую асинхронность. Выглядит это так:


 Веб сервису приходит запрос, в заголовке которого содержится адрес, на который нужно отослать ответ. Можно сейчас, а можно через неделю.
Одна беда, поддержка WS-Adressing есть не во всех контейнерах сервлетов, и существуют трудности с генерацией правильной обвертки. Как это сделать в Jetty, я так и не нашел.

Суть вопроса, хорошо раскрыта здесь, применительно к Oracle SOA Suite:
http://samolisov.blogspot.com/2012/05/blog-post.html

и здесь, для WebSphere
http://www.ibm.com/developerworks/ru/library/ws-JAXsupport/

А это для WebLogic:
http://docs.oracle.com/cd/E17904_01/web.1111/e13734/asynch.htm

А здесь просто много теории
http://docs.oracle.com/cd/E17904_01/web.1111/e13734/asynch.htm