MyUplink API -> InfluxDB -> Grafana

  • Viestiketjun aloittaja Viestiketjun aloittaja zyrppa
  • Aloituspäivämäärä Aloituspäivämäärä

zyrppa

Aktiivinen jäsen
MyUplink -> InfluxDB 1.x -> Grafana
Hyödynnetään MyUplinkin rajapintaa ja piirretään Grafanalla haluttuja käppyröitä.
Tämä ohjenuora on tehty sillä olettamalla että käytössä on linux ja influxista vanhempi 1.x versio.
(Skriptien teossa on käytetty osittain tekoälyä apuna)

MyUplink-rajapinnan käyttöönotto:
1. Mene osoitteeseen https://dev.myuplink.com
2. Kirjaudu myUplink-tililläsi
3. Luo uusi Application
4. Kopioi Client ID ja Client Secret itsellesi ylös. Nämä on henk. koht. eikä niitä kannata levitellä muille!
5. Aseta Redirect URI: https://www.marshflattsfarm.org.uk/nibeuplink/oauth2callback/index.php

Asenna tarvittavat python-moduulit mikäli niitä ei ole jo asennettu:
sudo apt update sudo apt install -y python3-pip pip3 install requests requests-oauthlib

Tallenna alla olevat skriptit.
Muuta skripteihin client id, client secret, device id* sekä influx-tietokannan tiedot.

Aja komento: python myuplink_oath.py
Tämä antaa nettiosoitteen jossa ohjelmalle annetaan käyttölupa.
Käyttöluvan antamisen jälkeen sinut uudelleen ohjataan toiselle sivulle, kopioi sen osoite terminaaliin.
Skripti tallentaa automaattisesti tokenin jolloin myuplink_to_influx.py -skripti toimii oikein.

* HAE DEVICE_ID:
Mene osoitteeseen https://api.myuplink.com/swagger/index.html
Anna client id sekä client secret. Ruksi READSYSTEM ja paina Authorize.
Rullaa sivua alemmas kohtaan /v2/devices/me, paina "Try it out" ja "Execute".
Kopioi kohdasta "id" merkkijono. Tämä on sinun DEVICE_ID mikä täytyy muokata skriptiin.


Halutessasi voit ajastaa skriptin hakemaan tiedot esimerkiksi 1 tai 5 minuutin välein:
crontab -e
*/5 * * * * /usr/bin/python3 /home/pi/myuplink_to_influx.py >> /home/pi/myuplink_to_influx.log 2>&1
(Vaihda oikea polku jos eri kuin esimerkissä.)

Tämän jälkeen pystyt piirtämään Grafanalla käppyröitä.

EDIT: Päivitetty koodia hakemaan API:lta tullut aikaleima ja Grafana piirtämään datapisteet sen mukaan eikä skriptin ajohetken mukaan.


Python:
# Skripti tallentaa tietokantaan API-vastauksen aikaleiman jonka mukaan
# grafana piirtää datapisteet.
# Mukana myös fallback mikäli aikaleimaa ei ole
# Grafanassa muista valita GROUP BY: time($_interval) ja fill(previous)
# mikäli et halua graafiin katkoksia.

import os
import json
import time
import requests
import re
from requests_oauthlib import OAuth2Session
from datetime import datetime, timezone

# ==================================================
# myUplink OAuth
# ==================================================
CLIENT_ID = "VAIHDA"
CLIENT_SECRET = "VAIHDA"

TOKEN_URL = "https://api.myuplink.com/oauth/token"
API_BASE = "https://api.myuplink.com/v3"

TOKEN_FILE = os.path.expanduser("~/.myuplink_token.json")

# ==================================================
# InfluxDB 1.x
# ==================================================
INFLUX_HOST = "http://localhost:8086"
INFLUX_DB = "VILP"
INFLUX_USER = "admin"
INFLUX_PASS = "admin"

# ==================================================
# myUplink Device ID
# ==================================================
DEVICE_ID = "VAIHDA"

