Blitze per MQTT

Begonnen von Salvi5, 14 Juli 2026, 12:52:00

Vorheriges Thema - Nächstes Thema

Salvi5

Ich habe mir von der freundlichen KI von nebenan mal ein Skript erstellen lassen, welches die Blitze innerhalb eines festzulegenden Radius um festzulegende Koordinaten herum per MQTT meldet. Ich finde das praktisch, weil man dann auf ein heranziehendes Gewitter reagieren kann. Vielleicht findet es ja noch jemand nützlich.

Gruß Mike

#!/usr/bin/env python3
"""
Blitzortung.org -> MQTT Bridge fuer FHEM
==========================================

Verbindet sich mit dem oeffentlichen Blitzortung-WebSocket (dasselbe
Netzwerk, das auch lightningmaps.org anzeigt), filtert Blitzeinschlaege
nach Entfernung zu einem festgelegten Standort und veroeffentlicht
Treffer innerhalb des Radius per MQTT.

Benoetigte Pakete:
    pip install websockets paho-mqtt --break-system-packages

Aufruf:
    python3 blitzortung_bridge.py

Zum Dauerbetrieb empfiehlt sich ein systemd-Service oder ein
Neustart-Wrapper (siehe Hinweise am Ende der Datei).
"""

import asyncio
import json
import logging
import math
import time
from datetime import datetime, timezone

import websockets
import paho.mqtt.client as mqtt

# ---------------------------------------------------------------------------
# Konfiguration - hier anpassen
# ---------------------------------------------------------------------------

HOME_LAT = 50.815      # eigene Breite
HOME_LON = 10.225      # eigene Laenge
RADIUS_KM = 20         # Umkreis, der als "hier" gilt

MQTT_HOST = "127.0.0.1"
MQTT_PORT = 1883
MQTT_USER = None       # z.B. "fhem" oder None wenn kein Login noetig
MQTT_PASS = None
MQTT_TOPIC = "blitz/strike"
MQTT_CLIENT_ID = "blitzortung_bridge"

# Mehrere Blitzortung-Relays, falls eines nicht erreichbar ist wird
# automatisch das naechste probiert
WS_SERVERS = [
    "wss://ws1.blitzortung.org/",
    "wss://ws2.blitzortung.org/",
    "wss://ws3.blitzortung.org/",
    "wss://ws4.blitzortung.org/",
    "wss://ws5.blitzortung.org/",
    "wss://ws6.blitzortung.org/",
    "wss://ws7.blitzortung.org/",
    "wss://ws8.blitzortung.org/",
]

# Handshake-Nachricht, die der Server erwartet bevor er Daten schickt
HANDSHAKE_MSG = '{"a": 111}'

# Blitzortung komprimiert die JSON-Daten mit einem einfachen
# Verfahren ("run length" auf Basis eines char-Offsets). Diese
# Funktion entpackt es wieder in normales JSON.
def unpack_blitzortung(raw: str) -> str:
    """Entpackt die LZW-aehnliche Kompression, die Blitzortung fuer
    die WebSocket-Rohdaten verwendet. Portiert aus einer verifizierten
    Referenzimplementierung (gkbrk.com/blitzortung)."""
    d = list(raw)
    if not d:
        return ""
    dictionary = {}
    c = d[0]
    f = c
    out = [c]
    next_code = 256
    for i in range(1, len(d)):
        code = ord(d[i])
        if code < 256:
            a = d[i]
        else:
            a = dictionary.get(code) or (f + c)
        out.append(a)
        c = a[0]
        dictionary[next_code] = f + c
        next_code += 1
        f = a
    return "".join(out)


def haversine_km(lat1, lon1, lat2, lon2) -> float:
    r = 6371.0
    dlat = math.radians(lat2 - lat1)
    dlon = math.radians(lon2 - lon1)
    a = (
        math.sin(dlat / 2) ** 2
        + math.cos(math.radians(lat1))
        * math.cos(math.radians(lat2))
        * math.sin(dlon / 2) ** 2
    )
    return r * 2 * math.asin(math.sqrt(a))


