Как использовать Neo4j в data engineering проекте: практический пример с Kafka, ClickHouse и Airflow

Когда говорят о Neo4j, чаще всего вспоминают social networks, fraud detection, recommendation systems или knowledge graphs. Но графовая база может быть полезна и внутри обычного data engineering проекта, где уже есть Kafka, аналитическая база типа ClickHouse, DBT, Airflow и BI. Мне было интересно проверить именно этот сценарий:
Как встроить Neo4j в современный data stack так, чтобы он не дублировал data warehouse, а решал те задачи, для которых графовая модель действительно удобнее SQL?
В качестве данных я использовал AIS — поток сообщений о движении судов, https://www.barentswatch.no/en/. Но основная тема этой статьи не AIS. Основная тема — место Neo4j в data engineering архитектуре.
Исходная архитектура
Проект строился вокруг довольно привычного набора компонентов:
Data sources
↓
Kafka
↓
ClickHouse
↓
dbt
↓
Airflow
↓
MetabaseClickHouse хранит историю AIS‑событий.
DBT строит аналитические модели.
Airflow управляет ежедневными pipeline.
Metabase показывает результаты.
На этом этапе возникает естественный вопрос: Зачем здесь еще Neo4j?
Ответ заключается в том, что часть данных имеет не только временную, но и сетевую структуру. Например:
Vessel
↓
VISITED
↓
Port
↓
CONNECTED_TO
↓
PortА дальше появляются вопросы:
какие порты наиболее важны в сети;
какие группы портов связаны между собой;
какие суда посещали связанные порты;
как устроена структура движения между портами;
можно ли выделить сообщества внутри сети.
Это уже не просто агрегация по таблице. Это граф.
Архитектура представлена на рисунке ниже.