# ==================================================
# Halutut sensorit
# ==================================================
TARGET_SENSORS = {
    "40004": "Outdoor temperature",
    "40008": "Supply line (BT2)",
    "40012": "Return line (BT3)",
    "40013": "Hot water top (BT7)",
    "40079": "Current (BE3)",
    "40081": "Current (BE2)",
    "40083": "Current (BE1)",
    "40941": "Degree minutes",
    "43009": "Calculated supply climate system 1",
    "43084": "Power internal add. heat",
    "44069": "number of starts",
    "44071": "total operating time",
    "44298": "Hot water, including int. add. heat",
    "44300": "Heating, including int. add. heat",
    "44306": "Hot water, compressor only",
    "44308": "Heating, compressor only",
    "44701": "Current compressor frequency (EB101)",
    "44703": "Defrosting (EB101-EP14)",
    "47028": "change in curve",
    "50096": "status",
    "47032": "climate system",
    "44396": "Heating medium pump speed (GP1)",
    "40940": "current value"
}

# ==================================================
# Token handling
# ==================================================
def load_token():
    with open(TOKEN_FILE) as f:
        return json.load(f)

def save_token(token):
    with open(TOKEN_FILE, "w") as f:
        json.dump(token, f, indent=2)

def get_oauth_session():
    token = load_token()
    oauth = OAuth2Session(CLIENT_ID, token=token)

    # refresh jos vanhentunut
    if token.get("expires_at", 0) <= time.time():
        print("🔄 Access token vanhentunut, refresh...")
        token = oauth.refresh_token(
            TOKEN_URL,
            refresh_token=token["refresh_token"],
            client_id=CLIENT_ID,
            client_secret=CLIENT_SECRET
        )
        save_token(token)
        oauth.token = token

    return oauth

# ==================================================
# myUplink API
# ==================================================
def get_sensor_points(oauth):
    url = f"{API_BASE}/devices/{DEVICE_ID}/points"
    r = oauth.get(url, headers={"accept": "application/json"})
    r.raise_for_status()
    return r.json()


# ==================================================
# InfluxDB write
# ==================================================
def escape_tag(s):
    return str(s).replace(" ", "\\ ").replace(",", "\\,").replace("=", "\\=")

def write_to_influx(lines):
    if not lines:
        print("⚠️ Ei uusia mittauspisteitä, ei tallennusta")
        return
    url = f"{INFLUX_HOST}/write"
    params = {
        "db": INFLUX_DB,
        "u": INFLUX_USER,
        "p": INFLUX_PASS
    }
    r = requests.post(url, params=params, data="\n".join(lines))
    r.raise_for_status()
    print(f"✅ Tallennettu {len(lines)} mittauspistettä InfluxDB:hen")

def parse_myuplink_timestamp(ts):
    """
    Muuntaa myUplink ISO 8601 aikaleiman InfluxDB:n nanosekunteihin
    """
    dt = datetime.fromisoformat(ts.replace("Z", "+00:00"))
    return int(dt.timestamp() * 1e9)


# ==================================================
# Main
# ==================================================
def clean_tag(s):
    """Poistaa kaikki merkit, joita ei saa InfluxDB tagissa olla, kuten SOFT HYPHEN"""
    if not s:
        return "unknown"
    # Poistaa U+00AD ja muut control-merkit
    s = re.sub(r"[\x00-\x1f\x7f\u00ad]", "", s)
    # Escapaa line protocolin tag-merkit
    s = s.replace(" ", "\\ ").replace(",", "\\,").replace("=", "\\=")
    return s

def main():
    oauth = get_oauth_session()
    points = get_sensor_points(oauth)

    influx_lines = []

    for p in points:
        param_id = str(p.get("parameterId"))
        if param_id not in TARGET_SENSORS:
            continue

        value = p.get("value")
        if value is None:
            continue

        # TIMESTAMP + FALLBACK
        api_ts = p.get("timestamp")

        if api_ts:
            try:
                timestamp = parse_myuplink_timestamp(api_ts)
            except Exception as e:
                timestamp = int(time.time() * 1e9)
                print(f"⚠️ Virhe aikaleiman parsimisessa ({api_ts}): {e}")
        else:
            timestamp = int(time.time() * 1e9)
            print(f"⚠️ Puuttuva timestamp sensorille {p.get('parameterName')}")

        # Käytetään parameterName tagina
        parameter_name = clean_tag(p.get("parameterName", TARGET_SENSORS[param_id]))
        device_tag = clean_tag(DEVICE_ID)

        measurement = "myuplink"
        tags = f"device={device_tag},parameter={parameter_name}"

        # numeric vs string
        if isinstance(value, (int, float)):
            fields = f"value={value}"
        else:
            value = str(value).replace("\\", "\\\\").replace('"', '\\"')
            fields = f'value="{value}"'

        line = f"{measurement},{tags} {fields} {timestamp}"
        influx_lines.append(line)

    write_to_influx(influx_lines)

