вторник, 4 марта 2014 г.

Сериализация java bean в JSON

Давно задумывался над сериализацией JSON в java сущности и обратно.
Чтобы красиво, грациозно, и главное - производительно. Еще одно требование, чтобы сериализация производилась в строки, это наиболее универсальный формат для межсистемного взаимодействия. Пo этой причине Protobuf не рассматривал. Как правило использовал org.codehaus.jettison.json.JSONObject и ручную сериализацию. Внутри него используется LinkedHashMap поэтому операции с полями происходят быстро.

public class Geo
{

    private static final String LAT_FIELD = "lat";
    private static final String LON_FIELD = "lon";

    private float lat;
    private float lon;

    public String toJSON()
    {
        JSONObject json = new JSONObject();

        try {
            json.put(LAT_FIELD, lat);
            json.put(LON_FIELD, lon);
        } catch (JSONException e) {
            LOG.error("Cannot convert Geo object into json.", e);
        }

        return json.toString();
    }


    public static Geo fromJSON(final String json) throws JSONException
    {
        JSONObject jsonObject = new JSONObject(json);
        Geo geo = new Geo();

        geo.setLon((float) jsonObject.optDouble(LON_FIELD, Float.NaN));
        geo.setLat((float) jsonObject.optDouble(LAT_FIELD, Float.NaN));

        return geo;
    }
}


Надо ли говорить, что при большом количестве полей, такой подход адски неудобен.
Не красиво и трудозатратно, зато сносно, с точки зрения производительности, плюс поддержка валидации.
Всегда знал о такой штуке как com.google.gson.Gson но обходил его стороной, т.к он использует рефлексию, казалось что это снизит производительность.
Но красота использования Gson,

public class Geo
{
    private float lat;
    private float lon;

    public String toJSON()
    {
        Gson gson = new Gson();
        return gson.toJson(this);
    }

    public static Geo fromJSON(final String json)
    {
        Gson gson = new Gson();
        return gson.fromJson(json, Geo.class);
    }
}

...и безобразность ручной сериализации не давали покоя.
Решил на днях замерить производительность этих двух подходов.
Всё как всегда, цикл на миллион, и время в наносекундах.

Результаты:

Сериализация сущности в строку
Gson              4088893178
JSONObject  4543766682

Десериализация из строки:
Gson              2728182025
JSONObject  4180603421

Приятно удивлен производительностью Gson, оказывается в некоторых случаях она в полтора раза выше. Радости нет предела. Теперь буду пользоваться им.

пятница, 29 ноября 2013 г.

Инструкция по установке Twitter Storm

Storm:  distributed and fault-tolerant realtime computation
По горячим следам, пока это еще в голове. На первый взгляд всё просто, но поприсидать пришлось.

Последовательность такая:

1.  Ставим Zookeeper
2.  Ставим зависимости для Nimbus и рабочих машин
3.  Скачиваем и ставим Storm на Nimbus и рабочие машины
4.  Правим конфиг Шторма  storm.yaml
5.  Запускаем демонов

Настраиваем мастер (Nimbus):


  • Zookeeper
    Для тестового кластера хватит одной ноды. За одно поставим утилиты и библиотеки, необходимые для Nimbus

apt-get install openjdk-6-jdk zookeeper make build-essential
apt-get install uuid-dev unzip pkg-config libtool automake

  • apache
    Для UI нужен веб сервер
apt-get install apache2


Далее настройка Нимбуса и рабочих нод одинакова.
Ставим зависимости. Storm требует ZeroMQ 2.1.7 и JZMQ.
  •  ZerroMQ 2.1.7Нужна версия 2.1.7 ибо 2.1.10 глюкава

sudo aptitude install build-essential
cd ~
wget http://download.zeromq.org/zeromq-2.1.7.tar.gz
tar zxvf zeromq-2.1.7.tar.gz
cd zeromq-2.1.7
./configure
make
sudo make install
sudo ldconfig


  • JZMQ
    С ним пришлось повозиться. Ниже приведен самый рабочий способ:

    Будет нужен GIT 
sudo apt-get install git
git clone --depth 1 https://github.com/nathanmarz/jzmq.git
cd jzmq
./autogen.sh