logging.basicConfig(
    level=logging.INFO,
    format="%(asctime)s %(levelname)s %(message)s",
)
log = logging.getLogger("blitz-bridge")


class MqttPublisher:
    def __init__(self):
        self.client = mqtt.Client(client_id=MQTT_CLIENT_ID)
        if MQTT_USER:
            self.client.username_pw_set(MQTT_USER, MQTT_PASS)
        self.client.connect(MQTT_HOST, MQTT_PORT, keepalive=60)
        self.client.loop_start()

    def publish(self, payload: dict):
        self.client.publish(MQTT_TOPIC, json.dumps(payload), qos=0, retain=False)


async def listen_once(publisher: MqttPublisher, start_index: int = 0):
    """Ein Verbindungsversuch. Wirft eine Exception bei Verbindungsabbruch."""
    servers = WS_SERVERS[start_index:] + WS_SERVERS[:start_index]
    for server in servers:
        try:
            log.debug("Verbinde mit %s ...", server)
            async with websockets.connect(server, ping_interval=20, ping_timeout=20) as ws:
                await ws.send(HANDSHAKE_MSG)
                log.debug("Verbunden, warte auf Blitzdaten ...")
                async for raw_msg in ws:
                    handle_message(raw_msg, publisher)
                return  # Verbindung sauber beendet -> aeussere Schleife reconnected
        except Exception as exc:
            log.debug("Verbindung zu %s fehlgeschlagen: %s", server, exc)
            continue
    # Wenn alle Server fehlgeschlagen sind, kurz warten bevor Retry
    await asyncio.sleep(5)


def handle_message(raw_msg: str, publisher: MqttPublisher):
    try:
        text = unpack_blitzortung(raw_msg)
        data = json.loads(text)
    except Exception as exc:
        log.debug("Konnte Nachricht nicht verarbeiten (%s): %.100s", exc, raw_msg)
        return  # ungueltige/nicht relevante Nachricht ignorieren

    lat = data.get("lat")
    lon = data.get("lon")
    if lat is None or lon is None:
        return

    dist = haversine_km(HOME_LAT, HOME_LON, lat, lon)
    if dist <= RADIUS_KM:
        ts_ns = data.get("time")
        ts = (
            datetime.fromtimestamp(ts_ns / 1e9, tz=timezone.utc).isoformat()
            if ts_ns
            else datetime.now(tz=timezone.utc).isoformat()
        )
        payload = {
            "distance_km": round(dist, 1),
            "lat": lat,
            "lon": lon,
            "time": ts,
        }
        log.info("Blitz im Radius: %.1f km entfernt", dist)
        publisher.publish(payload)


async def main():
    publisher = MqttPublisher()
    server_index = 0
    while True:
        try:
            await listen_once(publisher, start_index=server_index)
        except Exception as exc:
            log.error("Unerwarteter Fehler: %s", exc)
        server_index = (server_index + 1) % len(WS_SERVERS)
        log.debug("Verbindung verloren, versuche erneut in 5s ...")
        await asyncio.sleep(5)


if __name__ == "__main__":
    try:
        asyncio.run(main())
    except KeyboardInterrupt:
        log.info("Beendet durch Benutzer.")


# ---------------------------------------------------------------------------
# Hinweise fuer den Dauerbetrieb (systemd)
# ---------------------------------------------------------------------------
#
# /etc/systemd/system/blitzortung-bridge.service
#
# [Unit]
# Description=Blitzortung MQTT Bridge fuer FHEM
# After=network-online.target
#
# [Service]
# ExecStart=/usr/bin/python3 /opt/blitzortung/blitzortung_bridge.py
# Restart=always
# RestartSec=5
# User=pi
#
# [Install]
# WantedBy=multi-user.target
#
# Dann:
#   sudo systemctl daemon-reload
#   sudo systemctl enable --now blitzortung-bridge

pink99panther

Hallo