# ==================================================
if __name__ == "__main__":
    main()

Python:
import os
import json
import secrets
import hashlib
import base64
from requests_oauthlib import OAuth2Session

CLIENT_ID = "TÄHÄN CLIENT ID"
CLIENT_SECRET = "TÄHÄN CLIENT SECRET"

REDIRECT_URI = "https://www.marshflattsfarm.org.uk/nibeuplink/oauth2callback/index.php"

AUTHORIZE_URL = "https://api.myuplink.com/oauth/authorize"
TOKEN_URL = "https://api.myuplink.com/oauth/token"

SCOPE = ["READSYSTEM", "offline_access"]

TOKEN_FILE = os.path.expanduser("~/.myuplink_token.json")


def generate_pkce():
    verifier = secrets.token_urlsafe(64)
    challenge = base64.urlsafe_b64encode(
        hashlib.sha256(verifier.encode()).digest()
    ).rstrip(b"=").decode()
    return verifier, challenge


def main():
    code_verifier, code_challenge = generate_pkce()

    oauth = OAuth2Session(
        client_id=CLIENT_ID,
    scope=SCOPE,
        redirect_uri=REDIRECT_URI
    )

    auth_url, state = oauth.authorization_url(
        AUTHORIZE_URL,
        code_challenge=code_challenge,
        code_challenge_method="S256"
    )

    print("\nAvaa selaimessa tämä URL:\n")
    print(auth_url)

    redirect_response = input("\nLiitä koko redirect-URL tähän:\n")

    token = oauth.fetch_token(
        TOKEN_URL,
        authorization_response=redirect_response,
        client_secret=CLIENT_SECRET,
        code_verifier=code_verifier,
        include_client_id=True
    )

    with open(TOKEN_FILE, "w") as f:
        json.dump(token, f, indent=2)

    print("\n✅ Token tallennettu:", TOKEN_FILE)


if __name__ == "__main__":
    main()
 
Viimeksi muokattu:
Ei näköjään voi aloitusta enää muokata. Uusi versio tehty jotta ei tule graafiin katkoksia vaikka data ei olisi päivittynyt pitkiin aikoihin.

Python:
import os
import json
import time
import requests
import re
from requests_oauthlib import OAuth2Session
from datetime import datetime, timezone

# Hakee API:n kautta tiedot, tarvittaessa uusii tokenin
# Vaihda omat tietosi kaikkiin kohtiin missä lukee VAIHDA
# Testattu toimivaksi InfluxDB 1.x kanssa
# Tallentaa jokaisen sensorin tiedot (TARGET_SENSORS) joka kerta kun skripti ajetaan.
# Aikaleima datapisteille on skriptin ajon ajankohta, ei API:n palauttama timestamp
# API:n timestampia käyttämällä Grafanan graafiin tulee katkoksia.
# API:n timestampin saat käyttöön kun poistat kommentoinnin VANHA_VERSIO -kohdasta ja kommentoit UUSI_VERSIO -koodin

CLIENT_ID = "VAIHDA"
CLIENT_SECRET = "VAIHDA"

TOKEN_URL = "https://api.myuplink.com/oauth/token"
API_BASE = "https://api.myuplink.com/v3"

TOKEN_FILE = os.path.expanduser("~/.myuplink_token.json")

INFLUX_HOST = "http://localhost:8086"
INFLUX_DB = "VAIHDA"
INFLUX_USER = "VAIHDA"
INFLUX_PASS = "VAIHDA"
INFLUX_MEASUREMENT = "VAIHDA"

DEVICE_ID = "VAIHDA"

TARGET_SENSORS = {
    "40004": "Outdoor temperature",
    "40008": "Supply line (BT2)",
    "40012": "Return line (BT3)",
    "40013": "Hot water top (BT7)",
    "40079": "Current (BE3)",
    "40081": "Current (BE2)",
    "40083": "Current (BE1)",
    "40941": "Degree minutes",
    "43009": "Calculated supply climate system 1",
    "43084": "Power internal add. heat",
    "44069": "number of starts",
    "44071": "total operating time",
    "44298": "Hot water, including int. add. heat",
    "44300": "Heating, including int. add. heat",
    "44306": "Hot water, compressor only",
    "44308": "Heating, compressor only",
    "44701": "Current compressor frequency (EB101)",
    "44703": "Defrosting (EB101-EP14)",
    "47028": "change in curve",
    "50096": "status",
    "47032": "climate system",
    "44396": "Heating medium pump speed (GP1)",
    "40940": "current value",
    "44362": "Outd temperature (EB101-BT28)"
}

