Esimerkki · Energia · Datavirrat ja integraatio

Ennusteet ajan tasalla, vaikka tietolähteet päivittyvät eri tahtiin

Energiayhtiön kulutus- ja hintaennusteet · Tapahtumapohjainen dataintegraatio

Energiayhtiön ennusteet perustuvat kulutusmittauksiin, sääennusteisiin ja pörssisähkön hintoihin. Osa tiedoista tulee ulkoisista, osa sisäisistä lähteistä, ja jokainen lähde päivittyy omaan tahtiinsa. Haasteena on pitää syötteet yhtenäisinä ja tuoreina, niin että kulutus- ja hintaennusteet lasketaan uudelleen sekunneissa syötteen muuttumisesta eikä vasta seuraavassa eräajossa.

Datavirrat
LIVE
Syötteet
Kulutus · sää · spot
Ennusteiden laskenta
Tapahtumapohjainen
Tavoite
Ennuste päivittyy, kun syöte muuttuu
Tietolähteet ja päivitysrytmit

Lähteet päivittyvät eri tahtiin

Yksittäisen rajapinnan haku ei ole tässä vaikea osa. Vaikeaa on, että lähteet päivittyvät eri tahtiin ja niiden aikaleima tarkoittaa eri asioita, joko tapahtumahetkeä (event-time) tai käsittelyhetkeä (processing-time). Ennuste on vain yhtä tuore kuin sen vanhin syöte.

Kulutusmittaukset

LähdeDatahub / AMR head-end
ProtokollaTapahtumavirta (Kafka / Event Hubs)
KadenssiVarttitaso, jatkuva
ErityistäTaseselvityksen korjaukset tulevat jälkikäteen, joten osa datasta myöhästyy

Sääennuste

LähdeSää-API (ECMWF / FMI)
ProtokollaREST-poll, ajastettu nouto
KadenssiUusi malliajo ~6 h välein
ErityistäHilamuotoista dataa, joka interpoloidaan paikkakohtaisesti

Pörssisähkön hinta

LähdeNord Pool day-ahead + intraday
ProtokollaREST / markkinasanoma
KadenssiDay-ahead 1×/vrk, intraday jatkuva
ErityistäJulkaisuhetki seuraa markkinan aikataulua (CET)

Lasketut ennusteet

Jokainen ennuste on oma datatuotteensa. Sillä on omat syötteensä ja omat ehtonsa sille, milloin se lasketaan uudelleen.

Inference
KulutusennusteSyötteet: kulutushistoria + sääennuste + kalenteri. Koneoppimis- tai aikasarjamalli, jonka laskennan käynnistää tapahtuma.
HintaennusteSyötteet: kulutusennuste + sääennuste (tuuli/aurinko) + spot-historia. Lasketaan uudelleen, kun syöte päivittyy.

Eräajo ei pysy mukana

Kun syötteet saapuvat eri tahtiin, kiinteä eräajo joko viivästyttää tuoretta dataa tai laskee ennusteen turhaan uudelleen, vaikka syöte ei ole muuttunut.

Tuoreus
  • Uuden sääennusteen pitäisi näkyä kulutus- ja hintaennusteessa sekunneissa eikä vasta tuntien päästä.
  • Kun mittausdata tulee myöhässä, ennuste lasketaan uudelleen sille ajanjaksolle, jolta mittaus on (event-time).
  • Syötteiden riippuvuuksia pitää hallita, jotta päivitys laskee uudelleen vain ne ennusteet, joihin se vaikuttaa.
Integraatioarkkitehtuuri

Putki lähteestä ennusteen jakeluun

Tapahtumaväylä erottaa tietolähteet ja datan käyttäjät toisistaan. Päivitykset kulkevat virtana, ja ennusteet lasketaan uudelleen, kun syöte muuttuu.

1VastaanottoAjastetut haut, virrat, tiedostot
2TapahtumaväyläKafka / Event Hubs
3VirtaprosessointiNormalisointi, event-time
4TallennusAikasarja + feature store
5LaskentaKulutus- ja hintaennusteet
6JakeluAPI + alavirran järjestelmät

Interaktiivinen demo

Näin ennusteet päivittyvät syötteiden mukana

Vasemmalla näkyvät syötteet ja se, kuinka tuoreita ne ovat. Kun syöte päivittyy, siitä riippuvat ennusteet lasketaan oikealla uudelleen. Versionumero kasvaa ja kortti välähtää. Voit myös syöttää päivityksen käsin tai simuloida datakatkon.

Simulaatiokello: 0 s
Syötteet (tapahtumavirrat)
⚡
Lasketut ennusteet
Simulaatio on käynnissä. Syötteet päivittyvät kukin omaan tahtiinsa. Ennuste lasketaan uudelleen vain, kun jokin sen syötteistä muuttuu, ja juuri sitä tapahtumapohjainen integraatio tarkoittaa.
Tekniset huomiot

Mitä tapahtumapohjaisuus vaatii käytännössä

Tapahtumapohjainen uudelleenlaskenta

Laskennan käynnistää uusi syöte eikä kellonaika. Riippuvuusgraafi määrää, mitkä ennusteet lasketaan uudelleen ja missä järjestyksessä.

Tuoreuden valvonta

Jokaiselle syötteelle on sovittu enimmäisikä (tuoreus-SLA). Kun raja ylittyy, syntyy hälytys. Ennuste merkitään vanhentuneen syötteen varassa olevaksi kaikissa sitä käyttävissä järjestelmissä.

Myöhässä ja väärässä järjestyksessä saapuva data

Jälkikäteen saapuva mittaus kohdistetaan sen tapahtumahetkeen (event-time), ja watermarkit rajaavat, kuinka kauan myöhästyvää dataa odotetaan. Korjausajo kohdistuu siihen ajanjaksoon, jolta mittaus on, eikä saapumishetkeen.

Turvallinen toisto ja versiointi

Saman tapahtuman voi käsitellä uudelleen ilman, että tulos muuttuu (idempotenssi). Jokaisella ennusteella on aikaleima ja versio, joten sitä käyttävä järjestelmä tietää aina, mihin syötteisiin tulos perustuu.

Löyhästi kytketty arkkitehtuuri

Tapahtumaväylän ansiosta järjestelmien välillä ei ole suoria point-to-point-kytkentöjä. Uuden lähteen tai datan käyttäjän voi liittää muuttamatta olemassa olevia integraatioita.

Takaisinlataus ja toistettavuus

Historian voi toistaa virrasta (replay). Näin mallit voidaan ajaa uudelleen ja korjaukset tehdä ilman erillistä eräpoimintaputkea.

Lopputulos

Kun kaikki lähteet kulkevat saman väylän kautta ja ennusteet lasketaan tapahtumien pohjalta, järjestelmän tila pysyy yhtenäisenä. Jokainen ennuste perustuu tuoreimpaan saatavilla olevaan syötteeseen, ja sen lähtötiedot voi jäljittää version ja aikaleiman avulla. Lähes reaaliaikainen integraatio ei ole nopeutettu eräajo. Siinä muuttunut data työnnetään eteenpäin sitä mukaa kuin se syntyy, eikä sitä haeta aikataulun mukaan.