ich habe es mit hilfe von KI mal um die Himmelsrichtung erweitert.
#!/usr/bin/env python3
"""
Blitzortung.org -> MQTT Bridge fuer FHEM
==========================================

Verbindet sich mit dem oeffentlichen Blitzortung-WebSocket (dasselbe
Netzwerk, das auch lightningmaps.org anzeigt), filtert Blitzeinschlaege
nach Entfernung zu einem festgelegten Standort und veroeffentlicht
Treffer innerhalb des Radius per MQTT.

Benoetigte Pakete:
    pip install websockets paho-mqtt --break-system-packages

Aufruf:
    python3 blitzortung_bridge.py

Zum Dauerbetrieb empfiehlt sich ein systemd-Service oder ein
Neustart-Wrapper (siehe Hinweise am Ende der Datei).
"""

import asyncio
import json
import logging
import math
import time
from datetime import datetime, timezone

import websockets
import paho.mqtt.client as mqtt

# ---------------------------------------------------------------------------
# Konfiguration - hier anpassen
# ---------------------------------------------------------------------------

HOME_LAT = 50.815      # eigene Breite
HOME_LON = 10.225      # eigene Laenge
RADIUS_KM = 20         # Umkreis, der als "hier" gilt

MQTT_HOST = "127.0.0.1"
MQTT_PORT = 1883
MQTT_USER = None       # z.B. "fhem" oder None wenn kein Login noetig
MQTT_PASS = None
MQTT_TOPIC = "blitz/strike"
MQTT_CLIENT_ID = "blitzortung_bridge"

# Mehrere Blitzortung-Relays, falls eines nicht erreichbar ist wird
# automatisch das naechste probiert
WS_SERVERS = [
    "wss://ws1.blitzortung.org/",
    "wss://ws2.blitzortung.org/",
    "wss://ws3.blitzortung.org/",
    "wss://ws4.blitzortung.org/",
    "wss://ws5.blitzortung.org/",
    "wss://ws6.blitzortung.org/",
    "wss://ws7.blitzortung.org/",
    "wss://ws8.blitzortung.org/",
]

# Handshake-Nachricht, die der Server erwartet bevor er Daten schickt
HANDSHAKE_MSG = '{"a": 111}'

# Blitzortung komprimiert die JSON-Daten mit einem einfachen
# Verfahren ("run length" auf Basis eines char-Offsets). Diese
# Funktion entpackt es wieder in normales JSON.
def unpack_blitzortung(raw: str) -> str:
    """Entpackt die LZW-aehnliche Kompression, die Blitzortung fuer
    die WebSocket-Rohdaten verwendet. Portiert aus einer verifizierten
    Referenzimplementierung (gkbrk.com/blitzortung)."""
    d = list(raw)
    if not d:
        return ""
    dictionary = {}
    c = d[0]
    f = c
    out = [c]
    next_code = 256
    for i in range(1, len(d)):
        code = ord(d[i])
        if code < 256:
            a = d[i]
        else:
            a = dictionary.get(code) or (f + c)
        out.append(a)
        c = a[0]
        dictionary[next_code] = f + c
        next_code += 1
        f = a
    return "".join(out)


def haversine_km(lat1, lon1, lat2, lon2) -> float:
    r = 6371.0
    dlat = math.radians(lat2 - lat1)
    dlon = math.radians(lon2 - lon1)
    a = (
        math.sin(dlat / 2) ** 2
        + math.cos(math.radians(lat1))
        * math.cos(math.radians(lat2))
        * math.sin(dlon / 2) ** 2
    )
    return r * 2 * math.asin(math.sqrt(a))


def calculate_bearing(lat1, lon1, lat2, lon2) -> float:
    """Berechnet den mathematischen Grosskreiswinkel (Azimut) von Punkt 1 zu Punkt 2.
    Ergebnis: 0° = Norden, 90° = Osten, 180° = Sueden, 270° = Westen"""
    lat1_rad = math.radians(lat1)
    lat2_rad = math.radians(lat2)
    diff_lon_rad = math.radians(lon2 - lon1)

    x = math.sin(diff_lon_rad) * math.cos(lat2_rad)
    y = math.cos(lat1_rad) * math.sin(lat2_rad) - (
        math.sin(lat1_rad) * math.cos(lat2_rad) * math.cos(diff_lon_rad)
    )

    initial_bearing = math.atan2(x, y)
    initial_bearing = math.degrees(initial_bearing)
    return (initial_bearing + 360) % 360


