Експеримент за проверка на приложимостта на графовата СУБД JanusGraph за решаване на задачата за намиране на подходящи пътища

Експеримент за проверка на приложимостта на графовата СУБД JanusGraph за решаване на задачата за намиране на подходящи пътища

Здравейте на всички. Разработваме продукт за анализ на офлайн трафик. В проекта има задача, свързана със статистическия анализ на пътищата на движение на посетителите в определени области.

В рамките на тази задача потребителите могат да задават на системата запитвания от следния вид:

  • колко посетители са преминали от област "A" в област "Б";
  • колко посетители са преминали от област "A" в област "Б" през област "C", а след това през област "Д";
  • колко време е отнело преминаването на посетител от определен тип от област "А" в област "Б".

и още редица подобни аналитични запитвания.

Движението на посетителите в областите представлява насочен граф. Четейки в интернет, открих, че графовите бази данни се използват и за аналитични отчети. Появи се желание да видя как ще се справят с подобни запитвания графовите бази данни (TL;DR; лошо).

Избрах да използвам базата данни JanusGraph, като изключителен представител на графовите open-source бази данни, която се основава на стек от зрели технологии, които (по мое мнение) би трябвало да ѝ осигурят прилични оперативни характеристики:

  • бекенд хранилище BerkeleyDB, Apache Cassandra, Scylla;
  • сложни индекси могат да се съхраняват в Lucene, Elasticsearch, Solr.

Авторите на JanusGraph пишат, че тя е подходяща както за OLTP, така и за OLAP.

Работил съм с BerkeleyDB, Apache Cassandra, Scylla и ES, освен това, тези продукти често се използват в нашите системи, затова с оптимизъм подхождах към тестването на тази графова база данни. Странно ми се стори изборът на BerkeleyDB, а не на RocksDB, но вероятно това е свързано с изискванията за транзакции. Във всеки случай, за мащабируемо, продуктовото използване се препоръчва да се използва бекенд с Cassandra или Scylla.

Не разглеждах Neo4j, тъй като за клъстеризация е необходима търговска версия, т.е. продуктът не е отворен.

Графовите бази данни казват: "Ако нещо изглежда като граф — обработвайте го като граф!" — красота!

Първо нарисувах граф, който е направен точно по каноните на графовите бази данни:

Експеримент за проверка на приложимостта на графовата СУБД JanusGraph за решаване на задачата за намиране на подходящи пътища

Има сущност Zone, отговаряща за областта. Ако ZoneStep принадлежи на тази Zone, то той я посочва. По сущностите Area, ZoneTrack, Person не обръщайте внимание, те принадлежат на домейна и в рамките на теста не се разглеждат. В крайна сметка, за такава графова структура, запитването за търсене на вериги изглеждаше така:

g.V().hasLabel('Zone').has('id',0).in_()
       .repeat(__.out()).until(__.out().hasLabel('Zone').has('id',19)).count().next()

На български това звучи така: намери Zone с ID=0, вземи всички върхове, от които тръгва ребро (ZoneStep), продължавай напред, без да се връщаш назад, докато не намериш ZoneStep, от който тръгва ребро към Zone с ID=19, преброй броя на такива вериги.