Главное архитектурное решение: Neo4j не заменяет ClickHouse
Первое, что я хотел избежать, — попытки использовать Neo4j как универсальную БД. В проекте обязанности разделены.
ClickHouse отвечает за:
raw events
historical data
time series
aggregations
daily features
historical snapshotsNeo4j отвечает за:
entities
relationships
current graph state
paths
graph topology
graph algorithmsТо есть Neo4j не хранит всю историю AIS‑позиций. Полная история остается в raw‑слое ClickHouse — таблице raw.ais_positions А Neo4j хранит то, что действительно имеет смысл представлять как граф. В текущей архитектуре Neo4j используется как слой текущего графового состояния, а рассчитанные исторические graph snapshots сохраняются обратно в ClickHouse. Для меня это стало главным принципом дизайна:
Графовая база данных должна дополнять аналитическую базу данных, а не пытаться ее заменить.
Что именно мы храним в Neo4j
Базовая модель получилась очень простой:
(Vessel)-[:VISITED]->(Port)
(Port)-[:CONNECTED_TO]->(Port)
(Port)-[:MEMBER_OF]->(Community)То есть вместо миллионов строк:
mmsi
timestamp
latitude
longitude
speed
...в Neo4j мы работаем с предметными сущностями:
Vessel
Port
Communityи отношениями между ними. Это принципиально разные уровни хранения данных.
Шаблон 1: Kafka → Neo4j в реальном времени
Первый способ использования Neo4j в проекте — поддержание текущего состояния объектов.
Поток данных в реальном времени выглядит так:
BarentsWatch
↓
Python producer
↓
Kafka
├────────────→ ClickHouse
│
└────────────→ Kafka Connect
↓
Neo4jProducer (Скрипт на Python) публикует AIS event в Kafka.
ClickHouse сохраняет полную историю событий.
Kafka Connect Neo4j Sink обновляет текущий узел
Vesselв Neo4j.
То есть одно и то же сообщение Kafka используется двумя потребителями с разными задачами:
Kafka event
├─→ ClickHouse
│ historical event
│
└─→ Neo4j
current entity stateProducer использует MMSI как ключ в сообщениях Kafka, а Neo4j sink connector обновляет текущее состояние Vessel по MMSI. Упрощенно логика выглядит так:
MERGE (v:Vessel {mmsi: event.mmsi})
SET
v.name = event.name,
v.shipType = event.shipType,
v.latitude = event.latitude,
v.longitude = event.longitude,
v.speedOverGround = event.speedOverGround,
v.courseOverGround = event.courseOverGround,
v.lastSeen = datetime(event.msgtime)Здесь Neo4j работает практически как операционное графовое хранилище.
Почему не создавать Position node для каждого AIS event?
Можно было бы сделать:
(Vessel)-[:HAS_POSITION]->(Position)для каждой координаты. Но тогда Neo4j быстро превратился бы в хранилище миллионов time‑series событий. Это уже не та задача, в которой графовая модель дает максимальную пользу.
Поэтому я оставил:
all events
→ ClickHouseа в Neo4j:
current vessel
ports
relationships
communitiesПолучается очень понятное разделение:
ClickHouse отвечает на вопрос «что происходило во времени?»
Neo4j отвечает на вопрос «как сущности связаны между собой?»
Шаблон 2: использовать Neo4j после SQL/data warehouse обработки
Следующий интересный вариант — использовать Neo4j не как хранилище для первичного приема данных, а после аналитической обработки. AIS позиции сначала обрабатываются в ClickHouse. Из них определяется факт посещения порта. Условно:
raw.ais_positions
↓
port detection
↓
port visitsВ ClickHouse можно хранить подробные записи о визите порта. Но затем из этих записей строится граф:
(Vessel)-[:VISITED]->(Port)В контексте данной задачи Neo4j может быть не первой системой после источника данных. Он может быть аналитический движок последующей обработки. То есть:
raw data
↓
warehouse
↓
transformations
↓
business entities
↓
Neo4jЭто делает Neo4j особенно интересным для data engineering.
Шаблон 3: превращаем факты в отношения
После появления VISITED можно построить переходы между портами. Например:
Vessel A:
BERGEN → TANANGER → EGERSUNDиз этого появляется:
(:Port)-[:CONNECTED_TO]->(:Port)и постепенно строится network:
Port B
/ \
Port A Port D
\ /
Port CВ реляционный базе данных это можно представить таблицей:
port_from
port_to
weightи для простых запросов этого вполне достаточно. Но как только появляются вопросы о структуре всей сети, графовая база данных становится намного интереснее и производительнее с ростом глубины уровней.
Шаблон 4: Алгоритмы графов в Neo4j
Neo4j в проекте используется не только для хранения отношений между узлами Neo4j.
Следующий слой — Neo4j Graph Data Science (https://neo4j.com/docs/graph‑data‑science/current/).
На графе портов мы запускаем алгоритмы графов. Сейчас основными алгоритмами являются:
PageRank
LouvainПроект действительно использует PageRank для оценки важности портов и Louvain для выявления community. И это уже совсем другая роль Neo4j. Он становится не просто базой данных, а: специализированный аналитический движок для графовых задач
PageRank как пример свойства узла графа Port
Представим обычную метрику в SQL базе данных:
connections_countМы можем посчитать:
Port A = 10 connections
Port B = 7 connections
Port C = 4 connectionsНо количество связей не всегда равно важности. Порт может иметь мало связей, но быть соединен с ключевыми хабами портов. PageRank учитывает структуру сети. В результате появляется свойство узла Port:
Port.pageRankНапример:
BERGEN 0.63
TANANGER 0.42
EGERSUND 0.42Эту метрику затем можно использовать:
Neo4j
→ PageRank
→ ClickHouse
→ dbt
→ MetabaseЭто хороший пример того, как алгоритмы графа становятся обычной аналитической метрикой внутри платформы данных.
Louvain как пример feature engineering через граф
Вторая задача — выявление сообществ (community). Мы хотим понять, какие порты образуют естественные сообщества (community). На входе мы имеем:
Port ↔ Port ↔ Port ↔ PortПосле Louvain:
Community A
├── Port 1
├── Port 2
└── Port 3
Community B
├── Port 4
├── Port 5
└── Port 6В реляционный базе данных мы получили бы:
port_id | community_idНо Neo4j позволяет сделать узел Community полноценной частью графовой модели данных:
(:Port)-[:MEMBER_OF]->(:Community)Это уже полезно для дальнейших обходов по узлам графа. Например:
Vessel
↓ VISITED
Port
↓ MEMBER_OF
CommunityПолучается цепочка:
Vessel → Port → Communityкоторая очень естественно читается как предметная модель.
Почему появились узлы Community
Сначала community была просто свойством узла Port:
Port.communityId = 1Но со временем появилась необходимость работать с community как самостоятельной сущностью. Например:
Community
├── community_id
├── community_name
├── community_label
└── port_countи связь:
Port-[:MEMBER_OF]->Communityделает модель данных в графе более расширяемой. Позже к узлу Community можно добавить свойства:
region
dominant traffic type
risk score
historical metrics
LLM summaryне меняя основную модель портов. Пример community представлен на рисунке ниже:

От технического community ID к смысловому значению community
Графовый алгоритм Louvain возвращает, например:
communityId = 1Для алгоритма этого достаточно. Для BI — нет. Поэтому я добавил смысловые текстовые значения. Например:
community_id = 1
community_name = BERGEN
community_label = BERGEN / TANANGER / EGERSUNDНаименования портов выбираются детерминированно, в том числе используя PageRank. Проект сохраняет технический ключ community ID, но показывает человеку понятное имя. Это хороший пример перехода:
graph algorithm output
↓
business-oriented data modelШаблон 5: возвращать результаты Neo4j обратно в ClickHouse
Очень важная часть архитектуры — результаты GDS не остаются только в Neo4j. Почему Допустим сегодня:
BERGEN.pageRank = 0.63а завтра:
BERGEN.pageRank = 0.58Если просто обновлять Neo4j property, историческое значение потеряется. Поэтому рассчитанные в Neo4j метрики экспортируются обратно в ClickHouse. Получается:
Neo4j
↓
GDS
↓
PageRank + Communities
↓
ClickHouse snapshots
↓
dbt
↓
MetabaseТо есть Neo4j выполняет расчеты, а ClickHouse сохраняет историю. Это архитектурно очень похоже на ML поток данных:
warehouse
↓
ML engine
↓
predictions/features
↓
warehouseТолько здесь:
warehouse
↓
graph engine
↓
graph features
↓
warehouseИменно этот шаблон я считаю одним из самых полезных способов интеграции Neo4j в data engineering.
Neo4j как движок для расчета метрик
Если посмотреть шире, PageRank и communityId — это по сути свойства узла Neo4j. Например:
port_id
page_rank
community_id
community_name
community_sizeЭти свойства узла затем можно использовать в:
BI
ML
anomaly detection
LLM context
rankingТо есть роль Neo4j можно сформулировать так: Neo4j генерирует свойства, которые трудно или неудобно получить обычными табличными трансформациями.
Это особенно интересно для Data Engineer.
Шаблон 6: Neo4j + Kafka Connect
Еще один хороший интеграционного шаблона — Kafka Connect. Вместо того чтобы писать:
Python
→ consume Kafka
→ open Neo4j connection
→ execute Cypherиспользуется:
Kafka
→ Kafka Connect
→ Neo4j SinkЭто разделяет ответственность. Producer знает только Kafka. Kafka Connect знает, как доставить данные в Neo4j. Neo4j отвечает за доставленные данные. Bootstrap проекта автоматически применяет constraints и создает или обновляет Neo4j Kafka Sink connector. Получается:
source application
↓
Kafka
↓
integration layer
↓
Neo4jШаблон 7: Neo4j + Airflow
Neo4j не обязательно должен жить отдельно от обычной orchestration инфраструктуры. В проекте поток данных в/из графа включен в Airflow. Основной daily flow концептуально выглядит так:
AIS ready
↓
dbt
↓
port visits
↓
Neo4j sync
↓
GDS
↓
graph metrics export
↓
ClickHouse
↓
BIТо есть Neo4j становится обычным звеном в потоке данных. Airflow отвечает за управление зависимостями:
→ data prepared
→ graph updated
→ algorithms executed
→ results exportedЭто важный момент. Neo4j не отдельное приложение с ручными Cypher запросами. Он встроен в поток данных.
Шаблон 8: Neo4j + dbt
dbt и Neo4j в проекте решают разные задачи.
dbt:
SQL transformations
data quality tests
dimensions
facts
daily features
aggregationsNeo4j:
relationships
network topology
centrality
community detection
graph traversalПоток данных может выглядеть так:
dbt
↓
prepare graph input
↓
Neo4j / GDS
↓
graph features
↓
ClickHouse
↓
dbt
↓
final analytics martsВ текущем проекте dbt действительно содержит метрики vessel features, anomalies, current port visits, enriched graph metrics и community reporting models. Получается очень полезная комбинация:
SQL analytics
+
graph analyticsШаблон 9: пользовательский Neo4j Docker image
Еще один практический момент: Neo4j в реальном проекте часто требует дополнительных plugins. В моей конфигурации используются:
APOC
APOC Extended
Graph Data Science
ClickHouse JDBCDocker Compose действительно загружает APOC, APOC Extended и GDS, а также настраивает JDBC connection к ClickHouse. Поэтому появился собственный Docker image:
ais-graph-analytics-neo4jКонцептуально:
Neo4j upstream image
+
APOC
+
APOC Extended
+
GDS
+
ClickHouse JDBCЭто важно для воспроизводимости. На локальном компьютере и в CI используется одна и та же конфигурация Docker image.
Шаблон 10: Neo4j должен участвовать в CI
Еще одна вещь, которую часто пропускают в дело проектах:
Cypher works locallyеще не означает:
graph pipeline is reproducibleПоэтому CI проекта содержит отдельный Kafka + Neo4j smoke test.
Он проверяет:
custom Neo4j image build
startup
plugins
authentication
constraints
write/readВ общем CI есть отдельные jobs для Compose, pipeline, producer, dbt+ClickHouse и Kafka+Neo4j. Для меня это важная часть демонстрации Neo4j именно как data engineering component, а не как база данных, которую вручную запускают в Neo4j Desktop.
А где здесь BI?
Результаты графовых алгоритмов в итоге отображаются в Metabase. Например:
Top ports by PageRank
Community
Community label
Ports per community
graph activityПользователю BI не обязательно знать, что значение было рассчитано при помощи Neo4j GDS. Для него это просто аналитические метрики. И мне нравится именно такой подход. Графовая база данных становится внутренним компонентом внутри платформы данных. Для визуализации использовался дашборд в Metabase.

Общая модель
Если отбросить AIS, архитектуру можно обобщить. Допустим есть обычный warehouse:
Raw
↓
Warehouse
↓
TransformationsВ данных обнаруживается реляционная структура:
Customer → Account
Account → Transaction
Person → Device
Product → Customer
Service → Dependency
Document → EntityТогда можно добавить:
Warehouse
↓
prepare entities/relationships
↓
Neo4j
↓
Graph algorithms
↓
graph features
↓
WarehouseПолучается архитектура с возможностью повторного использования:
┌───────────────┐
│ Sources │
└───────┬───────┘
↓
Kafka
↓
Data warehouse
↓
dbt/SQL
↓
Graph preparation
↓
Neo4j
↓
Graph Data Science
↓
graph features
↓
Data warehouse
↓
BI / ML / AIВот именно этот шаблон я и хотел продемонстрировать проектом.
Когда Neo4j действительно имеет смысл
Я бы не добавлял Neo4j в проект только потому, что графовая база данных выглядит интересно. Он имеет смысл, когда бизнес‑вопрос связан с отношениями. Например:
fraudPerson
→ Account
→ Transaction
→ Merchantили:
customer identityEmail
→ Person
→ Device
→ Accountили:
supply chainSupplier
→ Component
→ Product
→ Warehouseили:
infrastructureService
→ Depends On
→ Serviceили, как в нашем случае:
Vessel
→ Port
→ Port
→ CommunityЕсли основной вопрос звучит: Как объекты связаны? Neo4j становится хорошим кандидатом.
Когда Neo4j не нужен
Если задача:
SUM(revenue)
GROUP BY monthили
AVG(speed)
GROUP BY vesselдля нее нужен ClickHouse.
Если задача:
daily KPIэто SQL/dbt.
Не стоит переносить такие расчеты в Neo4j только ради использования графовой базы данных.
Я бы сформулировал правило так:
Если структура отношений сама по себе является частью аналитики — стоит рассмотреть Neo4j.
Что Neo4j дал этому проекту
Без Neo4j проект оставался бы полноценным data engineering stack:
Kafka
ClickHouse
dbt
Airflow
MetabaseНо Neo4j добавил отдельное аналитическое измерение.
Вместо:
какие порты посещалисьмы можем спрашивать:
как порты связаныВместо:
сколько связей имеет портможем:
какова важность порта в сетиВместо:
какие пары портов встречаютсяможем:
какие сообщества портов образует весь графИ результаты этих вычислений снова становятся частью обычного analytics stack.
Что дальше
Следующий этап проекта — streaming processing через Apache Flink. Это позволит построить near real‑time дашборд в Metabase.
Первый этап уже выглядит так:
Kafka
↓
Flink
↓
stream transformation
↓
KafkaДальше я планирую:
Historical data
↓
Iceberg
↓
Trinoи затем попробовать соединить:
streaming features
+
graph features
+
ML
+
LLMОсобенно интересным мне кажется AI Agent, который сможет обращаться одновременно к:
ClickHouse
Neo4j
Iceberg / TrinoНапример, запрос:
Какие суда имели необычное поведение и при этом посещали наиболее значимые порты?
может потребовать:
time-series analytics
+
graph features
+
anomaly detection
+
LLM interpretationИ здесь Neo4j становится уже не отдельной графовой базой данных, а одним из аналитических движков внутри более крупной платформы данных.
Итог
Главная идея проекта для меня выглядит так:
Neo4j можно использовать в data engineering не вместо data warehouse, а рядом с ним.
ClickHouse хранит историю.
DBT выполняет трансформации.
Kafka транспортирует сообщения от API.
Airflow управляет потоком данных.
Neo4j строит графовую модель.
GDS рассчитывает свойства узлов графа.
Результаты возвращаются в ClickHouse.
BI, ML или AI используют уже объединенные данные.
В моем проекте это выглядит так:
Kafka
↓
ClickHouse
↓
dbt
↓
Neo4j
↓
GDS
↓
ClickHouse
↓
MetabaseИменно такой подход я хотел показать. Neo4j становится особенно полезным не тогда, когда мы просто переносим таблицы в графовую форму, а тогда, когда:
связи между объектами становятся самостоятельным источником аналитической информации.
Для AIS это:
Vessel
↓
Port
↓
Connected Port
↓
CommunityДля другого data engineering проекта это может быть:
Customer → Account → Transactionили:
Device → User → Identityили:
Service → Dependency → ServiceТехнологии и domain меняются. Сами шаблоны остается тем же.
Проект выполнен на ноутбуке с использованием Docker. Проект доступен в GitHub:
https://github.com/KonstantinLofichenko/ais‑graph‑analytics
KioskNews shows a cleaned-up reading view extracted from the publisher’s page — the original always lives on their site, not ours.