
Cześć wszystkim. Opracowujemy produkt do analizy ruchu offline. W projekcie istnieje zadanie związane z statystyczną analizą ścieżek poruszania się odwiedzających w obszarach.
W ramach tego zadania użytkownicy mogą zadawać systemowi zapytania następującego rodzaju:
- ile odwiedzających przeszło z obszaru "A" do obszaru "B";
- ile odwiedzających przeszło z obszaru "A" do obszaru "B" przez obszar "C", a następnie przez obszar "D";
- ile czasu zajęło przejście odwiedzającemu określonego typu z obszaru "A" do obszaru "B".
i jeszcze wiele podobnych zapytań analitycznych.
Ruch odwiedzających po obszarach stanowi skierowany graf. Po przeszukaniu internetu odkryłem, że grafowe bazy danych są również używane do raportów analitycznych. Zainteresowałem się, jak te bazy danych radzą sobie z takimi zapytaniami (TL;DR; słabo).
Wybrałem do użycia bazę danych , jako wybitnego przedstawiciela otwartych baz danych grafowych, która opiera się na zestawie dojrzałych technologii, które (moim zdaniem) powinny zapewnić jej przyzwoite osiągi operacyjne:
- backend magazynu BerkeleyDB, Apache Cassandra, Scylla;
- złożone indeksy można przechowywać w Lucene, Elasticsearch, Solr.
Autorzy JanusGraph piszą, że nadaje się ona zarówno do OLTP, jak i OLAP.
Pracowałem z BerkeleyDB, Apache Cassandra, Scylla i ES, ponadto te produkty są często używane w naszych systemach, więc z optymizmem podszedłem do testowania tej grafowej bazy danych. Wydawał mi się dziwny wybór BerkeleyDB zamiast RocksDB, ale prawdopodobnie ma to związek z wymaganiami dotyczącymi transakcji. W każdym razie, do skalowalnego, produkcyjnego użytku zaleca się użyvertrag backend na Cassandra lub Scylla.
Neo4j nie brałem pod uwagę, ponieważ do klasteryzacji wymagana jest wersja komercyjna, więc produkt nie jest otwarty.
Grafowe bazy danych mówią: "Jeśli coś wygląda jak graf — traktuj to jak graf!" — piękne!
Najpierw narysowałem graf, który był zrobiony zgodnie z kanonami baz danych grafowych:

Istnieje byt Zone, odpowiadający za obszar. Jeśli ZoneStep należy do tej Zone, to się na nią odnosi. Na byty Area, ZoneTrack, Person nie zwracaj uwagi, należą do domeny i w ramach testu nie są rozważane. Na koniec, do takiej struktury grafowej zapytanie o wyszukiwanie łańcuchów wyglądałoby następująco:
g.V().hasLabel('Zone').has('id',0).in_()
.repeat(__.out()).until(__.out().hasLabel('Zone').has('id',19)).count().next()Po polsku to mniej więcej: znajdź Zone o ID=0, weź wszystkie wierzchołki, z których wychodzi krawędź (ZoneStep), idź bez powrotu, aż znajdziesz takie ZoneStep, z których wychodzi krawędź do Zone o ID=19, policz liczbę takich łańcuchów.
Nie pretensje do znania wszystkich szczegółów wyszukiwania w grafach, ale to zapytanie zostało wygenerowane na podstawie tej książki ().
Załadowałem 50 tysięcy utworów o długości od 3 do 20 punktów do bazy grafowej JanusGraph, korzystającej z backendu BerkeleyDB, stworzyłem indeksy zgodnie z .
Skrypt do ładowania w Pythonie:
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)Używano VM z 4 rdzeniami i 16 GB RAM na SSD. JanusGraph został uruchomiony za pomocą następującego polecenia:
docker run --name janusgraph -p8182:8182 janusgraph/janusgraph:latestW tym przypadku dane i indeksy używane do wyszukiwania według ścisłego dopasowania są przechowywane w BerkeleyDB. Wykonując wcześniej podane zapytanie, otrzymałem czas równy kilku dziesiątkom sekund.
Uruchamiając 4 powyższe skrypty równolegle, udało mi się zamienić DB na dyni z wesołym strumieniem stacktrace'ów Java (a wszyscy lubimy czytać stacktrace'y Java) w logach Docker.
Pomyślałem i postanowiłem uprościć schemat grafu do następnego:

Stwierdziłem, że wyszukiwanie według atrybutów encji będzie szybsze niż wyszukiwanie według krawędzi. W rezultacie moje zapytanie przekształciło się w następujące:
g.V().hasLabel('ZoneStep').has('id',0).repeat(__.out().simplePath()).until(__.hasLabel('ZoneStep').has('id',19)).count().next()Po polsku to mniej więcej: znajdź ZoneStep o ID=0, idź bez powrotu, aż znajdziesz ZoneStep o ID=19, policz liczbę takich łańcuchów.
Skrypt ładowania, podany powyżej, również uprościłem, aby nie tworzyć zbędnych powiązań, ograniczając się do atrybutów.
Zapytanie wciąż trwało kilka sekund, co było całkowicie nieakceptowalne dla naszego zadania, ponieważ do celów zapytań AdHoc, tego typu rozwiązanie zupełnie się nie nadaje.
Spróbowałem uruchomić JanusGraph z użyciem Scylla, jako najszybszej implementacji Cassandry, ale to również nie przyniosło żadnych istotnych zmian w wydajności.
Tak więc, mimo że "to wygląda jak graf", nie udało mi się skłonić bazę danych do grafów do szybkiej obróbki danych. Jestem przekonany, że nie znam czegoś i można zmusić JanusGraph do wykonania tego zapytania w ułamku sekundy, ale mi się to nie udało.
Ponieważ i tak trzeba było rozwiązać problem, zacząłem myśleć o złączeniach (JOIN) i tabelach przestawnych (Pivot), co nie nastrajało optymistycznie z perspektywy elegancji, ale mogło się okazać całkiem praktycznym rozwiązaniem.
W naszym projekcie wykorzystujemy już Apache ClickHouse, dlatego postanowiłem sprawdzić moje badania na tym systemie analitycznym.
Uruchomiłem ClickHouse według prostego przepisu:
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-serverStworzyłem w nim bazę danych i tabelę typu:
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 = 8192Wypełniłem ją danymi za pomocą następującego skryptu:
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
)Ponieważ wstawienia odbywają się w partiach, wypełnienie danych było znacznie szybsze niż w przypadku JanusGraph.
Sformułowałem dwa zapytania za pomocą JOIN. Aby przejść z punktu A do punktu B:
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.whenAby przejść przez 3 punkty:
SELECT s3.person,
s1z,
s1w,
s2z,
s2w,
s3.zone,
s3.when
FROM
(SELECT s1.person AS person,
s1.zone AS s1z,
s1.when AS s1w,
s2.zone AS s2z,
s2.when AS s2w
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 = 3)) AS s2 USING person
WHERE s1.when <= s2.when) p ANY INNER JOIN
(SELECT *
FROM steps
WHERE (area = 0)
AND (zone = 19)) AS s3 USING person
WHERE p.s2w <= s3.whenZapytania rzeczywiście wyglądają dość strasznie, aby je skutecznie wykorzystać, potrzebny jest programowy interfejs - generator. Niemniej jednak działają i działają szybko. Zarówno pierwsze, jak i drugie zapytanie wykonują się w mniej niż 0,1 sekundy. Oto przykład czasu wykonania zapytania dla count(*) przechodzącego przez 3 punkty:
SELECT count(*)
FROM
(
SELECT
s1.person AS person,
s1.zone AS s1z,
s1.when AS s1w,
s2.zone AS s2z,
s2.when AS s2w
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 = 3)
) AS s2 USING (person)
WHERE s1.when <= s2.when
) AS p
ANY INNER JOIN
(
SELECT *
FROM steps
WHERE (area = 0) AND (zone = 19)
) AS s3 USING (person)
WHERE p.s2w <= s3.when
┌─count()─┐
│ 11592 │
└─────────┘1 wierszy w zbiorze. Czas: 0.068 sek. Przetworzono 250.03 tys. wierszy, 8.00 MB (3.69 miliona wierszy/s., 117.98 MB/s.)Uwagi na temat IOPS. Podczas wypełniania danych, JanusGraph generował dość wysoką liczbę IOPS (1000-1300 dla czterech strumieni wypełniania danymi), a IOWAIT był dość wysoki. Jednocześnie ClickHouse generował minimalne obciążenie systemu dyskowego.
Podsumowanie
Zdecydowaliśmy się użyć ClickHouse do obsługi tego typu zapytań. Zawsze możemy jeszcze bardziej zoptymalizować zapytania, korzystając z materializowanych widoków i równoległości, wykonując wstępną obróbkę strumienia zdarzeń przy użyciu Apache Flink przed załadowaniem ich do ClickHouse.
Wydajność jest na tyle dobra, że prawdopodobnie nie będziemy musieli nawet myśleć o ręcznych pivotach tabel. Wcześniej musieliśmy robić pivoty danych wyciąganych z Vertica poprzez eksport do Apache Parquet.
Niestety, kolejna próba użycia bazy danych grafowej nie zakończyła się sukcesem. Nie odkryłem, że JanusGraph ma przyjazne dla użytkownika ekosystem, który pozwala szybko zapoznać się z produktem. Co więcej, konfigurowanie serwera odbywa się w tradycyjny sposób Java, co sprawi, że osoby nieznające Javy będą płakać krwawymi łzami:
host: 0.0.0.0
port: 8182
threadPoolWorker: 1
gremlinPool: 8
scriptEvaluationTimeout: 30000
channelizer: org.janusgraph.channelizers.JanusGraphWsAndHttpChannelizer
graphManager: org.janusgraph.graphdb.management.JanusGraphManager
graphs: {
ConfigurationManagementGraph: conf/janusgraph-cql-configurationgraph.properties,
airlines: conf/airlines.properties
}
scriptEngines: {
gremlin-groovy: {
plugins: { 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]}}}}
serializers:
# GraphBinary ma zastąpić Gryo i 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 i Graphson, najnowsze wersje
- { 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] }}
# Starsze wersje serializacji dla zgodności:
- { 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] }}
processors:
- { 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 }}
metrics: {
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}}
threadPoolBoss: 1
maxInitialLineLength: 4096
maxHeaderSize: 8192
maxChunkSize: 8192
maxContentLength: 65536
maxAccumulationBufferComponents: 1024
resultIterationBatchSize: 64
writeBufferHighWaterMark: 32768
writeBufferHighWaterMark: 65536
ssl: {
enabled: false}Przypadkowo udało mi się "położyć" BerkeleyDB w wersji JanusGraph.
Dokumentacja jest dość krzywa w zakresie indeksów, ponieważ zarządzanie indeksami wymaga dość dziwnych sztuczek w Groovy. Na przykład, tworzenie indeksu powinno być realizowane poprzez kod w konsoli Gremlin (które, nawiasem mówiąc, nie działa od razu). Z oficjalnej dokumentacji JanusGraph:
graph.tx().rollback() //Nigdy nie twórz nowych indeksów, gdy transakcja jest aktywna
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()
//Czekaj na dostępność indeksu
ManagementSystem.awaitGraphIndexStatus(graph, 'byNameComposite').call()
ManagementSystem.awaitGraphIndexStatus(graph, 'byNameAndAgeComposite').call()
//Zindeksuj ponownie istniejące dane
mgmt = graph.openManagement()
mgmt.updateIndex(mgmt.getGraphIndex("byNameComposite"), SchemaAction.REINDEX).get()
mgmt.updateIndex(mgmt.getGraphIndex("byNameAndAgeComposite"), SchemaAction.REINDEX).get()
mgmt.commit()Epilog
W pewnym sensie powyższy eksperyment to porównanie ciepłego z miękkim. Jeśli się nad tym zastanowić, to grafowa baza danych wykonuje różne operacje, aby uzyskać te same rezultaty. Jednak w ramach testów prowadziłem również eksperyment z zapytaniem w postaci:
g.V().hasLabel('ZoneStep').has('id',0)
.repeat(__.out().simplePath()).until(__.hasLabel('ZoneStep').has('id',1)).count().next(), które odzwierciedla dostępność kroku. Jednak na takich danych grafowa baza danych pokazywała wyniki, które wykraczały poza kilka sekund… Z pewnością było to związane z tym, że istniały ścieżki w postaci 0 -> X -> Y ... -> 1, które silnik grafowy również sprawdzał.
Nawet dla zapytania w postaci:
g.V().hasLabel('ZoneStep').has('id',0).out().has('id',1)).count().next()nie udało mi się uzyskać wydajnej odpowiedzi z czasem przetwarzania poniżej sekundy.
Morał tej opowieści jest taki, że piękny pomysł i paradygmatyczne modelowanie nie prowadzą do pożądanego rezultatu, który z wyraźnie większą wydajnością demonstrowany jest na przykładzie ClickHouse. Przykład zastosowania przedstawiony w tym artykule to wyraźny antywzór dla grafowych baz danych, chociaż wygląda na odpowiedni do modelowania w ich paradygmatzie.
Źródło: habr.com