export JAVA_HOME=/usr/lib/jvm/java-1.6.0-openjdk-amd64

./configure
touch src/classdist_noinst.stamp
cd src/

CLASSPATH=.:./.:$CLASSPATH javac -d . org/zeromq/ZMQ.java org/zeromq/App.java org/zeromq/ZMQForwarder.java org/zeromq/EmbeddedLibraryTools.java org/zeromq/ZMQQueue.java org/zeromq/ZMQStreamer.java org/zeromq/ZMQException.java

cd ..
make
sudo make install

  • Storm     
    Собственной персоны
wget https://github.com/downloads/nathanmarz/storm/storm-0.8.1.zip
unzip storm-0.8.1.zip
sudo mkdir /mnt/storm
sudo chmod 777 /mnt/storm

         правим conf/storm.yaml
storm.zookeeper.servers:
- "111.222.333.444"
- "555.666.777.888"

storm.local.dir: "/mnt/storm"

nimbus.host: "111.222.333.44"

supervisor.slots.ports:
- 6700
- 6701
- 6702
- 6703

Запускаем:
  • Nimbus на мастере "bin/storm nimbus"
  • Supervisor на каждой рабочей машине "bin/storm supervisor"
  • UI на мастере "bin/storm ui" http://{nimbus host}:8080

вторник, 24 сентября 2013 г.

HBase загрузка больших массивов данных через bulk-load

Привет коллеги.
Хочу поделиться своим опытом использования HBase, а именно рассказать про bulk loading. Это еще один метод загрузки данных. Он принципиально отличается от обычного подхода (записи в таблицу через клиента). Есть мнение, что с помощью bulk load можно очень быстро загружать огромные массивы данных. Именно в этом я решил разобраться.

Итак, обо всём по порядку. Загрузка через bulk load происходит в три этапа:
  1. Помещаем файлы с данными в HDFS
  2. Запускаем MapReduce задачу, которая преобразует исходные данные непосредственно в файлы формата  HFile, посути HBase хранит свои данные именно в таких файлах.
  3. Запускаем bulk load функцию, которая зальёт (привяжет) полученные файлы в таблицу HBase.




В данном случае мне было необходимо прочувствовать эту технологию и понять её в цифрах: чему равна скорость, как она зависит от количества и размера файлов. Эти числа слишком зависимы от внешних условий, но помогают понять порядки между обычной загрузкой и bulk load.

Исходные данные:

  • Кластер под управлением Cloudera CDH4, HBase 0.94.6-cdh4.3.0.
    Три виртуальных хоста (на гипервизоре), в конфигурации CentOS / 4CPU / RAM 8GB / HDD 50GB
  • Тестовые данные хранились в CSV файлах различных размеров, суммарным объёмом 2GB, 3.5GB, 7.1GB и 14.2GB

Сначала о результатах:


Bulk loading
суммарный размер данных
количество файлов
суммарное количество записей (rows)
количество map
количество reduce
время задачи (Job, сек)
время общее (сек)
2GB
16
4 000 000
16
10
88
160
2GB
100
4 000 000
100
10
137
207
3.5GB
28
7 000 000
28
10
120
191
3.5GB
100
7 000 000
100
10
183
253
3.5GB
1
7 000 000
28
10
123
192
7.1GB
100
14 000 000
100
10
314
380
7.1GB
1
14 000 000
55
10
258
330
14.2G
1
28 000 000
109
10
510
583

Cкорость:
  • Max 29.2 Mb/sec или 58K rec/sec (3.5GB в 28 файлах)
  • Average 27 Mb/sec или 54K rec/sec(рабочая скорость, более приближенная к реальности )
  • Min  14.5 Mb/sec или 29K rec/sec (2GB в 100 файлах)
  • 1 файл загружается на 20% быстрее чем 100

Размер одной записи (row): 0.5Кb
Время инициализации MapReduce Job: 70 sec
Время загрузки файлов в HDFS с локальной файловой системы:
  • 3.5GB / 1 файл - 65 sec 
  • 7.5GB / 100 - 150 sec 
  • 14.2G / 1 файл - 285 sec

