diff --git a/.gitignore b/.gitignore new file mode 100644 index 0000000..e3ff419 --- /dev/null +++ b/.gitignore @@ -0,0 +1,47 @@ +# Compiled files +target/ +*.class + +# Maven +pom.xml.tag +pom.xml.releaseBackup +pom.xml.versionsBackup +release.properties + +# IDE specific +.idea/ +*.iml +*.iws +.idea_modules/ +out/ +.settings/ +.project +.classpath +.vscode/ + +# OS generated files +.DS_Store +.DS_Store? +._* +.Spotlight-V100 +.Trashes +Thumbs.db +ehthumbs.db +*.swp +*.swo + +# Lock files (created by application) +app.lock +*.lock + +# Logs +*.log + +# Temporary files +*.tmp +*.temp + +# Local configuration (if you want to keep your own copies) +# config.properties and mapping.json might be kept, but usually you want to keep them +# If you have local versions, you can ignore them, but the repo should contain sample files. +# So we don't ignore them by default. \ No newline at end of file diff --git a/README.md b/README.md new file mode 100644 index 0000000..18dbc14 --- /dev/null +++ b/README.md @@ -0,0 +1,259 @@ +# MQTT to Zabbix Bridge + +**MQTT to Zabbix Bridge** — это легковесное Java-приложение, которое подписывается на MQTT-топики, обрабатывает входящие сообщения и отправляет данные в Zabbix через протокол Trapper (используя утилиту `zabbix_sender`). Программа поддерживает гибкий маппинг топиков на элементы данных Zabbix, автоматическое переподключение к MQTT-брокеру, мониторинг собственного состояния и может быть запущена как systemd-сервис. + +## Возможности + +- **MQTT-подписка** — подписка на один или несколько топиков (поддерживаются wildcard-шаблоны `+` и `#`). +- **Гибкий маппинг** — сопоставление входящих топиков с хостами и ключами Zabbix с помощью регулярных выражений и подстановок (`$1`, `$2`…). +- **Извлечение значений** — из JSON-полей (по имени поля) или использование всего payload как значения. +- **Подписка на конкретные топики** — можно указать для каждого правила отдельный топик для подписки (если не указан, используется глобальный). +- **Автоматическое переподключение к MQTT** — при потере соединения выполняется повторное подключение с экспоненциальной задержкой. +- **Повторные попытки отправки в Zabbix** — при ошибках отправки данные повторяются до указанного числа раз с экспоненциальной задержкой. +- **Мониторинг работы программы** — каждые N секунд в Zabbix отправляются метрики: количество обработанных сообщений, ошибок, повторных попыток и размер очереди. +- **Валидация конфигурации при старте** — проверяются наличие `zabbix_sender`, обязательные параметры и корректность правил маппинга. +- **Блокировка двойного запуска** — предотвращает одновременный запуск нескольких экземпляров. +- **Поддержка systemd** — готовый unit-файл для запуска как сервиса с автоматическим перезапуском. +- **Внешняя конфигурация** — файлы `config.properties` и `mapping.json` могут лежать рядом с JAR-файлом или внутри него (для удобства изменения без пересборки). + +## Требования + +- **Java 11** или выше. +- **Maven** (для сборки). +- **Утилита `zabbix_sender`** должна быть установлена и доступна в `PATH` (обычно входит в пакет `zabbix-sender`). Используется для отправки данных. +- **MQTT-брокер** (например, Mosquitto, EMQX и др.). +- **Zabbix-сервер** (с настроенным trapper-портом, по умолчанию 10051). + +## Установка и сборка + +### 1. Клонирование репозитория + +```bash +git clone https://gitbucket.mfnd.ru/git/y.varenkov/MQTT-to-Zabbix-Bridge.git +cd mqtt-to-zabbix +``` + +### 2. Сборка проекта + +```bash +mvn clean package +``` + +В каталоге `target/` появится файл `mqtt-to-zabbix-1.0-SNAPSHOT-jar-with-dependencies.jar` — это самодостаточный JAR со всеми зависимостями. + +### 3. Установка (опционально) + +Скопируйте JAR и файлы конфигурации в рабочую директорию: + +```bash +sudo mkdir -p /opt/mqtt-to-zabbix +sudo cp target/mqtt-to-zabbix-1.0-SNAPSHOT-jar-with-dependencies.jar /opt/mqtt-to-zabbix/ +sudo cp src/main/resources/config.properties /opt/mqtt-to-zabbix/ +sudo cp src/main/resources/mapping.json /opt/mqtt-to-zabbix/ +``` + +Убедитесь, что `zabbix_sender` установлен: + +```bash +sudo apt install zabbix-sender # для Debian/Ubuntu +sudo yum install zabbix-sender # для RHEL/CentOS +``` + +## Настройка + +### Файл `config.properties` + +Основные настройки приложения: + +```properties +# MQTT Broker +mqtt.broker=tcp://localhost:1883 +mqtt.topic=# # глобальный топик для подписки (используется, если не указан subscribeTopic в маппинге) +mqtt.user= +mqtt.password= + +# Zabbix Trapper +zabbix.server=127.0.0.1 +zabbix.port=10051 + +# Пул потоков +thread.pool.core=2 +thread.pool.max=4 +thread.pool.queue=1000 + +# Повторные попытки +retry.max.attempts=3 +retry.initial.delay.ms=1000 +retry.backoff.multiplier=2.0 + +# Мониторинг +metrics.host=MQTT-Bridge # хост в Zabbix для метрик самой программы +metrics.interval=60 # интервал отправки метрик (секунд) +``` + +### Файл `mapping.json` + +Определяет правила преобразования топиков в элементы Zabbix. Формат: + +```json +[ + { + "topicPattern": "sensors/temperature", + "host": "Server1", + "key": "temperature", + "valueField": null, + "subscribeTopic": "sensors/temperature" + }, + { + "topicPattern": "devices/([^/]+)/status", + "host": "$1", + "key": "status", + "valueField": "state", + "subscribeTopic": "devices/+/status" + }, + { + "topicPattern": "metrics/([^/]+)/([^/]+)", + "host": "$1", + "key": "$2", + "valueField": "value" + } +] +``` + +**Поля:** + +- `topicPattern` — регулярное выражение для сопоставления с входящим MQTT-топиком. Может содержать группы, которые затем используются в `$1`, `$2`... +- `host` — имя хоста в Zabbix. Может содержать подстановки `$1`, `$2` из групп регулярного выражения. +- `key` — ключ элемента данных Zabbix. Аналогично поддерживает подстановки. +- `valueField` — имя поля в JSON, из которого берётся значение. Если `null` или отсутствует, используется весь payload как строка. +- `subscribeTopic` — (опционально) MQTT-топик для подписки. Если не указан, используется глобальный `mqtt.topic` из `config.properties` (или `#`). Можно указывать wildcard-символы `+` и `#`. + +### Настройка элементов данных в Zabbix + +Для каждого ключа, указанного в маппинге, необходимо создать элемент данных типа **Zabbix trapper** на соответствующем хосте. В поле **Разрешенные хосты** укажите IP-адрес, с которого будут приходить данные (например, `127.0.0.1`). + +Для метрик мониторинга создайте хост `MQTT-Bridge` (или измените `metrics.host`) и элементы с ключами: + +- `mqtt.bridge.processed` — количество обработанных сообщений (целое) +- `mqtt.bridge.errors` — количество ошибок обработки (целое) +- `mqtt.bridge.retries` — количество повторных попыток отправки (целое) +- `mqtt.bridge.queue_size` — текущий размер очереди задач (целое) + +## Запуск + +### Ручной запуск + +```bash +java -jar /opt/mqtt-to-zabbix/mqtt-to-zabbix-1.0-SNAPSHOT-jar-with-dependencies.jar +``` + +Если файлы конфигурации лежат рядом с JAR, они будут загружены из файловой системы. Если нет — будут использованы встроенные (из classpath). + +### systemd (рекомендуется для продакшена) + +Создайте unit-файл `/etc/systemd/system/mqtt-to-zabbix.service`: + +```ini +[Unit] +Description=MQTT to Zabbix Bridge +After=network.target zabbix-server.service +Wants=zabbix-server.service + +[Service] +Type=simple +User=zabbix +Group=zabbix +WorkingDirectory=/opt/mqtt-to-zabbix +ExecStart=/usr/bin/java -jar /opt/mqtt-to-zabbix/mqtt-to-zabbix-1.0-SNAPSHOT-jar-with-dependencies.jar +Restart=always +RestartSec=10 +StandardOutput=journal +StandardError=journal +SyslogIdentifier=mqtt-to-zabbix + +[Install] +WantedBy=multi-user.target +``` + +После создания выполните: + +```bash +sudo systemctl daemon-reload +sudo systemctl enable mqtt-to-zabbix +sudo systemctl start mqtt-to-zabbix +``` + +Проверьте статус: + +```bash +sudo systemctl status mqtt-to-zabbix +``` + +Логи можно просматривать через: + +```bash +sudo journalctl -u mqtt-to-zabbix -f +``` + +## Примеры работы + +### Пример MQTT-сообщения (JSON) + +**Топик:** `devices/room1/temperature` +**Payload:** `{"value":23.5, "timestamp":"2026-07-10T14:30:00Z"}` + +**Правило:** +```json +{ + "topicPattern": "devices/([^/]+)/temperature", + "host": "$1", + "key": "temperature", + "valueField": "value", + "subscribeTopic": "devices/+/temperature" +} +``` + +**Результат:** В Zabbix для хоста `room1` элемент с ключом `temperature` получит значение `23.5`. + +### Пример MQTT-сообщения (просто число) + +**Топик:** `sensors/temperature` +**Payload:** `20` + +**Правило:** +```json +{ + "topicPattern": "sensors/temperature", + "host": "Server1", + "key": "temp", + "valueField": null, + "subscribeTopic": "sensors/temperature" +} +``` + +**Результат:** В Zabbix для хоста `Server1` элемент `temp` получит значение `20`. + +## Логирование и диагностика + +Программа выводит подробные логи в консоль (или `journalctl` при использовании systemd). Основные события: + +- Получение MQTT-сообщения +- Обработка по правилу или legacy-режим +- Отправка через `zabbix_sender` +- Повторные попытки и ошибки +- Отправка метрик мониторинга + +## Архитектура и принцип работы + +1. **Загрузка конфигурации** — читается `config.properties` и `mapping.json`. +2. **Валидация** — проверяется наличие `zabbix_sender`, корректность параметров. +3. **Подключение к MQTT** — устанавливается соединение с брокером, подписка на топики (из маппинга или глобальный). +4. **Обработка сообщений** — каждое сообщение обрабатывается в отдельном потоке из пула. Поиск подходящего правила, извлечение host, key и value, отправка в Zabbix с повторными попытками. +5. **Переподключение** — при разрыве соединения запускается фоновый поток, который пытается восстановить связь с экспоненциальной задержкой. +6. **Мониторинг** — периодически отправляются метрики работы в Zabbix. +7. **Завершение** — при остановке программы (Ctrl+C или systemd) корректно закрываются все ресурсы. + +## Лицензия + +Этот проект распространяется под лицензией MIT. Подробнее см. файл LICENSE. + diff --git a/pom.xml b/pom.xml new file mode 100644 index 0000000..f691ce8 --- /dev/null +++ b/pom.xml @@ -0,0 +1,71 @@ + + + 4.0.0 + + com.example + mqtt-to-zabbix + 1.0-SNAPSHOT + jar + + + 11 + 11 + UTF-8 + + + + + + org.eclipse.paho + org.eclipse.paho.client.mqttv3 + 1.2.5 + + + + com.fasterxml.jackson.core + jackson-databind + 2.15.2 + + + + + + + + org.apache.maven.plugins + maven-compiler-plugin + 3.13.0 + + 11 + + + + org.apache.maven.plugins + maven-assembly-plugin + 3.5.0 + + + jar-with-dependencies + + + + MqttToZabbix + + + + + + make-assembly + package + + single + + + + + + + \ No newline at end of file diff --git a/src/main/java/MqttToZabbix.java b/src/main/java/MqttToZabbix.java new file mode 100644 index 0000000..7548a17 --- /dev/null +++ b/src/main/java/MqttToZabbix.java @@ -0,0 +1,583 @@ +import com.fasterxml.jackson.databind.JsonNode; +import com.fasterxml.jackson.databind.ObjectMapper; +import org.eclipse.paho.client.mqttv3.*; +import org.eclipse.paho.client.mqttv3.persist.MemoryPersistence; + +import java.io.*; +import java.nio.channels.FileLock; +import java.util.*; +import java.util.concurrent.*; +import java.util.concurrent.atomic.AtomicLong; +import java.util.regex.Matcher; +import java.util.regex.Pattern; + +public class MqttToZabbix { + + // ---------- Конфигурация ---------- + private static String mqttBroker; + private static String mqttTopic; + private static String mqttUser; + private static String mqttPassword; + private static String zabbixServer; + private static int zabbixPort; + private static int corePoolSize; + private static int maxPoolSize; + private static int queueCapacity; + private static int maxRetries; + private static long initialDelayMs; + private static double backoffMultiplier; + + // Параметры мониторинга + private static String metricsHost; + private static int metricsInterval; + + private static final ObjectMapper objectMapper = new ObjectMapper(); + private static ExecutorService executorService; + private static ScheduledExecutorService metricsScheduler; + private static List mappingRules = new ArrayList<>(); + + // MQTT клиент + private static MqttClient mqttClient; + + // Счётчики для мониторинга + private static final AtomicLong processedCounter = new AtomicLong(0); + private static final AtomicLong errorCounter = new AtomicLong(0); + private static final AtomicLong retryCounter = new AtomicLong(0); + + // ---------- Вспомогательный класс правила ---------- + static class MappingRule { + Pattern topicPattern; + String hostTemplate; + String keyTemplate; + String valueField; + String subscribeTopic; + + MappingRule(String topicPattern, String hostTemplate, String keyTemplate, + String valueField, String subscribeTopic) { + this.topicPattern = Pattern.compile(topicPattern); + this.hostTemplate = hostTemplate; + this.keyTemplate = keyTemplate; + this.valueField = valueField; + this.subscribeTopic = subscribeTopic; + } + + Matcher matcher(String topic) { + return topicPattern.matcher(topic); + } + } + + // ---------- Точка входа ---------- + public static void main(String[] args) { + // Блокировка от двойного запуска + try (FileOutputStream lockFile = new FileOutputStream("app.lock"); + FileLock lock = lockFile.getChannel().tryLock()) { + if (lock == null) { + System.err.println("Another instance is already running. Exiting."); + System.exit(1); + } + } catch (IOException e) { + System.err.println("Failed to acquire lock: " + e.getMessage()); + System.exit(1); + } + + if (!loadConfig("config.properties")) { + System.err.println("Failed to load configuration. Exiting."); + System.exit(1); + } + if (!loadMapping("mapping.json")) { + System.err.println("Failed to load mapping. Exiting."); + System.exit(1); + } + + // Валидация окружения + validateEnvironment(); + + // Создаём пул потоков + BlockingQueue workQueue = new ArrayBlockingQueue<>(queueCapacity); + executorService = new ThreadPoolExecutor( + corePoolSize, + maxPoolSize, + 60L, TimeUnit.SECONDS, + workQueue, + new ThreadPoolExecutor.CallerRunsPolicy() + ); + + // Запускаем планировщик метрик + metricsScheduler = Executors.newSingleThreadScheduledExecutor(); + metricsScheduler.scheduleAtFixedRate(() -> { + try { + long processed = processedCounter.get(); + long errors = errorCounter.get(); + long retries = retryCounter.get(); + int queueSize = ((ThreadPoolExecutor) executorService).getQueue().size(); + + sendMetric("mqtt.bridge.processed", processed); + sendMetric("mqtt.bridge.errors", errors); + sendMetric("mqtt.bridge.retries", retries); + sendMetric("mqtt.bridge.queue_size", queueSize); + + System.out.println("Metrics sent: processed=" + processed + + ", errors=" + errors + + ", retries=" + retries + + ", queue=" + queueSize); + } catch (Exception e) { + System.err.println("Failed to send metrics: " + e.getMessage()); + } + }, metricsInterval, metricsInterval, TimeUnit.SECONDS); + + // Подключаемся к MQTT + connectMQTT(); + } + + // ---------- Валидация окружения ---------- + private static void validateEnvironment() { + // Проверка zabbix_sender + try { + ProcessBuilder pb = new ProcessBuilder("zabbix_sender", "--help"); + pb.redirectErrorStream(true); + Process p = pb.start(); + int exitCode = p.waitFor(); + if (exitCode != 0) { + System.err.println("WARNING: zabbix_sender not found or not executable. Sending will fail."); + } else { + System.out.println("zabbix_sender: OK"); + } + } catch (Exception e) { + System.err.println("WARNING: Failed to check zabbix_sender: " + e.getMessage()); + } + + // Проверка обязательных параметров (дополнительно) + if (zabbixServer == null || zabbixServer.trim().isEmpty()) { + System.err.println("ERROR: zabbix.server is not set."); + System.exit(1); + } + if (zabbixPort <= 0 || zabbixPort > 65535) { + System.err.println("ERROR: Invalid zabbix.port."); + System.exit(1); + } + if (mqttBroker == null || mqttBroker.trim().isEmpty()) { + System.err.println("ERROR: mqtt.broker is not set."); + System.exit(1); + } + + // Проверка маппинга на дублирующиеся topicPattern (опционально) + Set patterns = new HashSet<>(); + for (MappingRule rule : mappingRules) { + String patternStr = rule.topicPattern.pattern(); + if (!patterns.add(patternStr)) { + System.err.println("WARNING: Duplicate topicPattern: " + patternStr); + } + if (rule.hostTemplate == null || rule.hostTemplate.trim().isEmpty()) { + System.err.println("ERROR: host is empty for pattern: " + patternStr); + System.exit(1); + } + if (rule.keyTemplate == null || rule.keyTemplate.trim().isEmpty()) { + System.err.println("ERROR: key is empty for pattern: " + patternStr); + System.exit(1); + } + if (rule.subscribeTopic != null && rule.subscribeTopic.trim().isEmpty()) { + System.err.println("WARNING: subscribeTopic is empty for pattern: " + patternStr + ", will use default."); + } + } + System.out.println("Environment validation completed."); + } + + // ---------- Подключение к MQTT ---------- + private static void connectMQTT() { + try { + mqttClient = new MqttClient(mqttBroker, MqttClient.generateClientId(), new MemoryPersistence()); + MqttConnectOptions options = new MqttConnectOptions(); + options.setCleanSession(true); + if (mqttUser != null && !mqttUser.trim().isEmpty()) { + options.setUserName(mqttUser); + options.setPassword(mqttPassword.toCharArray()); + } + + mqttClient.setCallback(new MqttCallback() { + @Override + public void connectionLost(Throwable cause) { + System.err.println("MQTT connection lost: " + cause.getMessage()); + new Thread(MqttToZabbix::reconnect).start(); + } + + @Override + public void messageArrived(String topic, MqttMessage message) { + String payload = new String(message.getPayload()); + System.out.println("Received from MQTT topic '" + topic + "': " + payload); + + executorService.submit(() -> { + try { + MappingRule rule = findRule(topic); + if (rule != null) { + processWithRule(topic, payload, rule); + } else { + processLegacy(payload); + } + processedCounter.incrementAndGet(); + } catch (Exception e) { + System.err.println("Error processing MQTT message: " + e.getMessage()); + errorCounter.incrementAndGet(); + } + }); + } + + @Override + public void deliveryComplete(IMqttDeliveryToken token) {} + }); + + mqttClient.connect(options); + subscribeTopics(); + System.out.println("Connected to MQTT broker."); + System.out.println("Zabbix target: " + zabbixServer + ":" + zabbixPort); + + // Shutdown hook + Runtime.getRuntime().addShutdownHook(new Thread(() -> { + System.out.println("Shutting down..."); + executorService.shutdown(); + metricsScheduler.shutdown(); + try { + if (!executorService.awaitTermination(30, TimeUnit.SECONDS)) { + executorService.shutdownNow(); + } + if (!metricsScheduler.awaitTermination(10, TimeUnit.SECONDS)) { + metricsScheduler.shutdownNow(); + } + } catch (InterruptedException e) { + executorService.shutdownNow(); + metricsScheduler.shutdownNow(); + } + try { + if (mqttClient != null && mqttClient.isConnected()) { + mqttClient.disconnect(); + mqttClient.close(); + } + } catch (MqttException e) { + e.printStackTrace(); + } + })); + + } catch (MqttException e) { + System.err.println("MQTT initialization error: " + e.getMessage()); + reconnect(); + } + } + + // ---------- Переподключение к MQTT (экспоненциальная задержка) ---------- + private static void reconnect() { + long delay = 1000; + int maxDelay = 60000; + while (!Thread.currentThread().isInterrupted()) { + try { + System.out.println("Attempting to reconnect to MQTT..."); + if (mqttClient != null && !mqttClient.isConnected()) { + MqttConnectOptions options = new MqttConnectOptions(); + options.setCleanSession(true); + if (mqttUser != null && !mqttUser.trim().isEmpty()) { + options.setUserName(mqttUser); + options.setPassword(mqttPassword.toCharArray()); + } + mqttClient.connect(options); + subscribeTopics(); + System.out.println("Reconnected to MQTT successfully."); + return; + } + } catch (MqttException e) { + System.err.println("Reconnect attempt failed: " + e.getMessage()); + } + try { + Thread.sleep(delay); + } catch (InterruptedException ie) { + Thread.currentThread().interrupt(); + return; + } + delay = Math.min(delay * 2, maxDelay); + } + } + + // ---------- Подписка на топики (из маппинга или глобальный) ---------- + private static void subscribeTopics() throws MqttException { + Set topics = new HashSet<>(); + for (MappingRule rule : mappingRules) { + if (rule.subscribeTopic != null && !rule.subscribeTopic.trim().isEmpty()) { + topics.add(rule.subscribeTopic.trim()); + } + } + if (topics.isEmpty()) { + String defaultTopic = (mqttTopic != null && !mqttTopic.trim().isEmpty()) ? mqttTopic : "#"; + topics.add(defaultTopic); + } + for (String topic : topics) { + mqttClient.subscribe(topic); + System.out.println("Subscribed to: " + topic); + } + } + + // ---------- Загрузка конфигурации ---------- + private static boolean loadConfig(String fileName) { + Properties props = new Properties(); + try (InputStream input = getConfigInputStream(fileName)) { + if (input == null) { + System.err.println("Configuration file not found: " + fileName); + return false; + } + props.load(input); + + mqttBroker = props.getProperty("mqtt.broker"); + mqttTopic = props.getProperty("mqtt.topic"); + mqttUser = props.getProperty("mqtt.user", ""); + mqttPassword = props.getProperty("mqtt.password", ""); + zabbixServer = props.getProperty("zabbix.server"); + String portStr = props.getProperty("zabbix.port", "10051"); + try { + zabbixPort = Integer.parseInt(portStr); + } catch (NumberFormatException e) { + System.err.println("Invalid zabbix.port: " + portStr); + return false; + } + + corePoolSize = Integer.parseInt(props.getProperty("thread.pool.core", "2")); + maxPoolSize = Integer.parseInt(props.getProperty("thread.pool.max", "4")); + queueCapacity = Integer.parseInt(props.getProperty("thread.pool.queue", "1000")); + maxRetries = Integer.parseInt(props.getProperty("retry.max.attempts", "3")); + initialDelayMs = Long.parseLong(props.getProperty("retry.initial.delay.ms", "1000")); + backoffMultiplier = Double.parseDouble(props.getProperty("retry.backoff.multiplier", "2.0")); + + metricsHost = props.getProperty("metrics.host", "MQTT-Bridge"); + metricsInterval = Integer.parseInt(props.getProperty("metrics.interval", "60")); + + return true; + } catch (IOException e) { + System.err.println("Error reading configuration file: " + e.getMessage()); + return false; + } + } + + // ---------- Загрузка маппинга ---------- + private static boolean loadMapping(String fileName) { + try (InputStream input = getConfigInputStream(fileName)) { + if (input == null) { + System.err.println("Mapping file not found: " + fileName + " - using legacy mode."); + return true; + } + JsonNode root = objectMapper.readTree(input); + if (root.isArray()) { + for (JsonNode node : root) { + String pattern = node.get("topicPattern").asText(); + String host = node.get("host").asText(); + String key = node.get("key").asText(); + String valueField = node.has("valueField") ? node.get("valueField").asText() : null; + String subscribeTopic = node.has("subscribeTopic") ? node.get("subscribeTopic").asText() : null; + mappingRules.add(new MappingRule(pattern, host, key, valueField, subscribeTopic)); + } + System.out.println("Loaded " + mappingRules.size() + " mapping rules."); + return true; + } else { + System.err.println("Mapping file must contain a JSON array."); + return false; + } + } catch (IOException e) { + System.err.println("Error reading mapping file: " + e.getMessage()); + return false; + } + } + + private static InputStream getConfigInputStream(String fileName) throws IOException { + File file = new File(fileName); + if (file.exists() && file.isFile()) { + return new FileInputStream(file); + } + InputStream is = MqttToZabbix.class.getClassLoader().getResourceAsStream(fileName); + if (is != null) { + return is; + } + return null; + } + + // ---------- Поиск правила по топику ---------- + private static MappingRule findRule(String topic) { + for (MappingRule rule : mappingRules) { + Matcher m = rule.matcher(topic); + if (m.matches()) { + return rule; + } + } + return null; + } + + // ---------- Обработка по правилу ---------- + private static void processWithRule(String topic, String payload, MappingRule rule) { + Matcher matcher = rule.matcher(topic); + if (!matcher.matches()) return; + + String host = interpolate(rule.hostTemplate, matcher); + String key = interpolate(rule.keyTemplate, matcher); + + String valueStr; + if (rule.valueField != null && !rule.valueField.isEmpty() && !"null".equalsIgnoreCase(rule.valueField)) { + try { + JsonNode jsonNode = objectMapper.readTree(payload); + if (jsonNode.isObject()) { + JsonNode fieldNode = jsonNode.path(rule.valueField); + if (!fieldNode.isMissingNode()) { + valueStr = fieldNode.asText(); + } else { + System.err.println("Field '" + rule.valueField + "' not found in JSON payload: " + payload); + return; + } + } else { + System.err.println("Payload is not a JSON object, using entire payload for topic " + topic); + valueStr = payload; + } + } catch (IOException e) { + System.err.println("Failed to parse JSON for topic " + topic + ": " + e.getMessage() + ", using entire payload"); + valueStr = payload; + } + } else { + valueStr = payload; + } + + sendToZabbixWithRetry(host, key, valueStr, System.currentTimeMillis() / 1000); + } + + // ---------- Legacy обработка (старый формат) ---------- + private static void processLegacy(String payload) { + if (payload == null || payload.trim().isEmpty() || !payload.trim().startsWith("{")) { + System.err.println("Legacy processing skipped: payload is not a JSON object: " + payload); + return; + } + try { + Map data = objectMapper.readValue(payload, Map.class); + String host = (String) data.get("host"); + String key = (String) data.get("key"); + Object value = data.get("value"); + Long clock = data.containsKey("clock") + ? ((Number) data.get("clock")).longValue() + : System.currentTimeMillis() / 1000; + sendToZabbixWithRetry(host, key, String.valueOf(value), clock); + } catch (Exception e) { + System.err.println("Legacy processing failed: " + e.getMessage()); + } + } + + // ---------- Интерполяция строк с $1, $2 ... ---------- + private static String interpolate(String template, Matcher matcher) { + if (template == null) return null; + String result = template; + for (int i = 1; i <= matcher.groupCount(); i++) { + String group = matcher.group(i); + if (group != null) { + result = result.replace("$" + i, group); + } + } + return result; + } + + // ---------- Отправка с повторными попытками ---------- + private static void sendToZabbixWithRetry(String host, String key, String value, long clock) { + long delayMs = initialDelayMs; + for (int attempt = 1; attempt <= maxRetries; attempt++) { + try { + boolean success = sendToZabbix(host, key, value, clock); + if (success) { + return; + } + if (attempt < maxRetries) { + System.err.println("Retry " + attempt + "/" + maxRetries + " for " + host + ":" + key + + " - waiting " + delayMs + " ms"); + retryCounter.incrementAndGet(); + Thread.sleep(delayMs); + delayMs = (long) (delayMs * backoffMultiplier); + } else { + System.err.println("All retries failed for " + host + ":" + key); + errorCounter.incrementAndGet(); + } + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + System.err.println("Retry interrupted for " + host + ":" + key); + return; + } catch (Exception e) { + if (attempt < maxRetries) { + System.err.println("Attempt " + attempt + "/" + maxRetries + " failed for " + host + ":" + key + + " - error: " + e.getMessage() + ", retrying in " + delayMs + " ms"); + retryCounter.incrementAndGet(); + try { + Thread.sleep(delayMs); + delayMs = (long) (delayMs * backoffMultiplier); + } catch (InterruptedException ie) { + Thread.currentThread().interrupt(); + return; + } + } else { + System.err.println("All retries failed for " + host + ":" + key + " - final error: " + e.getMessage()); + errorCounter.incrementAndGet(); + } + } + } + } + + // ---------- Отправка в Zabbix (через внешний zabbix_sender) ---------- + private static boolean sendToZabbix(String host, String key, String value, long clock) { + try { + ProcessBuilder pb = new ProcessBuilder( + "zabbix_sender", + "-z", zabbixServer, + "-p", String.valueOf(zabbixPort), + "-s", host, + "-k", key, + "-o", String.valueOf(value) + ); + pb.redirectErrorStream(true); + Process p = pb.start(); + int exitCode = p.waitFor(); + if (exitCode == 0) { + try (BufferedReader reader = new BufferedReader(new InputStreamReader(p.getInputStream()))) { + String line; + while ((line = reader.readLine()) != null) { + if (line.contains("processed:")) { + System.out.println("zabbix_sender: " + line); + } + } + } + return true; + } else { + try (BufferedReader reader = new BufferedReader(new InputStreamReader(p.getInputStream()))) { + String line; + while ((line = reader.readLine()) != null) { + System.err.println("zabbix_sender error: " + line); + } + } + return false; + } + } catch (Exception e) { + System.err.println("zabbix_sender execution error: " + e.getMessage()); + return false; + } + } + + // ---------- Отправка метрик (без учёта счётчиков) ---------- + private static void sendMetric(String key, long value) { + try { + ProcessBuilder pb = new ProcessBuilder( + "zabbix_sender", + "-z", zabbixServer, + "-p", String.valueOf(zabbixPort), + "-s", metricsHost, + "-k", key, + "-o", String.valueOf(value) + ); + pb.redirectErrorStream(true); + Process p = pb.start(); + int exitCode = p.waitFor(); + if (exitCode != 0) { + try (BufferedReader reader = new BufferedReader(new InputStreamReader(p.getInputStream()))) { + String line; + while ((line = reader.readLine()) != null) { + System.err.println("metric send error: " + line); + } + } + } + } catch (Exception e) { + System.err.println("Failed to send metric " + key + ": " + e.getMessage()); + } + } +} \ No newline at end of file diff --git a/src/main/resources/config.properties b/src/main/resources/config.properties new file mode 100644 index 0000000..ab2c046 --- /dev/null +++ b/src/main/resources/config.properties @@ -0,0 +1,23 @@ +# MQTT Broker settings +mqtt.broker=tcp://localhost:1883 +mqtt.topic=# +mqtt.user= +mqtt.password= + +# Zabbix Trapper settings +zabbix.server=127.0.0.1 +zabbix.port=10051 + +# Asynchronous processing +thread.pool.core=2 +thread.pool.max=4 +thread.pool.queue=1000 + +# Retry settings (exponential backoff) +retry.max.attempts=3 +retry.initial.delay.ms=1000 +retry.backoff.multiplier=2.0 + +# Monitoring (metrics) +metrics.host=MQTT-Bridge +metrics.interval=60 \ No newline at end of file diff --git a/src/main/resources/mapping.json b/src/main/resources/mapping.json new file mode 100644 index 0000000..15fb585 --- /dev/null +++ b/src/main/resources/mapping.json @@ -0,0 +1,29 @@ +[ + { + "topicPattern": "sensors/temperature", + "host": "Server1", + "key": "temperature", + "valueField": null, + "subscribeTopic": "sensors/temperature" + }, + { + "topicPattern": "sensors/humidity", + "host": "Server1", + "key": "humidity", + "valueField": null, + "subscribeTopic": "sensors/humidity" + }, + { + "topicPattern": "devices/([^/]+)/status", + "host": "$1", + "key": "status", + "valueField": "state", + "subscribeTopic": "devices/+/status" + }, + { + "topicPattern": "metrics/([^/]+)/([^/]+)", + "host": "$1", + "key": "$2", + "valueField": "value" + } +] \ No newline at end of file diff --git a/src/main/resources/mqtt-to-zabbix.service b/src/main/resources/mqtt-to-zabbix.service new file mode 100644 index 0000000..22dbe36 --- /dev/null +++ b/src/main/resources/mqtt-to-zabbix.service @@ -0,0 +1,19 @@ +[Unit] +Description=MQTT to Zabbix Bridge +After=network.target zabbix-server.service +Wants=zabbix-server.service + +[Service] +Type=simple +User=zabbix # или ваш пользователь +Group=zabbix +WorkingDirectory=/opt/mqtt-to-zabbix +ExecStart=/usr/bin/java -jar /opt/mqtt-to-zabbix/mqtt-to-zabbix-1.0-SNAPSHOT-jar-with-dependencies.jar +Restart=always +RestartSec=10 +StandardOutput=journal +StandardError=journal +SyslogIdentifier=mqtt-to-zabbix + +[Install] +WantedBy=multi-user.target \ No newline at end of file