Не претендирам за познания по всичките нюанси на графовото търсене, но този запит е генериран на база на тази книга (https://kelvinlawrence.net/book/Gremlin-Graph-Guide.html).

Заредих 50 хиляди трака с дължина от 3 до 20 точки в графовата база JanusGraph, използваща бекенд BerkeleyDB, създадох индекси съгласно ръководството.

Скрипт за зареждане на Python:


from random import random
from time import time

from init import g, graph

if __name__ == '__main__':

    points = []
    max_zones = 19
    zcache = dict()
    for i in range(0, max_zones + 1):
        zcache[i] = g.addV('Zone').property('id', i).next()

    startZ = zcache[0]
    endZ = zcache[max_zones]

    for i in range(0, 10000):

        if not i % 100:
            print(i)

        start = g.addV('ZoneStep').property('time', int(time())).next()
        g.V(start).addE('belongs').to(startZ).iterate()

        while True:
            pt = g.addV('ZoneStep').property('time', int(time())).next()
            end_chain = random()
            if end_chain < 0.3:
                g.V(pt).addE('belongs').to(endZ).iterate()
                g.V(start).addE('goes').to(pt).iterate()
                break
            else:
                zone_id = int(random() * max_zones)
                g.V(pt).addE('belongs').to(zcache[zone_id]).iterate()
                g.V(start).addE('goes').to(pt).iterate()

            start = pt

    count = g.V().count().next()
    print(count)

Използваше се ВМ с 4 ядра и 16 GB RAM на SSD. JanusGraph беше разгорнат с следната команда:

docker run --name janusgraph -p8182:8182 janusgraph/janusgraph:latest

В този случай данните и индексите, които се използват за търсене по точно съвпадение, се съхраняват в BerkeleyDB. След като изпълних горепосоченото запитване, получих време, равно на няколко десетки секунди.

Изпълнявайки 4 от горепосочените скриптове паралелно, успях да превърна СУБД в тиква с весела вълна от Java стектрейсове (а ние всички обичаме да четем Java стектрейсове) в Docker логовете.

След размисъл реших да опростя схемата на графа до следното:

Експеримент за проверка на приложимостта на графовата СУБД JanusGraph за решаване на задачата за намиране на подходящи пътища

Реших, че търсенето по атрибутите на сущността ще бъде по-бързо, отколкото търсенето по ребрата. В крайна сметка моето запитване се преобразува в следното:

g.V().hasLabel('ZoneStep').has('id',0).repeat(__.out().simplePath()).until(__.hasLabel('ZoneStep').has('id',19)).count().next()

На български това звучеше така: намери ZoneStep с ID=0, продължавай напред, без да се връщаш назад, докато не намериш ZoneStep с ID=19, преброй броя на такива вериги.

Скриптът за зареждане, посочен по-горе, също опростих, за да не създавам излишни връзки, ограничявайки се до атрибутите.

Запитът все още се извършваше няколко секунди, което беше напълно неприемливо за нашата задача, тъй като за целите на AdHoc запитванията произволен тип, това изобщо не беше подходящо.

Опитах се да разширя JanusGraph, използвайки Scylla, като най-бързото реализиране на Cassandra, но това също не доведе до съществени промени в производителността.

Така че, въпреки че "изглежда като граф", не успях да накарам графовата база данни да обработи това бързо. Пълно е възможно да не знам нещо и да мога да накарам JanusGraph да извърши това търсене за части от секундата, обаче, това не ми се удари.

Тъй като задачата все пак трябваше да бъде решена, започнах да мисля за JOIN операции и Pivot таблици, което не вдъхваше оптимизъм от гледна точка на елегантност, но можеше да се окаже напълно работещ вариант на практика.

В нашия проект вече се използва Apache ClickHouse, затова реших да проверя своите изследвания на тази аналитична база данни.

Разширих ClickHouse по простата рецепта:

sudo docker run -d --name clickhouse_1 
     --ulimit nofile=262144:262144 
     -v /opt/clickhouse/log:/var/log/clickhouse-server 
     -v /opt/clickhouse/data:/var/lib/clickhouse 
     yandex/clickhouse-server

Създадох в него база данни и таблица от вида:

CREATE TABLE 
db.steps (`area` Int64, `when` DateTime64(1, 'Europe/Moscow') DEFAULT now64(), `zone` Int64, `person` Int64) 
ENGINE = MergeTree() ORDER BY (area, zone, person) SETTINGS index_granularity = 8192

Попълних я с данни със следния скрипт:

from time import time

from clickhouse_driver import Client
from random import random

client = Client('vm-12c2c34c-df68-4a98-b1e5-a4d1cef1acff.domain',
                database='db',
                password='secret')

max = 20

for r in range(0, 100000):

    if r % 1000 == 0:
        print("CNT: {}, TS: {}".format(r, time()))

    data = [{
            'area': 0,
            'zone': 0,
            'person': r
        }]

    while True:
        if random() < 0.3:
            break

        data.append({
                'area': 0,
                'zone': int(random() * (max - 2)) + 1,
                'person': r
            })

    data.append({
            'area': 0,
            'zone': max - 1,
            'person': r
        })

    client.execute(
        'INSERT INTO steps (area, zone, person) VALUES',
        data
    )

Тъй като вставките стават на партиди, попълването беше много по-бързо, отколкото при JanusGraph.

Конструирах два запита с помощта на JOIN. За преминаване от точка A в точка Б:

SELECT s1.person AS person,
       s1.zone,
       s1.when,
       s2.zone,
       s2.when
FROM
  (SELECT *
   FROM steps
   WHERE (area = 0)
     AND (zone = 0)) AS s1 ANY INNER JOIN
  (SELECT *
   FROM steps AS s2
   WHERE (area = 0)
     AND (zone = 19)) AS s2 USING person
WHERE s1.when <= s2.when

За преминаване през 3 точки:

ИЗБОР s3.person,
       s1z,
       s1w,
       s2z,
       s2w,
       s3.zone,
       s3.when
ИЗ
  (ИЗБРАТЬ s1.person AS person,
          s1.zone AS s1z,
          s1.when AS s1w,
          s2.zone AS s2z,
          s2.when AS s2w
   ИЗ
     (ИЗБРАТЬ *
      ИЗ steps
      ГДЕ (area = 0)
        И (zone = 0)) AS s1 ANY INNER JOIN
     (ИЗБРАТЬ *
      ИЗ steps AS s2
      ГДЕ (area = 0)
        И (zone = 3)) AS s2 ИСПОЛЬЗОВАТЬ person
   ГДЕ s1.when <= s2.when) p ANY INNER JOIN
  (ИЗБРАТЬ *
   ИЗ steps
   ГДЕ (area = 0)
     И (zone = 19)) AS s3 ИСПОЛЬЗОВАТЬ person
ГДЕ p.s2w <= s3.when

Запросите, конечно, выглядят довольно сложными. Для реального использования требуется программная обвязка-генератор. Тем не менее, они функционируют и очень быстро. Первый и второй запросы выполняются менее чем за 0.1 секунды. Вот пример времени исполнения запроса для count(*) по трем точкам:

ИЗБРАТЬ count(*)
ИЗ 
(
    ИЗБРАТЬ 
        s1.person AS person, 
        s1.zone AS s1z, 
        s1.when AS s1w, 
        s2.zone AS s2z, 
        s2.when AS s2w
    ИЗ 
    (
        ИЗБРАТЬ *
        ИЗ steps
        ГДЕ (area = 0) И (zone = 0)
    ) AS s1
    ANY INNER JOIN 
    (
        ИЗБРАТЬ *
        ИЗ steps AS s2
        ГДЕ (area = 0) И (zone = 3)
    ) AS s2 ИСПОЛЬЗОВАТЬ (person)
    ГДЕ s1.when <= s2.when
) AS p
ANY INNER JOIN 
(
    ИЗБРАТЬ *
    ИЗ steps
    ГДЕ (area = 0) И (zone = 19)
) AS s3 ИСПОЛЬЗОВАТЬ (person)
ГДЕ p.s2w <= s3.when

┌─count()─┐
│   11592 │
└─────────┘

1 строка в результате. Затраченное время: 0.068 секунд. Обработано 250.03 тысяч строк, 8.00 МБ (3.69 миллионов строк/с., 117.98 МБ/с.)

Замечание по IOPS. При заполнении данных, JanusGraph генерировал довольно высокое количество IOPS (1000-1300 для четырех потоков заполнения данными), в то время как IOWAIT был достаточно высоким. Тем не менее, ClickHouse создал минимальную нагрузку на дисковую подсистему.

Заключение

Мы решили использовать ClickHouse для обработки запросов этого типа. Мы всегда можем дополнительно оптимизировать запросы, используя материализованные представления и параллелизацию, предварительно обрабатывая поток событий с помощью Apache Flink перед загрузкой их в ClickHouse.

Производительность настолько хороша, что, вероятно, нам даже не придется думать о программных пивотах таблиц. Ранее нам пришлось делать пивоты данных, извлекаемых из Vertica через выгрузку в Apache Parquet.

К сожалению, очередная попытка использования графовой СУБД не удалась. Я не нашел, что у JanusGraph дружелюбная экосистема, которая бы позволила быстро освоиться с продуктом. При этом для конфигурирования сервера применяется традиционный Java-метод, который заставит людей, незнакомых с Java, плакать кровавыми слезами:

хост: 0.0.0.0
порт: 8182
пулПотоков: 1
пулГрэмлина: 8
таймаутОценкиСкрипта: 30000
канализатор: org.janusgraph.channelizers.JanusGraphWsAndHttpChannelizer

менеджерГрафов: org.janusgraph.graphdb.management.JanusGraphManager
graphs: {
  ConfigurationManagementGraph: conf/janusgraph-cql-configurationgraph.properties,
  airlines: conf/airlines.properties
}

движкиСкриптов: {
  gremlin-groovy: {
    плагины: { org.janusgraph.graphdb.tinkerpop.plugin.JanusGraphGremlinPlugin: {},
               org.apache.tinkerpop.gremlin.server.jsr223.GremlinServerGremlinPlugin: {},
               org.apache.tinkerpop.gremlin.tinkergraph.jsr223.TinkerGraphGremlinPlugin: {},
               org.apache.tinkerpop.gremlin.jsr223.ImportGremlinPlugin: {classImports: [java.lang.Math], methodImports: [java.lang.Math#*]},
               org.apache.tinkerpop.gremlin.jsr223.ScriptFileGremlinPlugin: {files: [scripts/airline-sample.groovy]}}}}

сериализаторы:
# GraphBinary здесь заменяет Gryo и Graphson
  - { className: org.apache.tinkerpop.gremlin.driver.ser.GraphBinaryMessageSerializerV1, config: { ioRegistries: [org.janusgraph.graphdb.tinkerpop.JanusGraphIoRegistry] }}
  - { className: org.apache.tinkerpop.gremlin.driver.ser.GraphBinaryMessageSerializerV1, config: { serializeResultToString: true }}
  # Gryo и Graphson, последние версии
  - { className: org.apache.tinkerpop.gremlin.driver.ser.GryoMessageSerializerV3d0, config: { ioRegistries: [org.janusgraph.graphdb.tinkerpop.JanusGraphIoRegistry] }}
  - { className: org.apache.tinkerpop.gremlin.driver.ser.GryoMessageSerializerV3d0, config: { serializeResultToString: true }}
  - { className: org.apache.tinkerpop.gremlin.driver.ser.GraphSONMessageSerializerV3d0, config: { ioRegistries: [org.janusgraph.graphdb.tinkerpop.JanusGraphIoRegistry] }}
  # Старые версии сериализации для обратной совместимости:
  - { className: org.apache.tinkerpop.gremlin.driver.ser.GryoMessageSerializerV1d0, config: { ioRegistries: [org.janusgraph.graphdb.tinkerpop.JanusGraphIoRegistry] }}
  - { className: org.apache.tinkerpop.gremlin.driver.ser.GryoMessageSerializerV1d0, config: { serializeResultToString: true }}
  - { className: org.apache.tinkerpop.gremlin.driver.ser.GryoLiteMessageSerializerV1d0, config: {ioRegistries: [org.janusgraph.graphdb.tinkerpop.JanusGraphIoRegistry] }}
  - { className: org.apache.tinkerpop.gremlin.driver.ser.GraphSONMessageSerializerGremlinV2d0, config: { ioRegistries: [org.janusgraph.graphdb.tinkerpop.JanusGraphIoRegistry] }}
  - { className: org.apache.tinkerpop.gremlin.driver.ser.GraphSONMessageSerializerGremlinV1d0, config: { ioRegistries: [org.janusgraph.graphdb.tinkerpop.JanusGraphIoRegistryV1d0] }}
  - { className: org.apache.tinkerpop.gremlin.driver.ser.GraphSONMessageSerializerV1d0, config: { ioRegistries: [org.janusgraph.graphdb.tinkerpop.JanusGraphIoRegistryV1d0] }}

процессоры:
  - { className: org.apache.tinkerpop.gremlin.server.op.session.SessionOpProcessor, config: { sessionTimeout: 28800000 }}
  - { className: org.apache.tinkerpop.gremlin.server.op.traversal.TraversalOpProcessor, config: { cacheExpirationTime: 600000, cacheMaxSize: 1000 }}

метрики: {
  consoleReporter: {enabled: false, interval: 180000},
  csvReporter: {enabled: false, interval: 180000, fileName: /tmp/gremlin-server-metrics.csv},
  jmxReporter: {enabled: false},
  slf4jReporter: {enabled: true, interval: 180000},
  gangliaReporter: {enabled: false, interval: 180000, addressingMode: MULTICAST},
  graphiteReporter: {enabled: false, interval: 180000}}
пулПотоковБосса: 1
maxInitialLineLength: 4096
maxHeaderSize: 8192
maxChunkSize: 8192
maxContentLength: 65536
maxAccumulationBufferComponents: 1024
размерпакетаИтерацииРезультатов: 64
writeBufferHighWaterMark: 32768
writeBufferHighWaterMark: 65536
ssl: {
  включен: false}

У меня случайно "упала" BerkeleyDB версия JanusGraph.

Документацията е доста несигурна по отношение на индексите, тъй като управлението на индексите изисква извършването на доста странни манипулации с Groovy. Например, създаването на индекс трябва да се осъществи чрез написването на код в конзолата Gremlin (която, между другото, не работи по подразбиране). От официалната документация на JanusGraph:

graph.tx().rollback() // Никога не създавайте нови индекси, докато транзакцията е активна
mgmt = graph.openManagement()
name = mgmt.getPropertyKey('name')
age = mgmt.getPropertyKey('age')
mgmt.buildIndex('byNameComposite', Vertex.class).addKey(name).buildCompositeIndex()
mgmt.buildIndex('byNameAndAgeComposite', Vertex.class).addKey(name).addKey(age).buildCompositeIndex()
mgmt.commit()

// Изчакайте индексът да стане наличен
ManagementSystem.awaitGraphIndexStatus(graph, 'byNameComposite').call()
ManagementSystem.awaitGraphIndexStatus(graph, 'byNameAndAgeComposite').call()
// Преиндемтирайте съществуващите данни
mgmt = graph.openManagement()
mgmt.updateIndex(mgmt.getGraphIndex("byNameComposite"), SchemaAction.REINDEX).get()
mgmt.updateIndex(mgmt.getGraphIndex("byNameAndAgeComposite"), SchemaAction.REINDEX).get()
mgmt.commit()

Следсловие

В известен смисъл представеният по-горе експеримент е сравнение на топло с меко. Ако се замислите, графовата СУБД извършва различни операции за получаване на същите резултати. Въпреки това, в хода на тестовете проведох и експеримент с заявка от вида:

g.V().hasLabel('ZoneStep').has('id',0)
    .repeat(__.out().simplePath()).until(__.hasLabel('ZoneStep').has('id',1)).count().next()

която отразява стъпковата достъпност. Въпреки това, и с такива данни графовата СУБД показваше резултати, които надвишаваха няколко секунди... Това определено е свързано с факта, че имаше пътища от вида 0 -> X -> Y ... -> 1, които графовият двигател също проверяваше.

Дори за заявка от вида:

g.V().hasLabel('ZoneStep').has('id',0).out().has('id',1)).count().next()

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

Морала на баснята е, че красивата идея и парадигматичното моделиране не водят до желания резултат, който с значително по-висока ефективност се демонстрира на примера на ClickHouse. Представената в тази статия употреба е явен антипатерн за графови СУБД, въпреки че изглежда подходяща за моделиране в тяхната парадигма.

Източник: habr.com

Купете надежден хостинг за сайтове с защита от DDoS, VPS VDS сървъри 🔥 Купете надежден хостинг за сайтове с защита от DDoS, VPS VDS сървъри | ProHoster