Загрузка через клиенты:

Загрузка осуществлялась с 2-х хостов по 8 потоков на каждом.
Клиенты запускались по крону в одно и тоже время, загрузка CPU не превышала 40%
Размер одной записи (row), как и в предыдущем случае был равен 0.5Кb.

количество записей (rows)
среднее время (сек)
4 000 000
109
7 000 000
208
14 000 000
380
28 000 000
836

Что в итоге?


Реализовать этот тест, я решил на волне разговоров о bulk load, как о способе сверхбыстрой загрузки данных. Стоит сказать, что в официальной документации речь идет только о снижении нагрузки на сеть и CPU. Как бы там ни было,  я не вижу выигрыша в скорости. Тесты показывают что bulk load быстре всего лишь в полтора раза, но не будем забывать, что это без учета инициализации m/r джобы. Кроме того, данные надо доставить в HDFS, на это тоже потребуется какое то время.
Думаю, стоит относиться к bulk load просто, как к еще одному способу загрузки данных, архитектурно иному (в некоторых случаях очень даже удобному).

А теперь о реализации

Теоретически всё довольно просто, но на практике возникает несколько технических нюансов.

//Создаём джоб
Job job = new Job(configuration, JOB_NAME);
job.setJarByClass(BulkLoadJob.class);

job.setMapOutputKeyClass(ImmutableBytesWritable.class);
job.setMapOutputValueClass(Put.class);

job.setMapperClass(DataMapper.class);
job.setNumReduceTasks(0);

job.setInputFormatClass(TextInputFormat.class);
job.setOutputFormatClass(HFileOutputFormat.class);

FileInputFormat.setInputPaths(job, inputPath);
HFileOutputFormat.setOutputPath(job, new Path(outputPath));

HTable dataTable = new HTable(jobConfiguration, TABLE_NAME);
HFileOutputFormat.configureIncrementalLoad(job, dataTable);

//Запускаем
ControlledJob controlledJob = new ControlledJob(
    job,
    null
);

JobControl jobController = new JobControl(JOB_NAME);
jobController.addJob(controlledJob);

Thread thread = new Thread(jobController);
thread.start();
.
.
.
//Даём права на output
setFullPermissions(JOB_OUTPUT_PATH);

//Запускаем функцию bulk-load
LoadIncrementalHFiles loader = new LoadIncrementalHFiles(jobConfiguration);
loader.doBulkLoad(
        new Path(JOB_OUTPUT_PATH),
        dataTable
);

  • MapReduce Job создаёт выходные файлы c правами пользователя, от имени которого он был запущен.
  • bulk load всегда запускается от имени пользователя hbase, поэтому не может прочитать подготовленные для него файлы, и валится вот с таким исключением:
    org.apache.hadoop.security.AccessControlException: Permission denied: user=hbase