def load_token():
    with open(TOKEN_FILE) as f:
        return json.load(f)

def save_token(token):
    with open(TOKEN_FILE, "w") as f:
        json.dump(token, f, indent=2)

def get_oauth_session():
    token = load_token()
    oauth = OAuth2Session(CLIENT_ID, token=token)

    if token.get("expires_at", 0) <= time.time():
        print("🔄 Access token vanhentunut, refresh...")
        token = oauth.refresh_token(
            TOKEN_URL,
            refresh_token=token["refresh_token"],
            client_id=CLIENT_ID,
            client_secret=CLIENT_SECRET
        )
        save_token(token)
        oauth.token = token

    return oauth

def get_sensor_points(oauth):
    url = f"{API_BASE}/devices/{DEVICE_ID}/points"
    r = oauth.get(url, headers={"accept": "application/json"})
    r.raise_for_status()
    return r.json()

def escape_tag(s):
    return str(s).replace(" ", "\\ ").replace(",", "\\,").replace("=", "\\=")

def write_to_influx(lines):
    if not lines:
        print("⚠️ Ei uusia mittauspisteitä, ei tallennusta")
        return
    url = f"{INFLUX_HOST}/write"
    params = {
        "db": INFLUX_DB,
        "u": INFLUX_USER,
        "p": INFLUX_PASS
    }
    r = requests.post(url, params=params, data="\n".join(lines))
    r.raise_for_status()
    print(f"✅ Tallennettu {len(lines)} mittauspistettä InfluxDB:hen")

#def parse_myuplink_timestamp(ts):
#    dt = datetime.fromisoformat(ts.replace("Z", "+00:00"))
#    return int(dt.timestamp() * 1e9)

def clean_tag(s):
    if not s:
        return "unknown"
    s = re.sub(r"[\x00-\x1f\x7f\u00ad]", "", s)
    s = s.replace(" ", "\\ ").replace(",", "\\,").replace("=", "\\=")
    return s

def main():
    received_parameters = set()
    oauth = get_oauth_session()
    points = get_sensor_points(oauth)

    influx_lines = []
    
    now_ns = int(time.time() * 1e9)

#UUSI_VERSIO ALKAA
    for p in points:
        param_id = p.get("parameterId")
        if param_id not in TARGET_SENSORS:
            continue

        value = p.get("value")
        if value is None:
            continue

        parameter_name = clean_tag(TARGET_SENSORS[param_id])

        line = (
            f"{INFLUX_MEASUREMENT},"
            f"device={DEVICE_ID},"
            f"parameter={parameter_name} "
            f"value={value} {now_ns}"
        )

        influx_lines.append(line)
# LOPPUU

# VANHA_VERSIO
#    for p in points:
#        param_id = p.get("parameterId")
#        if param_id not in TARGET_SENSORS:
#            continue
#
#        parameter_name = clean_tag(TARGET_SENSORS[param_id])
#        api_ts = p.get("timestamp")
#        value = p.get("value")
#
#        if value is not None and api_ts:
#            timestamp = parse_myuplink_timestamp(api_ts)
#
#            line = (
#                f"{INFLUX_MEASUREMENT},"
#                f"device={DEVICE_ID},"
#                f"parameter={parameter_name},"
#                f"source=api "
#                f"value={value} {timestamp}"
#            )
#
#            influx_lines.append(line)
# LOPPUU           
        received_parameters.add(param_id)
        continue

        last_value = get_last_value_from_influx(parameter_name)
        if last_value is None:
            continue

        now_ns = int(time.time() * 1e9)

        line = (
            f"{INFLUX_MEASUREMENT},"
            f"device={DEVICE_ID},"
            f"parameter={parameter_name},"
            f"source=fallback "
            f"value={last_value} {now_ns}"
        )

        influx_lines.append(line)
        received_parameters.add(param_id)

    write_to_influx(influx_lines)

if __name__ == "__main__":
    main()
 
Takaisin
Ylös Bottom