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