Предположим, что от маршрутизатора Cisco, который выполняет функции BRAS и NAT подают на наш сервер два потока
1) NetFlow в формате IPFIX, порт UDP 2055, это статистика (Кто, куда, сколько байт передал). Формат полей данных следующий:
1) NetFlow в формате IPFIX, порт UDP 2055, это статистика (Кто, куда, сколько байт передал). Формат полей данных следующий:
{
"ip_src": "10.10.145.224",
"ip_dst": "54.191.232.8",
"port_src": 44354,
"port_dst": 443,
"bytes": 4934,
"packets": 5,
"ip_proto": 6,
"timestamp_start": 1787910941512000000,
"timestamp_end": 1787910941737000000,
"peer_ip_src": "172.19.0.2"
}
2) NEL в формате NetFlow v9, порт UDP 2056, это события состояния сессии NAT (создалась сессия NAT, удалилась сессия NAT).
NEL может передаваться разными шаблонами. Будем использовать самый простой шаблон, который маршрутизатор отдает постоянно. Формат данных следующий:
{
"peer_ip_src": "172.19.0.2",
"ip_src": "10.10.148.170",
"port_src": 49534,
"post_nat_ip_src": "97.221.112.206",
"post_nat_port_src": 58001,
"nat_event": 1, // 1=CREATE, 2=DELETE
"ip_proto": "udp",
"timestamp_start": "2026-09-07 16:36:10.415000"
}
Настроим связку которая будет принимать данные и стабильно складывать в базу Clickhouse.
Данные будут приниматься и разбираются на лету сервисом pmacct. Каждое сообщение будет преобразовываться в JSON формат и передается в агрегатор Kafka. Kafka будет выступать буфером, сглаживающим пики и объединяющим порции данных перед вставкой в Clickhouse.
Clickhouse будет постоянно читать топики Kafka и забирать новые данные в базу. Поток flow и поток nel будут очевидно хранится в разных таблицах базы, так как имеют разную структуру.
"ip_dst": "54.191.232.8",
"port_src": 44354,
"port_dst": 443,
"bytes": 4934,
"packets": 5,
"ip_proto": 6,
"timestamp_start": 1787910941512000000,
"timestamp_end": 1787910941737000000,
"peer_ip_src": "172.19.0.2"
}
2) NEL в формате NetFlow v9, порт UDP 2056, это события состояния сессии NAT (создалась сессия NAT, удалилась сессия NAT).
NEL может передаваться разными шаблонами. Будем использовать самый простой шаблон, который маршрутизатор отдает постоянно. Формат данных следующий:
{
"peer_ip_src": "172.19.0.2",
"ip_src": "10.10.148.170",
"port_src": 49534,
"post_nat_ip_src": "97.221.112.206",
"post_nat_port_src": 58001,
"nat_event": 1, // 1=CREATE, 2=DELETE
"ip_proto": "udp",
"timestamp_start": "2026-09-07 16:36:10.415000"
}
Настроим связку которая будет принимать данные и стабильно складывать в базу Clickhouse.
Данные будут приниматься и разбираются на лету сервисом pmacct. Каждое сообщение будет преобразовываться в JSON формат и передается в агрегатор Kafka. Kafka будет выступать буфером, сглаживающим пики и объединяющим порции данных перед вставкой в Clickhouse.
Clickhouse будет постоянно читать топики Kafka и забирать новые данные в базу. Поток flow и поток nel будут очевидно хранится в разных таблицах базы, так как имеют разную структуру.
*** Установка ClickHouse
# apt-get install -y apt-transport-https ca-certificates curl gnupg
# cd /usr/src/
У меня под рукой был старый сервер с процессором Intel Xeon E5620, который не поддерживает современные версии Clickhouse.
Процессор Intel Xeon E5620 относится к семейству Westmere (2010 год). Поддерживает инструкции SSE4.2, но не поддерживает AVX/AVX2. Это важно, так как современные версии ClickHouse требуют AVX2.
Оптимальным выбором будет 22.3 LTS, которые все еще поддерживаются и работают с SSE4.2.
# wget https://packages.clickhouse.com/tgz/lts/clickhouse-common-static-22.3.19.6-amd64.tgz
# wget https://packages.clickhouse.com/tgz/lts/clickhouse-server-22.3.19.6-amd64.tgz
# wget https://packages.clickhouse.com/tgz/lts/clickhouse-client-22.3.19.6-amd64.tgz
Установка:
# tar -xzvf clickhouse-common-static-22.3.19.6-amd64.tgz
# ./clickhouse-common-static-22.3.19.6/install/doinst.sh
# tar -xzvf clickhouse-server-22.3.19.6-amd64.tgz
# ./clickhouse-server-22.3.19.6/install/doinst.sh
# tar -xzvf clickhouse-client-22.3.19.6-amd64.tgz
# ./clickhouse-client-22.3.19.6/install/doinst.sh
Настройка Clickhouse
1. Задать пользовательские настройки
# nano /etc/clickhouse-server/users.xml
1.1. Задать пароль по умолчанию для доступа к Clickhouse
Секция <users> -> <default> -> <password>
Прописываем пароль:
<password>*******</password>
1.2. Задать ограничения по использованию памяти и ядер процессора
Секция <profiles> -> <default>
Размер памяти определяем как 50% от памяти, которую выделим Clickhouse (в моем случае - 32 ГБ)
Для сервера с 16 ядрами и выделенными 32 ГБ RAM под ClickHouse, рекомендация max_threads = 4.
<max_memory_usage>17179869184</max_memory_usage>
<max_threads>4</max_threads>
2.Правим конфигурацию сервера:
# nano /etc/clickhouse-server/config.xml
2.1. Отключаем использование IPv6, раскомментируем строку (тем самым отключаем IPv6)
<listen_host>0.0.0.0</listen_host>
2.2. Изменяем уровень логирования – оставляем только ошибки
Так же устанавливаем размер лог-файла в 100Мбайт (вместо 1G)
<logger>
<level>warning</level>
<size>100M</size>
<logger>
2.3. Ограничиваем сервер полкой в использовании памяти - в 32G, иначе сервер съест другие приложения
<max_server_memory_usage>34359738368</max_server_memory_usage>
2.4. Отключаем процентное указание в использование памяти
<max_server_memory_usage_to_ram_ratio>0</max_server_memory_usage_to_ram_ratio>
2.5. Проверить, что Clickhouse работает с директорией, которая для этого была создана в примонтированном хранилище:
<path>/var/lib/clickhouse/</path>
2.6. Из блока <openSSL> вырезать раздел <server>, оставить только раздел клиента
<openSSL>
<client>
...
</client>
</openSSL>
2.7. Добавляем в конец конфигурации перед закрывающим тегом </clickhouse> директивы, определяющие протокол взаимодействия между Kafka и Clickhouse.
Это необходимо, так как версия 22.3 не совместима с новыми сервисами Kafka.
<kafka>
<api_version_request>false</api_version_request>
<broker_version_fallback>2.4.1</broker_version_fallback>
</kafka>
</clickhouse>
Создаем пользователя clickhouse и передаем ему права на директорию, в которой он должен работать
# groupadd --system clickhouse
# useradd --system --gid clickhouse --shell /usr/sbin/nologin --home-dir /var/lib/clickhouse clickhouse
Передаем clickhouse права на директорию:
# chown -R clickhouse:clickhouse /var/lib/clickhouse/
# chmod -R 750 /var/lib/clickhouse/
Создаем каталог логов
# mkdir -p /var/log/clickhouse-server
# chown clickhouse:clickhouse /var/log/clickhouse-server
# chmod 0750 /var/log/clickhouse-server
Запускаем сервис:
# service clickhouse-server start
Проверка:
# systemctl status clickhouse-server
# ss -tlnp | grep clickhouse
Подключение к серверу с паролем:
# clickhouse-client --query "SELECT 1" --password
Если что-то идет не так смотрим логи:
# journalctl -u clickhouse-server -b -n 100 --no-pager
# tail -f /var/log/clickhouse-server/clickhouse-server.err.log
Вносим Clickhouse в автозагрузку
# systemctl enable clickhouse-server
*** Установка Kafka
# apt update
Устанавливаем Java
# apt install -y openjdk-21-jdk
# java --version
openjdk 21.0.12.1 2026-08-18
Затем скачать и распаковать Kafka. Для работы с Clickhouse 22.3 нет возможности использовать новую версию Kafka. Используем версию 3.7.1, которая совместима с Clickhouse 22.3.
# wget https://archive.apache.org/dist/kafka/3.7.1/kafka_2.12-3.7.1.tgz
# tar -xzf kafka_2.12-3.7.1.tgz -C /opt/
# mv /opt/kafka_2.12-3.7.1 /opt/kafka
Конфигурации kafka тут: /opt/kafka/config/
Нас интересует конфигурация в режиме kraft: /opt/kafka/config/kraft/server.properties
# nano /opt/kafka/config/kraft/server.properties
Тут важно указать директорию, что бы данные были именно в ней и не забивали нужный диск.
Выбираем директорию /home/kafka, так как тут есть место.
Так же необходимо указать формат взаимодействия с Clickhouse - 3.0
Дополнительно указываем размер топиков на дисков и время жизни топиков:
log.dirs=/home/kafka
log.message.format.version=3.0
# Хранить данные не более 24 часов
log.retention.hours=24
# Ограничение на один топик - 10G
log.segment.bytes=10737418240
# Оганичение размера сегмента
log.segment.bytes=536870912 # 512MB
# Время жизни сегмента
log.segment.ms=3600000 # 1 час
Создаем директорию для данных, которую определили в конфигурации и директорию для логов:
# mkdir -p /home/kafka
# chmod -R 755 /home/kafka
# mkdir -p /var/log/kafka
Генерация идентификатора кластера Kafka и форматирование директории с даными Kafka
# CLUSTER_ID=$(/opt/kafka/bin/kafka-storage.sh random-uuid)
# /opt/kafka/bin/kafka-storage.sh format \
-t $CLUSTER_ID \
-c /opt/kafka/config/kraft/server.properties
Создаем пользователя под которым будет работать kafka и передаем ему права на директории:
# useradd -r -s /bin/false -d /opt/kafka kafka
# chown -R kafka:kafka /opt/kafka /var/log/kafka /home/kafka
Сервис systemd
# nano /etc/systemd/system/kafka.service
Содержание файла:
[Unit]
Description=Apache Kafka Server (KRaft mode)
Documentation=http://kafka.apache.org/documentation.html
Requires=network.target remote-fs.target
After=network.target remote-fs.target
[Service]
Type=simple
User=kafka
Group=kafka
# Переменные окружения
Environment="JAVA_HOME=/usr/lib/jvm/java-21-openjdk-amd64"
Environment="KAFKA_HEAP_OPTS=-Xmx1G -Xms1G"
Environment="KAFKA_OPTS=-Djava.net.preferIPv4Stack=true"
Environment="LOG_DIR=/var/log/kafka"
# Команда запуска
ExecStart=/opt/kafka/bin/kafka-server-start.sh /opt/kafka/config/kraft/server.properties
# Команда остановки (мягкая)
ExecStop=/opt/kafka/bin/kafka-server-stop.sh
# Политика перезапуска
Restart=on-failure
RestartSec=10
# Лимиты
LimitNOFILE=100000
LimitNPROC=100000
[Install]
WantedBy=multi-user.target
# systemctl daemon-reload
# systemctl start kafka
Проверка:
# systemctl status kafka
# journalctl -u kafka -f -n 100
# ss -tlnp | grep -E '9092|9093'
# systemctl enable kafka
# apt-get install -y apt-transport-https ca-certificates curl gnupg
# cd /usr/src/
У меня под рукой был старый сервер с процессором Intel Xeon E5620, который не поддерживает современные версии Clickhouse.
Процессор Intel Xeon E5620 относится к семейству Westmere (2010 год). Поддерживает инструкции SSE4.2, но не поддерживает AVX/AVX2. Это важно, так как современные версии ClickHouse требуют AVX2.
Оптимальным выбором будет 22.3 LTS, которые все еще поддерживаются и работают с SSE4.2.
# wget https://packages.clickhouse.com/tgz/lts/clickhouse-common-static-22.3.19.6-amd64.tgz
# wget https://packages.clickhouse.com/tgz/lts/clickhouse-server-22.3.19.6-amd64.tgz
# wget https://packages.clickhouse.com/tgz/lts/clickhouse-client-22.3.19.6-amd64.tgz
Установка:
# tar -xzvf clickhouse-common-static-22.3.19.6-amd64.tgz
# ./clickhouse-common-static-22.3.19.6/install/doinst.sh
# tar -xzvf clickhouse-server-22.3.19.6-amd64.tgz
# ./clickhouse-server-22.3.19.6/install/doinst.sh
# tar -xzvf clickhouse-client-22.3.19.6-amd64.tgz
# ./clickhouse-client-22.3.19.6/install/doinst.sh
Настройка Clickhouse
1. Задать пользовательские настройки
# nano /etc/clickhouse-server/users.xml
1.1. Задать пароль по умолчанию для доступа к Clickhouse
Секция <users> -> <default> -> <password>
Прописываем пароль:
<password>*******</password>
1.2. Задать ограничения по использованию памяти и ядер процессора
Секция <profiles> -> <default>
Размер памяти определяем как 50% от памяти, которую выделим Clickhouse (в моем случае - 32 ГБ)
Для сервера с 16 ядрами и выделенными 32 ГБ RAM под ClickHouse, рекомендация max_threads = 4.
<max_memory_usage>17179869184</max_memory_usage>
<max_threads>4</max_threads>
2.Правим конфигурацию сервера:
# nano /etc/clickhouse-server/config.xml
2.1. Отключаем использование IPv6, раскомментируем строку (тем самым отключаем IPv6)
<listen_host>0.0.0.0</listen_host>
2.2. Изменяем уровень логирования – оставляем только ошибки
Так же устанавливаем размер лог-файла в 100Мбайт (вместо 1G)
<logger>
<level>warning</level>
<size>100M</size>
<logger>
2.3. Ограничиваем сервер полкой в использовании памяти - в 32G, иначе сервер съест другие приложения
<max_server_memory_usage>34359738368</max_server_memory_usage>
2.4. Отключаем процентное указание в использование памяти
<max_server_memory_usage_to_ram_ratio>0</max_server_memory_usage_to_ram_ratio>
2.5. Проверить, что Clickhouse работает с директорией, которая для этого была создана в примонтированном хранилище:
<path>/var/lib/clickhouse/</path>
2.6. Из блока <openSSL> вырезать раздел <server>, оставить только раздел клиента
<openSSL>
<client>
...
</client>
</openSSL>
2.7. Добавляем в конец конфигурации перед закрывающим тегом </clickhouse> директивы, определяющие протокол взаимодействия между Kafka и Clickhouse.
Это необходимо, так как версия 22.3 не совместима с новыми сервисами Kafka.
<kafka>
<api_version_request>false</api_version_request>
<broker_version_fallback>2.4.1</broker_version_fallback>
</kafka>
</clickhouse>
Создаем пользователя clickhouse и передаем ему права на директорию, в которой он должен работать
# groupadd --system clickhouse
# useradd --system --gid clickhouse --shell /usr/sbin/nologin --home-dir /var/lib/clickhouse clickhouse
Передаем clickhouse права на директорию:
# chown -R clickhouse:clickhouse /var/lib/clickhouse/
# chmod -R 750 /var/lib/clickhouse/
Создаем каталог логов
# mkdir -p /var/log/clickhouse-server
# chown clickhouse:clickhouse /var/log/clickhouse-server
# chmod 0750 /var/log/clickhouse-server
Запускаем сервис:
# service clickhouse-server start
Проверка:
# systemctl status clickhouse-server
# ss -tlnp | grep clickhouse
Подключение к серверу с паролем:
# clickhouse-client --query "SELECT 1" --password
Если что-то идет не так смотрим логи:
# journalctl -u clickhouse-server -b -n 100 --no-pager
# tail -f /var/log/clickhouse-server/clickhouse-server.err.log
Вносим Clickhouse в автозагрузку
# systemctl enable clickhouse-server
*** Установка Kafka
# apt update
Устанавливаем Java
# apt install -y openjdk-21-jdk
# java --version
openjdk 21.0.12.1 2026-08-18
Затем скачать и распаковать Kafka. Для работы с Clickhouse 22.3 нет возможности использовать новую версию Kafka. Используем версию 3.7.1, которая совместима с Clickhouse 22.3.
# wget https://archive.apache.org/dist/kafka/3.7.1/kafka_2.12-3.7.1.tgz
# tar -xzf kafka_2.12-3.7.1.tgz -C /opt/
# mv /opt/kafka_2.12-3.7.1 /opt/kafka
Конфигурации kafka тут: /opt/kafka/config/
Нас интересует конфигурация в режиме kraft: /opt/kafka/config/kraft/server.properties
# nano /opt/kafka/config/kraft/server.properties
Тут важно указать директорию, что бы данные были именно в ней и не забивали нужный диск.
Выбираем директорию /home/kafka, так как тут есть место.
Так же необходимо указать формат взаимодействия с Clickhouse - 3.0
Дополнительно указываем размер топиков на дисков и время жизни топиков:
log.dirs=/home/kafka
log.message.format.version=3.0
# Хранить данные не более 24 часов
log.retention.hours=24
# Ограничение на один топик - 10G
log.segment.bytes=10737418240
# Оганичение размера сегмента
log.segment.bytes=536870912 # 512MB
# Время жизни сегмента
log.segment.ms=3600000 # 1 час
Создаем директорию для данных, которую определили в конфигурации и директорию для логов:
# mkdir -p /home/kafka
# chmod -R 755 /home/kafka
# mkdir -p /var/log/kafka
Генерация идентификатора кластера Kafka и форматирование директории с даными Kafka
# CLUSTER_ID=$(/opt/kafka/bin/kafka-storage.sh random-uuid)
# /opt/kafka/bin/kafka-storage.sh format \
-t $CLUSTER_ID \
-c /opt/kafka/config/kraft/server.properties
Создаем пользователя под которым будет работать kafka и передаем ему права на директории:
# useradd -r -s /bin/false -d /opt/kafka kafka
# chown -R kafka:kafka /opt/kafka /var/log/kafka /home/kafka
Сервис systemd
# nano /etc/systemd/system/kafka.service
Содержание файла:
[Unit]
Description=Apache Kafka Server (KRaft mode)
Documentation=http://kafka.apache.org/documentation.html
Requires=network.target remote-fs.target
After=network.target remote-fs.target
[Service]
Type=simple
User=kafka
Group=kafka
# Переменные окружения
Environment="JAVA_HOME=/usr/lib/jvm/java-21-openjdk-amd64"
Environment="KAFKA_HEAP_OPTS=-Xmx1G -Xms1G"
Environment="KAFKA_OPTS=-Djava.net.preferIPv4Stack=true"
Environment="LOG_DIR=/var/log/kafka"
# Команда запуска
ExecStart=/opt/kafka/bin/kafka-server-start.sh /opt/kafka/config/kraft/server.properties
# Команда остановки (мягкая)
ExecStop=/opt/kafka/bin/kafka-server-stop.sh
# Политика перезапуска
Restart=on-failure
RestartSec=10
# Лимиты
LimitNOFILE=100000
LimitNPROC=100000
[Install]
WantedBy=multi-user.target
# systemctl daemon-reload
# systemctl start kafka
Проверка:
# systemctl status kafka
# journalctl -u kafka -f -n 100
# ss -tlnp | grep -E '9092|9093'
# systemctl enable kafka
=>>> ПРОБНЫЕ ЗАПИСИ <<<=
=>>> Тестовый топик
=>>> # /opt/kafka/bin/kafka-topics.sh --create \
--topic test-topic \
--bootstrap-server localhost:9092 \
--partitions 3 \
--replication-factor 1
=>>> Список топиков
=>>> # /opt/kafka/bin/kafka-topics.sh --list --bootstrap-server localhost:9092
=>>> Тестовое сообщение
=>>> # /opt/kafka/bin/kafka-console-producer.sh --topic test-topic --bootstrap-server localhost:9092
Теперь вводим несколько сообщений и для выхода Ctrl+C
=>>> Информация о топике
=>>> # /opt/kafka/bin/kafka-topics.sh --describe --topic test-topic --bootstrap-server localhost:9092
=>>> Принят сообщение:
=>>> # /opt/kafka/bin/kafka-console-consumer.sh --topic test-topic --from-beginning --bootstrap-server localhost:9092
=>>> Список всех топиков:
=>>> # /opt/kafka/bin/kafka-topics.sh --bootstrap-server localhost:9092 --list
=>>> Удаление тестового топика:
=>>> # /opt/kafka/bin/kafka-topics.sh --delete --topic test-topic --bootstrap-server localhost:9092
=>>> Тестовый топик
=>>> # /opt/kafka/bin/kafka-topics.sh --create \
--topic test-topic \
--bootstrap-server localhost:9092 \
--partitions 3 \
--replication-factor 1
=>>> Список топиков
=>>> # /opt/kafka/bin/kafka-topics.sh --list --bootstrap-server localhost:9092
=>>> Тестовое сообщение
=>>> # /opt/kafka/bin/kafka-console-producer.sh --topic test-topic --bootstrap-server localhost:9092
Теперь вводим несколько сообщений и для выхода Ctrl+C
=>>> Информация о топике
=>>> # /opt/kafka/bin/kafka-topics.sh --describe --topic test-topic --bootstrap-server localhost:9092
=>>> Принят сообщение:
=>>> # /opt/kafka/bin/kafka-console-consumer.sh --topic test-topic --from-beginning --bootstrap-server localhost:9092
=>>> Список всех топиков:
=>>> # /opt/kafka/bin/kafka-topics.sh --bootstrap-server localhost:9092 --list
=>>> Удаление тестового топика:
=>>> # /opt/kafka/bin/kafka-topics.sh --delete --topic test-topic --bootstrap-server localhost:9092
Теперь создаем топики pmacct-flows и pmacct-nel для приема сообщений Netflow (сам поток Netflow и NEL):
# /opt/kafka/bin/kafka-topics.sh --create \
--topic pmacct-flows \
--bootstrap-server localhost:9092 \
--partitions 3 \
--replication-factor 1
# /opt/kafka/bin/kafka-topics.sh --create \
--topic pmacct-nel \
--bootstrap-server localhost:9092 \
--partitions 3 \
--replication-factor 1
*** Сборка и установка pmacct с поддержкой kafka
# /usr/src/
# apt install git
# apt update && sudo apt install -y libtool automake autoconf
# apt install -y pkgconf build-essential
# apt install -y libpcap-dev
# apt install -y libcurl4-openssl-dev
# apt install -y libmariadb-dev-compat libmariadb-dev
# apt install -y libnuma-dev
# apt install -y librdkafka-dev
# apt install -y libjansson-dev
# git clone https://github.com/pmacct/pmacct.git
# cd pmacct/
# ./autogen.sh
# ./configure --enable-jansson --enable-kafka --disable-bgp-bins --disable-bmp-bins --disable-st-bins
# make -j$(nproc)
# make install
# nfacctd -V
В выводе мы должны увидеть версию программы и информацию о том, что плагин kafka подключен
Libs:
cdada 0.6.4
libpcap version 1.10.5 (with TPACKET_V3)
rdkafka 2.8.0
jansson 2.14
Plugins:
...
kafka
Создаем директорию для логов:
# mkdir -p /var/log/pmacct
Конфиги лежат тут: /etc/pmacct/
Создаем конфиг для приема основного потока Netflow: pmacct-flows.conf
# nano /etc/pmacct/pmacct-flows.conf
##### Для теста с выводом в консоль:
#daemonize: false
#pidfile: /var/run/pmacct-flows.pid
#logfile: /var/log/pmacct/pmacct-flows.log
#debug: true
#nfacctd_port: 2055
#nfacctd_time_new: true
#plugins: print[flows]
#aggregate[flows]: src_host, dst_host, src_port, dst_port, proto, timestamp_start, peer_src_ip
#print_output[flows]: json
#print_refresh_time[flows]: 5
##### Для записи в Kafka
#daemonize: true # Для запуска без systemd
daemonize: false
pidfile: /var/run/pmacct-flows.pid
logfile: /var/log/pmacct/pmacct-flows.log
debug: false
nfacctd_port: 2055
nfacctd_time_new: true
plugins: kafka[flows]
aggregate[flows]: src_host, dst_host, src_port, dst_port, proto, timestamp_start, peer_src_ip
kafka_broker_host[flows]: localhost
kafka_broker_port[flows]: 9092
kafka_topic[flows]: pmacct-flows
kafka_partition_key[flows]: src_host
kafka_refresh_time[flows]: 10
[Пробный запуск /usr/local/sbin/nfacctd -f /etc/pmacct/pmacct-flows.conf ]
# /opt/kafka/bin/kafka-topics.sh --create \
--topic pmacct-flows \
--bootstrap-server localhost:9092 \
--partitions 3 \
--replication-factor 1
# /opt/kafka/bin/kafka-topics.sh --create \
--topic pmacct-nel \
--bootstrap-server localhost:9092 \
--partitions 3 \
--replication-factor 1
*** Сборка и установка pmacct с поддержкой kafka
# /usr/src/
# apt install git
# apt update && sudo apt install -y libtool automake autoconf
# apt install -y pkgconf build-essential
# apt install -y libpcap-dev
# apt install -y libcurl4-openssl-dev
# apt install -y libmariadb-dev-compat libmariadb-dev
# apt install -y libnuma-dev
# apt install -y librdkafka-dev
# apt install -y libjansson-dev
# git clone https://github.com/pmacct/pmacct.git
# cd pmacct/
# ./autogen.sh
# ./configure --enable-jansson --enable-kafka --disable-bgp-bins --disable-bmp-bins --disable-st-bins
# make -j$(nproc)
# make install
# nfacctd -V
В выводе мы должны увидеть версию программы и информацию о том, что плагин kafka подключен
Libs:
cdada 0.6.4
libpcap version 1.10.5 (with TPACKET_V3)
rdkafka 2.8.0
jansson 2.14
Plugins:
...
kafka
Создаем директорию для логов:
# mkdir -p /var/log/pmacct
Конфиги лежат тут: /etc/pmacct/
Создаем конфиг для приема основного потока Netflow: pmacct-flows.conf
# nano /etc/pmacct/pmacct-flows.conf
##### Для теста с выводом в консоль:
#daemonize: false
#pidfile: /var/run/pmacct-flows.pid
#logfile: /var/log/pmacct/pmacct-flows.log
#debug: true
#nfacctd_port: 2055
#nfacctd_time_new: true
#plugins: print[flows]
#aggregate[flows]: src_host, dst_host, src_port, dst_port, proto, timestamp_start, peer_src_ip
#print_output[flows]: json
#print_refresh_time[flows]: 5
##### Для записи в Kafka
#daemonize: true # Для запуска без systemd
daemonize: false
pidfile: /var/run/pmacct-flows.pid
logfile: /var/log/pmacct/pmacct-flows.log
debug: false
nfacctd_port: 2055
nfacctd_time_new: true
plugins: kafka[flows]
aggregate[flows]: src_host, dst_host, src_port, dst_port, proto, timestamp_start, peer_src_ip
kafka_broker_host[flows]: localhost
kafka_broker_port[flows]: 9092
kafka_topic[flows]: pmacct-flows
kafka_partition_key[flows]: src_host
kafka_refresh_time[flows]: 10
[Пробный запуск /usr/local/sbin/nfacctd -f /etc/pmacct/pmacct-flows.conf ]
Создаем unit systemd для запуска через systemd
# nano /etc/systemd/system/pmacct-flows.service
[Unit]
Description=pmacct NetFlow Accounting Daemon
Documentation=http://www.pmacct.net/
After=network.target kafka.service
Wants=kafka.service
[Service]
Type=simple
User=root
Group=root
ExecStart=/usr/local/sbin/nfacctd -f /etc/pmacct/pmacct-flows.conf
Restart=on-failure
RestartSec=5
LimitNOFILE=65536
[Install]
WantedBy=multi-user.target
# nano /etc/systemd/system/pmacct-flows.service
[Unit]
Description=pmacct NetFlow Accounting Daemon
Documentation=http://www.pmacct.net/
After=network.target kafka.service
Wants=kafka.service
[Service]
Type=simple
User=root
Group=root
ExecStart=/usr/local/sbin/nfacctd -f /etc/pmacct/pmacct-flows.conf
Restart=on-failure
RestartSec=5
LimitNOFILE=65536
[Install]
WantedBy=multi-user.target
Создаем конфиг для приема потока NEL: pmacct-nel.conf
# nano /etc/pmacct/pmacct-nel.conf
##### Для теста с выводом в консоль:
#daemonize: false
#pidfile: /var/run/pmacct-nel.pid
#logfile: /var/log/pmacct/pmacct-nel.log
#debug: true
#nfacctd_port: 2056
#nfacctd_time_new: true
#plugins: print[cgn]
#aggregate[cgn]: src_host, post_nat_src_host, src_port, post_nat_src_port, proto, nat_event, timestamp_start, peer_src_ip
#aggregate_filter[cgn]: dst_host == ""
#print_output[cgn]: json
#print_refresh_time[cgn]: 5
#print_output_file[cgn]: /dev/stdout
##### Для записи в Kafka
#daemonize: true
daemonize: false
pidfile: /var/run/pmacct-nel.pid
logfile: /var/log/pmacct/pmacct-nel.log
debug: false
nfacctd_port: 2056
nfacctd_time_new: true
plugins: kafka[cgn]
aggregate[cgn]: src_host, post_nat_src_host, src_port, post_nat_src_port, proto, nat_event, timestamp_start, peer_src_ip
aggregate_filter[cgn]: dst_host == ""
kafka_broker_host[cgn]: localhost
kafka_broker_port[cgn]: 9092
kafka_topic[cgn]: pmacct-nel
kafka_partition_key[cgn]: src_host
kafka_refresh_time[cgn]: 10
[Пробный запуск /usr/local/sbin/nfacctd -f /etc/pmacct/pmacct-nel.conf ]
# nano /etc/pmacct/pmacct-nel.conf
##### Для теста с выводом в консоль:
#daemonize: false
#pidfile: /var/run/pmacct-nel.pid
#logfile: /var/log/pmacct/pmacct-nel.log
#debug: true
#nfacctd_port: 2056
#nfacctd_time_new: true
#plugins: print[cgn]
#aggregate[cgn]: src_host, post_nat_src_host, src_port, post_nat_src_port, proto, nat_event, timestamp_start, peer_src_ip
#aggregate_filter[cgn]: dst_host == ""
#print_output[cgn]: json
#print_refresh_time[cgn]: 5
#print_output_file[cgn]: /dev/stdout
##### Для записи в Kafka
#daemonize: true
daemonize: false
pidfile: /var/run/pmacct-nel.pid
logfile: /var/log/pmacct/pmacct-nel.log
debug: false
nfacctd_port: 2056
nfacctd_time_new: true
plugins: kafka[cgn]
aggregate[cgn]: src_host, post_nat_src_host, src_port, post_nat_src_port, proto, nat_event, timestamp_start, peer_src_ip
aggregate_filter[cgn]: dst_host == ""
kafka_broker_host[cgn]: localhost
kafka_broker_port[cgn]: 9092
kafka_topic[cgn]: pmacct-nel
kafka_partition_key[cgn]: src_host
kafka_refresh_time[cgn]: 10
[Пробный запуск /usr/local/sbin/nfacctd -f /etc/pmacct/pmacct-nel.conf ]
Создаем unit systemd для запуска через systemd
# nano /etc/systemd/system/pmacct-nel.service
[Unit]
Description=pmacct NEL
Documentation=http://www.pmacct.net/
After=network.target kafka.service
Wants=kafka.service
[Service]
Type=simple
User=root
Group=root
ExecStart=/usr/local/sbin/nfacctd -f /etc/pmacct/pmacct-nel.conf
Restart=on-failure
RestartSec=5
LimitNOFILE=65536
[Install]
WantedBy=multi-user.target
Запускаем:
[Если запускал руками, то сначала нужно убить процессы nfacctd]
# systemctl daemon-reload
# systemctl start pmacct-flows
# systemctl status pmacct-flows
# systemctl start pmacct-nel
# systemctl status pmacct-nel
Проверка того, что данные залетают в Kafka:
# ps aux | grep nfacctd
# /opt/kafka/bin/kafka-console-consumer.sh \
--bootstrap-server localhost:9092 \
--topic pmacct-flows \
--from-beginning \
--max-messages 5
# /opt/kafka/bin/kafka-console-consumer.sh \
--bootstrap-server localhost:9092 \
--topic pmacct-nel \
--from-beginning \
--max-messages 5
# systemctl enable pmacct-flows
# systemctl enable pmacct-nel
** Подключаем Kafka к Clickhouse
1 часть. Поток Netflow.
# clickhouse-client --password (нужен пароль)
1. Создаем базу данных:
CREATE DATABASE IF NOT EXISTS netflow;
2. Создание Kafka Engine таблицы для подключения к Kafka (удаление старой таблицы DROP TABLE IF EXISTS netflow.flows_kafka;)
CREATE TABLE netflow.flows_kafka
(
`value` String
)
ENGINE = Kafka()
SETTINGS
kafka_broker_list = 'localhost:9092',
kafka_topic_list = 'pmacct-flows',
kafka_group_name = 'clickhouse-flows',
kafka_format = 'JSONAsString';
3. Создание таблицы хранения данных: (удаление старой таблицы DROP TABLE IF EXISTS netflow.flows;)
CREATE TABLE netflow.flows
(
`timestamp_start` DateTime CODEC(Delta(4), ZSTD(1)),
`ip_src` IPv4,
`ip_dst` IPv4,
`port_src` UInt16 DEFAULT 0 CODEC(ZSTD(1)),
`port_dst` UInt16 DEFAULT 0 CODEC(ZSTD(1)),
`proto` LowCardinality(String) DEFAULT '' CODEC(ZSTD(1)),
`peer_ip_src` IPv4,
`bytes` UInt64 DEFAULT 0 CODEC(ZSTD(1)),
`packets` UInt32 DEFAULT 0 CODEC(ZSTD(1))
)
ENGINE = MergeTree
PARTITION BY toYYYYMMDD(timestamp_start)
ORDER BY (timestamp_start, ip_src, ip_dst)
TTL timestamp_start + toIntervalMonth(6)
SETTINGS index_granularity = 8192;
4. Создание Materialized View для форматирования принимаемых из Kafka данных (удаление старой таблицы DROP TABLE IF EXISTS netflow.flows_mv;)
CREATE MATERIALIZED VIEW netflow.flows_mv TO netflow.flows AS
SELECT
toDateTime(substring(JSONExtractString(value, 'timestamp_start'), 1, 19)) AS timestamp_start,
toIPv4(nullIf(JSONExtractString(value, 'ip_src'), '')) AS ip_src,
toIPv4(nullIf(JSONExtractString(value, 'ip_dst'), '')) AS ip_dst,
toUInt16(JSONExtractInt(value, 'port_src')) AS port_src,
toUInt16(JSONExtractInt(value, 'port_dst')) AS port_dst,
JSONExtractString(value, 'ip_proto') AS proto,
toIPv4(nullIf(JSONExtractString(value, 'peer_ip_src'), '')) AS peer_ip_src,
toUInt64(JSONExtractInt(value, 'bytes')) AS bytes,
toUInt32(JSONExtractInt(value, 'packets')) AS packets
FROM netflow.flows_kafka
WHERE (JSONExtractString(value, 'ip_src') != '') AND (JSONExtractString(value, 'ip_dst') != '')
Проверка логов на предмет чего-то критичного:
# tail -n 20 /var/log/clickhouse-server/clickhouse-server.err.log | grep -i kafka
2 часть. Поток NEL
5. Создание Kafka Engine таблицы для подключения к Kafka (удаление старой таблицы DROP TABLE IF EXISTS netflow.nel_kafka;)
CREATE TABLE IF NOT EXISTS netflow.nel_kafka
(
`value` String
)
ENGINE = Kafka()
SETTINGS
kafka_broker_list = 'localhost:9092',
kafka_topic_list = 'pmacct-nel',
kafka_group_name = 'clickhouse-nel',
kafka_format = 'JSONAsString';
6. Создание таблицы хранения данных NEL (удаление старой таблицы DROP TABLE IF EXISTS netflow.nel;)
CREATE TABLE netflow.nel
(
`timestamp_start` DateTime CODEC(Delta(4), ZSTD(1)),
`ip_src` IPv4,
`port_src` UInt16 DEFAULT 0 CODEC(ZSTD(1)),
`proto` LowCardinality(String) DEFAULT '' CODEC(ZSTD(1)),
`peer_ip_src` IPv4,
`post_nat_ip_src` IPv4,
`post_nat_port_src` UInt16 DEFAULT 0 CODEC(ZSTD(1)),
`nat_event` UInt8 DEFAULT 0 CODEC(ZSTD(1))
)
ENGINE = MergeTree
PARTITION BY toYYYYMMDD(timestamp_start)
ORDER BY (timestamp_start, ip_src, post_nat_ip_src)
TTL timestamp_start + toIntervalMonth(6)
SETTINGS index_granularity = 8192;
7. Создание Materialized View для форматирования принимаемых из Kafka данных (удаление старой таблицы DROP TABLE IF EXISTS netflow.nel_mv;)
CREATE MATERIALIZED VIEW IF NOT EXISTS netflow.nel_mv TO netflow.nel AS
SELECT
toDateTime(substring(JSONExtractString(value, 'timestamp_start'), 1, 19)) AS timestamp_start,
toIPv4(nullIf(JSONExtractString(value, 'ip_src'), '')) AS ip_src,
toUInt16(JSONExtractInt(value, 'port_src')) AS port_src,
JSONExtractString(value, 'ip_proto') AS proto,
toIPv4(nullIf(JSONExtractString(value, 'peer_ip_src'), '')) AS peer_ip_src,
toIPv4(nullIf(JSONExtractString(value, 'post_nat_ip_src'), '')) AS post_nat_ip_src,
toUInt16(JSONExtractInt(value, 'post_nat_port_src')) AS post_nat_port_src,
toUInt8(JSONExtractInt(value, 'nat_event')) AS nat_event
FROM netflow.nel_kafka
WHERE
JSONExtractString(value, 'ip_src') != '' AND
JSONExtractString(value, 'post_nat_ip_src') != '';
*** Включаем firewall
# apt install ufw
# ufw allow ssh
# ufw allow 2055/udp
# ufw allow 2056/udp
# ufw allow from 10.0.0.0/8 to any port 8123 proto tcp
# ufw allow from 10.0.0.0/8 to any port 9000 proto tcp
# ufw enable
Проверка:
# ufw status verbose
# systemctl status ufw
*** Включение ротации логов
# apt install ufw
# ufw allow ssh
# ufw allow 2055/udp
# ufw allow 2056/udp
# ufw allow from 10.0.0.0/8 to any port 8123 proto tcp
# ufw allow from 10.0.0.0/8 to any port 9000 proto tcp
# ufw enable
Проверка:
# ufw status verbose
# systemctl status ufw
*** Включение ротации логов
У сlickhouse ротацию логов уже включили через конфиг при насройке clickhouse
Логи демона pmacct храняться тут /var/log/pmacct
Используем стандартный logrotate для ротации логов
# nano /etc/logrotate.d/pmacct
Содержимое файла:
/var/log/pmacct/*.log {
hourly
rotate 2
compress
delaycompress
missingok
notifempty
copytruncate
}
Этот конфиг копирует текущий лог в архив и очищает исходный файл. pmacct продолжает писать в тот же файловый дескриптор, поэтому перезапуск сервиса не нужен.
Проверка синтаксиса:
# logrotate -d /etc/logrotate.d/pmacct
Принудительная ротация логов прямо сейчас:
# logrotate -f /etc/logrotate.d/pmacct
Ротация будет проводиться раз в сутки. Когда следующая ротация покажет команда:
# systemctl list-timers logrotate.timer
Используем стандартный logrotate для ротации логов
# nano /etc/logrotate.d/pmacct
Содержимое файла:
/var/log/pmacct/*.log {
hourly
rotate 2
compress
delaycompress
missingok
notifempty
copytruncate
}
Этот конфиг копирует текущий лог в архив и очищает исходный файл. pmacct продолжает писать в тот же файловый дескриптор, поэтому перезапуск сервиса не нужен.
Проверка синтаксиса:
# logrotate -d /etc/logrotate.d/pmacct
Принудительная ротация логов прямо сейчас:
# logrotate -f /etc/logrotate.d/pmacct
Ротация будет проводиться раз в сутки. Когда следующая ротация покажет команда:
# systemctl list-timers logrotate.timer
Логи kafka расположены в директории /var/log/kafka
Kafka использует log4j для логирования. Файл конфигурации логирования: /opt/kafka/config/log4j.properties
А вот удаление старых файлов kafka делать не умеет.
Для удаления старых файлов используем команду в crontab:
# Очищение логов kafka старше 5 суток
42 * * * * root find /var/log/kafka -type f -name "*.log.20*" -mmin +7200 -delete > /dev/null 2>&1
Kafka использует log4j для логирования. Файл конфигурации логирования: /opt/kafka/config/log4j.properties
А вот удаление старых файлов kafka делать не умеет.
Для удаления старых файлов используем команду в crontab:
# Очищение логов kafka старше 5 суток
42 * * * * root find /var/log/kafka -type f -name "*.log.20*" -mmin +7200 -delete > /dev/null 2>&1

Комментариев нет:
Отправить комментарий