Поэтому надо запускать Job от имени пользователя hbase или раздать права на выходные файлы (именно так я сделал).
  • Необходимо правильно создать таблицу HBase. По умолчанию она создается с одним Region-ом. Это приводит к тому, что создается только один редьюсер и запись идет только на одну ноду, загружая её на 100%, остальные при этом курят.
    Поэтому при создании новой таблицы надо сделать pre-split. В моём случае таблица разбивалась на 10 Region-ов равномерно разбросанных по всему кластеру.

    //Создаём таблицу и делаем пре-сплит
    HTableDescriptor descriptor = new HTableDescriptor(
            Bytes.toBytes(tableName)
    );
    
    descriptor.addFamily(
            new HColumnDescriptor(Constants.COLUMN_FAMILY_NAME)
    );
    
    HBaseAdmin admin = new HBaseAdmin(config);
    
    byte[] startKey = new byte[16];
    Arrays.fill(startKey, (byte) 0);
    
    byte[] endKey = new byte[16];
    Arrays.fill(endKey, (byte)255);
    
    admin.createTable(descriptor, startKey, endKey, REGIONS_COUNT);
    admin.close();
    

    • MapReduce Job пишет в выходную директорию, которую мы ему указываем, но при этом создает  субдиректории, одноименные с column family. Файлы создаются именно там.
    В целом, это всё. Хочется сказать, что это довольно грубый тест, без хитрых оптимизаций, поэтому если у вас есть что добавить, буду рад услышать.

    Весь код проекта доступен на GitHub: https://github.com/2anikulin/hbase-bulk-load


    воскресенье, 25 августа 2013 г.

    Cloudera oozie WebUI error


    Не заботает web консоль oozie в Клаудере.
    Ругается:

    Oozie web console is disabled.
    To enable Oozie web console install the Ext JS library.

    ставим библиотеку:  

    wget 'http://extjs.com/deploy/ext-2.2.zip'
    unzip ext-2.2.zip
    sudo cp ext-2.2 /usr/lib/oozie/libext
    даём права:
    sudo chmod -R ugo+rwx /usr/lib/oozie/libext/ext-2.2

    прописываем симлинк:
    cd /var/lib/oozie/
    sudo ln -s /usr/lib/oozie/libext/ext-2.2 ext-2.2

    перезапускаем oozie

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

    вторник, 16 июля 2013 г.

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

    Hadoop MapReduce "Can't read partitions file"

    Во время выполнения джобы, падает вот такое исключение "Can't read partitions file"


    13/07/15 18:48:42 WARN mapred.LocalJobRunner: job_local38891965_0001
    java.lang.Exception: java.lang.IllegalArgumentException: Can't read partitions file
    at org.apache.hadoop.mapred.LocalJobRunner$Job.run(LocalJobRunner.java:404)
    Caused by: java.lang.IllegalArgumentException: Can't read partitions file
    at org.apache.hadoop.mapreduce.lib.partition.TotalOrderPartitioner.setConf(TotalOrderPartitioner.java:108)
    at org.apache.hadoop.util.ReflectionUtils.setConf(ReflectionUtils.java:70)
    at org.apache.hadoop.util.ReflectionUtils.newInstance(ReflectionUtils.java:130)
    at org.apache.hadoop.mapred.MapTask$NewOutputCollector.<init>(MapTask.java:588)
    at org.apache.hadoop.mapred.MapTask.runNewMapper(MapTask.java:657)
    at org.apache.hadoop.mapred.MapTask.run(MapTask.java:331)
    at org.apache.hadoop.mapred.LocalJobRunner$Job$MapTaskRunnable.run(LocalJobRunner.java:266)
    at java.util.concurrent.Executors$RunnableAdapter.call(Executors.java:441)
    at java.util.concurrent.FutureTask$Sync.innerRun(FutureTask.java:303)
    at java.util.concurrent.FutureTask.run(FutureTask.java:138)
    at java.util.concurrent.ThreadPoolExecutor$Worker.runTask(ThreadPoolExecutor.java:886)
    at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:908)
    at java.lang.Thread.run(Thread.java:662)
    Caused by: java.io.FileNotFoundException: File file:/home/anikulin/_partition.lst does not exist
    at org.apache.hadoop.fs.RawLocalFileSystem.getFileStatus(RawLocalFileSystem.java:468)
    at org.apache.hadoop.fs.FilterFileSystem.getFileStatus(FilterFileSystem.java:373)
    at org.apache.hadoop.io.SequenceFile$Reader.<init>(SequenceFile.java:1704)
    at org.apache.hadoop.io.SequenceFile$Reader.<init>(SequenceFile.java:1728)
    at org.apache.hadoop.mapreduce.lib.partition.TotalOrderPartitioner.readPartitions(TotalOrderPartitioner.java:293)
    at org.apache.hadoop.mapreduce.lib.partition.TotalOrderPartitioner.setConf(TotalOrderPartitioner.java:80)
    ... 12 more


    Это значит что наш кластер функционирует в режиме "pseudo distributed mode"
    Хотя никто его об этом не просил!
    Чинится прописыванием адреса JobTracker в mapred-site.xml

    В Клоудеровском исполнении это делается из консоли, добавлением свойств в
    "MapReduce Service Configuration Safety Valve for mapred-site.xml"

    <property>
           <name>mapred.job.tracker</name>
           <value>categorizer-hadoop-1.mywork:8021</value>
    </property>


    http://archive.cloudera.com/cdh4/cdh/4/hbase/book.html