ESPN DeportesEl Atlético se complica y el Alavés le empata con un penaltiRTP Desporto"Foi uma graçola." Varandas responde ao primeiro-ministro com referência à corrupçãoThe Jerusalem PostUS pushed Israel to attack Iran alone in effort to avoid political fallout before midterms - reportInquirer6,000 openings up in Eastern Visayas job fairESPNAnalysis, highlights from sweeps by Valkyries, DreamPunchBayelsa CP urges community leaders to help curb cult violenceZDF heuteAktuelle Pressemitteilungen des ZDFObservador DesportoPepa avisa: Arouca "dificulta" mas foco é o EstrelaBBC عربيخمسة جرحى في انفجار قوي بمطار الملك خالد بالرياضRai NewsLa coalizione dei ricorsi: il melonellum si sposta dal parlamento ai tribunaliESPN CricinfoRenshaw 190 lifts Australia to 390 before Hazlewood strikesSRF NewsErstes Monument geholt – Jungstar Seixas triumphiert erstmals bei der Lombardei-Rundfahrt
The Daily Newsstand · Free, Always
Saturday, October 10, 2026

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

Translate

Когда говорят о 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
    ↓
Metabase
  1. ClickHouse хранит историю AIS‑событий.

  2. DBT строит аналитические модели.

  3. Airflow управляет ежедневными pipeline.

  4. Metabase показывает результаты.

На этом этапе возникает естественный вопрос: Зачем здесь еще Neo4j?

Ответ заключается в том, что часть данных имеет не только временную, но и сетевую структуру. Например:

Vessel
  ↓
VISITED
  ↓
Port
  ↓
CONNECTED_TO
  ↓
Port

А дальше появляются вопросы:

  • какие порты наиболее важны в сети;

  • какие группы портов связаны между собой;

  • какие суда посещали связанные порты;

  • как устроена структура движения между портами;

  • можно ли выделить сообщества внутри сети.

Это уже не просто агрегация по таблице. Это граф.

Архитектура представлена на рисунке ниже.

Главное архитектурное решение: Neo4j не заменяет ClickHouse

Первое, что я хотел избежать, — попытки использовать Neo4j как универсальную БД. В проекте обязанности разделены.

ClickHouse отвечает за:

raw events
historical data
time series
aggregations
daily features
historical snapshots

Neo4j отвечает за:

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
                      ↓
                    Neo4j
  1. Producer (Скрипт на Python) публикует AIS event в Kafka.

  2. ClickHouse сохраняет полную историю событий.

  3. Kafka Connect Neo4j Sink обновляет текущий узел Vessel в Neo4j.

То есть одно и то же сообщение Kafka используется двумя потребителями с разными задачами:

Kafka event
   ├─→ ClickHouse
   │     historical event
   │
   └─→ Neo4j
         current entity state

Producer использует 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
aggregations

Neo4j:

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 JDBC

Docker 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 в проект только потому, что графовая база данных выглядит интересно. Он имеет смысл, когда бизнес‑вопрос связан с отношениями. Например:

fraud
Person
→ Account
→ Transaction
→ Merchant

или:

customer identity
Email
→ Person
→ Device
→ Account

или:

supply chain
Supplier
→ Component
→ Product
→ Warehouse

или:

infrastructure
Service
→ 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

View the original on Хабр →

KioskNews shows a cleaned-up reading view extracted from the publisher’s page — the original always lives on their site, not ours.