def bearing_to_text(bearing: float) -> str:
    """Konvertiert Gradzahlen (0-360) in Himmelsrichtungs-Kuerzel (8er-Teilung)."""
    directions = ["N", "NE", "E", "SE", "S", "SW", "W", "NW"]
    index = int((bearing + 22.5) / 45) % 8
    return directions[index]


logging.basicConfig(
    level=logging.INFO,
    format="%(asctime)s %(levelname)s %(message)s",
)
log = logging.getLogger("blitz-bridge")


class MqttPublisher:
    def __init__(self):
        self.client = mqtt.Client(client_id=MQTT_CLIENT_ID)
        if MQTT_USER:
            self.client.username_pw_set(MQTT_USER, MQTT_PASS)
        self.client.connect(MQTT_HOST, MQTT_PORT, keepalive=60)
        self.client.loop_start()

    def publish(self, payload: dict):
        self.client.publish(MQTT_TOPIC, json.dumps(payload), qos=0, retain=False)


async def listen_once(publisher: MqttPublisher, start_index: int = 0):
    """Ein Verbindungsversuch. Wirft eine Exception bei Verbindungsabbruch."""
    servers = WS_SERVERS[start_index:] + WS_SERVERS[:start_index]
    for server in servers:
        try:
            log.debug("Verbinde mit %s ...", server)
            async with websockets.connect(server, ping_interval=20, ping_timeout=20) as ws:
                await ws.send(HANDSHAKE_MSG)
                log.debug("Verbunden, warte auf Blitzdaten ...")
                async for raw_msg in ws:
                    handle_message(raw_msg, publisher)
                return  # Verbindung sauber beendet -> aeussere Schleife reconnected
        except Exception as exc:
            log.debug("Verbindung zu %s fehlgeschlagen: %s", server, exc)
            continue
    # Wenn alle Server fehlgeschlagen sind, kurz warten bevor Retry
    await asyncio.sleep(5)


def handle_message(raw_msg: str, publisher: MqttPublisher):
    try:
        text = unpack_blitzortung(raw_msg)
        data = json.loads(text)
    except Exception as exc:
        log.debug("Konnte Nachricht nicht verarbeiten (%s): %.100s", exc, raw_msg)
        return  # ungueltige/nicht relevante Nachricht ignorieren

    lat = data.get("lat")
    lon = data.get("lon")
    if lat is None or lon is None:
        return

    dist = haversine_km(HOME_LAT, HOME_LON, lat, lon)
    if dist <= RADIUS_KM:
        ts_ns = data.get("time")
        ts = (
            datetime.fromtimestamp(ts_ns / 1e9, tz=timezone.utc).isoformat()
            if ts_ns
            else datetime.now(tz=timezone.utc).isoformat()
        )
       
        # Berechnung der Richtung (Grad und Text)
        bearing = calculate_bearing(HOME_LAT, HOME_LON, lat, lon)
        bearing_txt = bearing_to_text(bearing)
       
        payload = {
            "distance_km": round(dist, 1),
            "direction_deg": round(bearing, 1),
            "direction_text": bearing_txt,
            "lat": lat,
            "lon": lon,
            "time": ts,
        }
        log.info("Blitz im Radius: %.1f km entfernt in Richtung %s (%.1f°)", dist, bearing_txt, bearing)
        publisher.publish(payload)


async def main():
    publisher = MqttPublisher()
    server_index = 0
    while True:
        try:
            await listen_once(publisher, start_index=server_index)
        except Exception as exc:
            log.error("Unerwarteter Fehler: %s", exc)
        server_index = (server_index + 1) % len(WS_SERVERS)
        log.debug("Verbindung verloren, versuche erneut in 5s ...")
        await asyncio.sleep(5)


if __name__ == "__main__":
    try:
        asyncio.run(main())
    except KeyboardInterrupt:
        log.info("Beendet durch Benutzer.")
FHEM alexa-fhem mariadb im Portainer Stack auf Synology NAS