понедельник, 24 июня 2013 г.

Переключаемся между разными JVM

Иногда требуется переключиться с одной версии джавы на другую.
Для этого есть парочка полезных команд. Запускаем и выбираем нужное...

sudo update-alternatives --config java
To update the Java compiler run:
sudo update-alternatives --config javac

среда, 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

четверг, 18 апреля 2013 г.

Unsupported major.minor version 51.0

Часто, в последнее время сталкиваюсь с таким исключением. Дабы не лазить на StackOverflaw - напишу здесь Это значит, что подключаемые внешние jar файлы, собраны для более новой версии java. А мы пытаемся запуститься на старой. Обычно речь идет о JRE 6 -7

Exception in thread "main" java.lang.UnsupportedClassVersionError: org/eclipse/jetty/server/Handler : Unsupported major.minor version 51.0
at java.lang.ClassLoader.defineClass1(Native Method)
at java.lang.ClassLoader.defineClassCond(ClassLoader.java:631)
at java.lang.ClassLoader.defineClass(ClassLoader.java:615)
at java.security.SecureClassLoader.defineClass(SecureClassLoader.java:141)
at java.net.URLClassLoader.defineClass(URLClassLoader.java:283)
at java.net.URLClassLoader.access$000(URLClassLoader.java:58)
at java.net.URLClassLoader$1.run(URLClassLoader.java:197)
at java.security.AccessController.doPrivileged(Native Method)
at java.net.URLClassLoader.findClass(URLClassLoader.java:190)
at java.lang.ClassLoader.loadClass(ClassLoader.java:306)
at sun.misc.Launcher$AppClassLoader.loadClass(Launcher.java:301)
at java.lang.ClassLoader.loadClass(ClassLoader.java:247)

суббота, 6 апреля 2013 г.

VirtualBox проблемы после обновления linux kernel

Почти после каждого обновления Ubuntu, отваливается VirtualBox.
Выкидывает, такое сообщение:
Kernel driver not installed (rc=-1908)
The VirtualBox Linux kernel driver (vboxdrv) is either not loaded or there is a permission problem with /dev/vboxdrv. Re-setup the kernel module by executing
'/etc/init.d/vboxdrv setup'
as root. Users of Ubuntu, Fedora or Mandriva should install the DKMS package first. This package keeps track of Linux kernel changes and recompiles the vboxdrv kernel module if necessary.
Напрягает!
Чинится, вот так:
sudo apt-get install linux-headers-`uname -r`
sudo apt-get remove dkms  
sudo apt-get install dkms virtualbox-dkms  
sudo modprobe vboxdrv



четверг, 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> {...


пятница, 22 февраля 2013 г.

MongoDB map-reduce, интересные нюансы

Задача:
Необходимо сгруппировать объекты по атрибуту source_name, и посчитать количество таких объектов для каждой группы.

Решение:
В MongoDB есть Aggregation framework, он позволяет в том числе запускать m/r задачи.
Надо сказать что в Монге какой то странный мэп-редьюс. Во первых он однопоточный, что противоречит самой парадигме m/r, а во вторых с редьюсами происходит что то не понятное, что запутывает окончательно.

Первое решение (логичное, но не правильное) было таким:
Всё как обычно, каждому ключу ставим значение 1, в редьюсе подсчитываем их количество
db.buffer.mapReduce(
    "function() {
         emit(this.source_name, 1);
    }",
    "function(k, vals) {
         var count = vals.length;
         return {count: count};
    }",
    {out: "buffer_results"}
);
out:
{
    "result" : "buffer_results",
    "timeMillis" : 70258,
    "counts" : {
        "input" : 2803200,
        "emit" : 2803200,
        "reduce" : 166542,
        "output" : 13
     },
   "ok" : 1,
}

{ "_id" : "source_A", "value" : { "count" : 2601 } }
{ "_id" : "source_B", "value" : { "count" : 50456 } }
{ "_id" : "source_C", "value" : { "count" : 23422 } }
{ "_id" : "source_D", "value" : { "count" : 725523 } }
{ "_id" : "source_E", "value" : { "count" : 234 } }
{ "_id" : "source_F", "value" : { "count" : 123111 } }
{ "_id" : "source_G", "value" : { "count" : 4 } }
{ "_id" : "source_H", "value" : { "count" : 1 } }
{ "_id" : "source_I", "value" : { "count" : 1 } }
{ "_id" : "source_J", "value" : { "count" : 7363 } }
{ "_id" : "source_K", "value" : { "count" : 560 } }
{ "_id" : "source_L", "value" : { "count" : 10 } }
Если суммировать все count, должно получиться  2803200, но не получается. Такое впечатление, что это результат какого то не полного редьюса.

И действительно,  выясняется интересный нюанс. Редьюс может быть несколько раз вызван для одного и того же ключа. Невероятно, но факт! Это связанно с архитектурой Монги, m/r может выполняться на кластере,  документы с одним и тем же ключом могут находится на разных серверах. Сначала редьюс по ключу выполняется на одном сервере, а затем на другом.
После чего происходит финальный редьюс, в который попадают результаты выполнения первых двух. Т.е входом в редьюс могут быть как значения выхода мэпа, так и редьюса.



Теперь всё встало на свои места. В сложившейся ситуации, код m/r должен быть таким:
db.buffer.mapReduce(
   "function() {
       emit(this.source_name, {count: 1});
   }",
   "function(key, values) {
      var count = 0; 
      values.forEach(function(v) {count += v['count'];}); 
      return {count: count};
   }", 
   {out: "buffer_results"}
);
Мэп и редьюс генерируют одинаковые объекты на выходе. Поэтому редьюсу без разницы, от кого пришли данные.
out:
{
    "result" : "buffer_results",
    "timeMillis" : 70258,
    "counts" : {
            "input" : 2803200,
                "emit" : 2803200,
            "reduce" : 166542,
            "output" : 13
    },
    "ok" : 1,
}

{ "_id" : "source_A", "value" : { "count" : 2601 } }

{ "_id" : "source_B", "value" : { "count" : 1380432 } }
{ "_id" : "source_C", "value" : { "count" : 106087 } }
{ "_id" : "source_D", "value" : { "count" : 854294 } }
{ "_id" : "source_E", "value" : { "count" : 100509 } }
{ "_id" : "source_F", "value" : { "count" : 151713 } }
{ "_id" : "source_G", "value" : { "count" : 4 } }
{ "_id" : "source_H", "value" : { "count" : 1 } }
{ "_id" : "source_I", "value" : { "count" : 1 } }
{ "_id" : "source_J", "value" : { "count" : 72363 } }
{ "_id" : "source_K", "value" : { "count" : 2560 } }
{ "_id" : "source_L", "value" : { "count" : 10 } }

В итоге получаем правильные результаты. Дебит сходится с кредитом.






Фтопку !

Здесь буду перечислять всякие неприятности, о которых стоит помнить в следующий раз...

-------------------------------------------------------------------------------------------
au.com.bytecode.opencsv.CSVReader в топку !

Жутко заглючил у заказчика.
Если от парсера не требуется мэппинг в POJO, то лучше свой написать.
Всего то 10 строк кода.
-------------------------------------------------------------------------------------------