Kirjoittaja: Vectura Solutions Oy · Data-analytiikan ja raportoinnin toteutuksia vuodesta 2018.
1. Aloita yhdestä lähteestä ja yhdestä taulusta
Rajapinta eli API tarjoaa sovellukselle sovitun tavan hakea tai lähettää tietoa. Tässä oppaassa Python tekee HTTPS-pyynnön ulkoiseen rajapintaan, käsittelee JSON-vastauksen ja lataa tietueet BigQuery-tauluun. BigQuery on tietovarasto, jossa aineistoa voidaan myöhemmin yhdistää ja muokata raportointia varten.
Aloitetaan pienestä kokonaisesta aineistosta, esimerkiksi tuoteluettelosta. Kaikki tietueet haetaan jokaisella ajolla, minkä jälkeen erillinen tilannekuvataulu korvataan uudella aineistolla. Tämä on ymmärrettävä lähtökohta silloin, kun koko aineiston hakeminen mahtuu lähteen käyttörajoihin, käytettävissä olevaan muistiin ja ajalle asetettuun tavoitteeseen.
Ennen koodaamista tarkistetaan rajapinnan tunnistautuminen, sivutus, käyttörajat ja muutostietojen saatavuus. Selvitetään myös, tarjoaako lähde jo tarpeeseen sopivan ylläpidetyn liittimen tai eräviennin. Oma Python-ohjelma on perusteltu, kun tiedon haku tai käsittely vaatii omaa logiikkaa.
2. Python-esimerkki: hae sivut ja lataa tilannekuva
Esimerkki käyttää kuvitteellista rajapintasopimusta. Korvatkaa API_URL oman lähteen osoitteella ja sovittakaa kentät sekä sivutus sen dokumentaatioon. API_TOKEN on Bearer-tunniste. Jokainen vastaus sisältää items-listan ja next_cursor-kentän; viimeisellä sivulla next_cursor on null. Jokaisella tietueella on yksilöllinen merkkijonomuotoinen id, aikavyöhykkeellinen updated_at ja name-kenttä.
Kaikki sivut haetaan ennen BigQuery-latausta. Puuttuva sivutustieto, toistuva tunniste tai toistuva kursori pysäyttää ajon. Sivumäärän yläraja estää loputtoman haun. Aikakatkaisut ja rajatut GET-pyyntöjen uusinnat auttavat tilapäisissä verkkovirheissä ja lähteen ruuhkatilanteissa.
Luo esimerkille oma BigQuery-dataset ja taulu, esimerkiksi oma-projekti.raw_api.products_snapshot. WRITE_TRUNCATE korvaa tämän kohteen sisällön ja käyttää latauksen skeemaa: älä osoita esimerkkiä olemassa olevaan liiketoimintatauluun. BigQueryn lataustyö julkaisee tuloksen atomisesti. Koodi odottaa työn valmistumista eikä ilmoita onnistumisesta heti lähetyksen jälkeen.
{
"items": [
{"id": "P-100", "updated_at": "2026-09-10T08:00:00Z", "name": "Tuote A"}
],
"next_cursor": null
}import json
import os
from datetime import datetime, timezone
import requests
from google.cloud import bigquery
from requests.adapters import HTTPAdapter
from urllib3.util.retry import Retry
def fetch_rows(api_url, token):
# Esimerkin API palauttaa kaikki tietueet sivutettuna.
rows, ids, cursors = [], set(), set()
cursor = None
fetched_at = datetime.now(timezone.utc).isoformat()
retry = Retry(
total=3, backoff_factor=1,
status_forcelist=[429, 500, 502, 503, 504],
allowed_methods=["GET"], respect_retry_after_header=True,
)
with requests.Session() as session:
session.mount("https://", HTTPAdapter(max_retries=retry))
session.headers["Authorization"] = f"Bearer {token}"
for _ in range(1000):
response = session.get(
api_url,
params={} if cursor is None else {"cursor": cursor},
timeout=(10, 60), allow_redirects=False,
)
if response.status_code != 200:
raise RuntimeError(f"API-vastaus: {response.status_code}")
body = response.json()
if not isinstance(body, dict) or not isinstance(body.get("items"), list):
raise ValueError("API-vastauksen items-lista puuttuu")
for item in body["items"]:
if not isinstance(item, dict):
raise ValueError("Virheellinen tietue")
record_id = item.get("id")
if not isinstance(record_id, str) or not record_id or record_id in ids:
raise ValueError("Puuttuva tai toistuva id")
updated_at = datetime.fromisoformat(item["updated_at"].replace("Z", "+00:00"))
if updated_at.tzinfo is None:
raise ValueError("updated_at tarvitsee aikavyöhykkeen")
ids.add(record_id)
rows.append({
"id": record_id,
"updated_at": updated_at.isoformat(),
"name": item["name"],
"fetched_at": fetched_at,
})
if "next_cursor" not in body:
raise ValueError("Sivutuksen päättymistieto puuttuu")
cursor = body["next_cursor"]
if cursor is None:
return rows
if not isinstance(cursor, str) or not cursor or cursor in cursors:
raise ValueError("Virheellinen tai toistuva kursori")
cursors.add(cursor)
raise RuntimeError("Sivuraja ylittyi; latausta ei aloitettu")
def main():
api_url = os.environ["API_URL"]
if not api_url.startswith("https://"):
raise ValueError("API_URL edellyttää HTTPS-yhteyttä")
rows = fetch_rows(api_url, os.environ["API_TOKEN"])
# Esimerkin lähteessä tyhjä koko aineisto on poikkeustilanne.
if not rows:
raise ValueError("Tyhjä aineisto: aiempi taulu säilytetään")
config = bigquery.LoadJobConfig(
schema=[
bigquery.SchemaField("id", "STRING", mode="REQUIRED"),
bigquery.SchemaField("updated_at", "TIMESTAMP", mode="REQUIRED"),
bigquery.SchemaField("name", "STRING"),
bigquery.SchemaField("fetched_at", "TIMESTAMP", mode="REQUIRED"),
],
write_disposition=bigquery.WriteDisposition.WRITE_TRUNCATE,
max_bad_records=0,
)
client = bigquery.Client(project=os.environ["BQ_PROJECT"])
job = client.load_table_from_json(
rows, os.environ["BQ_TABLE"], job_config=config,
location=os.environ["BQ_LOCATION"],
)
job.result() # Odota valmistumista; virhe keskeyttää ohjelman.
print(json.dumps({"status": "loaded", "rows": len(rows), "job_id": job.job_id}))
if __name__ == "__main__":
main()Tyhjää kokonaistulosta käsitellään tässä virheenä, jotta poikkeuksellinen API-vastaus ei tyhjennä aiempaa taulua. Jos tyhjä luettelo on lähteessä normaali tila, sille tarvitaan erikseen sovittu ja testattu käsittely. Puuttuvia yksittäisiä tietueita koodi ei pysty tunnistamaan ilman vertailutietoa lähteen odotetusta kattavuudesta.
Esimerkki säilyttää vain nykytilan ja kokoaa aineiston muistiin. Se ei tallenna muutoshistoriaa, varmista lähteen kaikkien sivujen muodostavan samaa ajanhetkeä eikä estä erillisten ajojen päällekkäisyyttä. Näihin palataan jäljempänä.
3. Vie ohjelma Cloud Run Jobsiin
Cloud Run Job ajaa konttiin pakatun ohjelman ja päättyy työn valmistuttua. Se sopii eräsiirtoon, jonka ei tarvitse ylläpitää HTTP-palvelinta. Cloud Run -palvelu puolestaan soveltuu esimerkiksi saapuvien webhook-pyyntöjen vastaanottamiseen. Molempia voi käyttää samassa kokonaisuudessa eri tehtäviin.
Tallenna main.py, requirements.txt ja Dockerfile samaan hakemistoon. Alla olevat riippuvuusversiot on käytetty esimerkin paikallisessa tarkistuksessa. Tuotantoprojektissa myös välilliset riippuvuudet lukitaan ja päivitykset testataan versionhallinnassa.
- Google Cloud: Python-työn rakentaminen ja julkaisu Cloud Runiin
- Google Cloud: Cloud Run Jobsin suorittaminen
requests==2.34.2
google-cloud-bigquery==3.45.0FROM python:3.12-slim
WORKDIR /app
COPY requirements.txt .
RUN pip install --no-cache-dir -r requirements.txt
COPY main.py .
ENV PYTHONUNBUFFERED=1
CMD ["python", "main.py"]gcloud run jobs deploy api-to-bigquery \
--project=oma-projekti \
--region=europe-north1 \
--source=. \
--service-account=api-loader@oma-projekti.iam.gserviceaccount.com \
--set-env-vars="API_URL=https://api.example.com/products,BQ_PROJECT=oma-projekti,BQ_TABLE=oma-projekti.raw_api.products_snapshot,BQ_LOCATION=EU" \
--set-secrets="API_TOKEN=partner-api-token:1" \
--tasks=1 --parallelism=1 --max-retries=1 --task-timeout=20m
gcloud run jobs execute api-to-bigquery \
--project=oma-projekti --region=europe-north1 --waitKomennot olettavat, että projektissa on laskutus, tarvittavat rajapinnat käytössä, BigQuery-dataset ja Secret Managerin salaisuus partner-api-token versiona 1. Julkaisijalla ja lähdekoodin kontiksi rakentavalla palvelutilillä täytyy olla lähdekoodijulkaisun edellyttämät oikeudet. BQ_LOCATION viittaa datasetin sijaintiin; esimerkissä se on EU. Sovittakaa myös Cloud Runin alue organisaation vaatimuksiin.
Testaa ensimmäinen ajo rajatulla aineistolla. --wait odottaa Cloud Run -ajon päättymistä, ja Pythonin job.result() odottaa sen sisällä BigQuery-latauksen valmistumista. Molemmat onnistumiset tarvitaan.
4. Erota tunnukset ja käyttöoikeudet tehtävien mukaan
Cloud Run Jobin ajonaikainen palvelutili api-loader tarvitsee oikeuden käynnistää BigQuery-töitä sekä kirjoittaa kohdedatasettiin. Tyypillinen rajaus on BigQuery Job User projektissa ja BigQuery Data Editor vain tarvittavassa datasetissä. Salaisuuden lukemiseen annetaan Secret Manager Secret Accessor vain kyseiselle salaisuudelle.
BigQueryn Python-kirjasto käyttää Cloud Runissa palvelutilin tunnistautumista Application Default Credentials -mekanismin kautta. Palvelutilin JSON-avainta ei tarvita konttiin. Ulkoisen API:n tunniste luetaan Secret Managerista; yllä oleva julkaisu sitoo sen ympäristömuuttujaan.
Lähdejärjestelmän tunnukselle rajataan vain siirrossa tarvittava lukuoikeus. Lokitetaan rivimääriä, ajon tunnisteita ja virheluokkia. Tunnuksia, henkilötietoja tai kokonaisia API-vastauksia ei käytetä tavallisena virhelokina. Avaimen vaihto ja vanhan version käytöstä poistaminen kuuluvat ylläpitoon.
5. Ajasta yksittäinen siirto Cloud Schedulerilla
Kun manuaalinen ajo on tarkistettu, Cloud Scheduler voi käynnistää saman työn esimerkiksi joka aamu. Ajastuksen ei tarvitse tietää Python-ohjelman sisäisiä yksityiskohtia. Se kutsuu Cloud Run Jobsin suoritusrajapintaa, ja työ hoitaa haun sekä latauksen.
Cloud Runin työn Triggers-näkymästä voidaan lisätä Scheduler-ajastus. Valitkaa ajastus, aikavyöhyke ja erillinen käynnistävä palvelutili. Käynnistäjälle annetaan Cloud Run Invoker -oikeus kyseiseen työhön. Kun HTTP-kohteena on run.googleapis.com-rajapinta, ajastimen tunnistautuminen käyttää OAuth-tokenia.
Schedulerin onnistunut käynnistyspyyntö ei vielä tarkoita, että tiedonsiirto valmistui. Seuratkaa erikseen Cloud Run -ajon lopputulosta ja BigQueryssa olevan aineiston tuoreutta. Scheduler voi myös toimittaa saman ajastetun kutsun useammin kuin kerran, joten siirron pitää kestää uusinta.
6. Tee uusinnasta turvallinen ennen ajastusten lisäämistä
Idempotenssi tarkoittaa tässä, että saman aineiston käsittely uudelleen ei monista liiketoimintatietoja. Esimerkin tilannekuvalataus korvaa taulun sisällön, joten sama id ei kerry sinne uudelleen append-lisäyksenä. Noutoaika voi silti muuttua, eikä lähteen sisältö välttämättä ole toisella ajolla sama.
Yksi tehtävä ja parallelism=1 rajoittavat vain yhden Cloud Run -suorituksen sisäistä rinnakkaisuutta. Ne eivät estä kahta erillistä suoritusta ajautumasta päällekkäin. Hidas vanha ajo voisi julkaista tilannekuvansa uuden jälkeen. Kun useita käynnistäjiä tai pitkiä ajoja on mahdollista, tarvitaan esimerkiksi jaetussa tallennuksessa atomisesti varattava lukko vanhenemisaikoineen tai julkaisun version tarkistus.
HTTP-pyynnön uusinta, Cloud Runin tehtäväuusinta ja koko työnkulun uusinta ovat eri tasoja. Rajaa niiden yhteisvaikutus, ettei lähderajapinta saa häiriössä hallitsematonta määrää pyyntöjä. Epäselvässä BigQuery-latauksen aikakatkaisussa selvitetään työn tunnisteen perusteella, valmistuiko lataus, ennen uuden julkaisun aloittamista.
7. Kun kokonaislataus kasvaa, siirry muutosten hakuun
Inkrementaalinen siirto hakee vain uuden tai muuttuneen aineiston. Lähde voi tarjota esimerkiksi updated_since-suodatuksen, muutoskursorin tai muutoslokin. Etenemispiste tallennetaan pysyvään tilaan ja vahvistetaan vasta sen jälkeen, kun kyseinen erä on turvallisesti tallessa. Kontin muisti ei ole seuraavan ajon tilavarasto.
Aikaleimarajauksessa käytetään usein päällekkäistä hakuikkunaa, jotta myöhässä näkyvät päivitykset eivät jää väliin. Rivien kaksoiskappaleet poistetaan avaimen ja lähteen version tai muutosajan perusteella. BigQueryssa muutokset voidaan ladata välitauluun ja yhdistää kohdetauluun MERGE-lauseella. Saman avaimen useat lähderivit käsitellään ennen yhdistämistä.
Pelkät päivitysajat eivät ratkaise poistoja. Jos lähde ei palauta poistotapahtumaa tai poistomerkintää, katoavat tietueet on tunnistettava esimerkiksi erillisellä täsmäytyksellä. Sivutuksesta tarkistetaan lisäksi, voiko aineisto muuttua haun aikana: kursori ei itsessään takaa yhtenäistä tilannekuvaa.
Historiatäydennys eli backfill tarkoittaa sovitun vanhan jakson lataamista uudelleen. Säilyttäkää ajossa käsiteltävä aikaväli erillään suoritusajasta, jotta esimerkiksi elokuun aineisto voidaan palauttaa syyskuussa. Suunnitelkaa, miten täydennys ja päivittäinen ajo välttävät saman kohteen ristiriitaiset päivitykset.
8. Erota raaka-aineisto, tarkistukset ja raportointimallit
Kun aineisto ei enää mahdu hyvin muistiin tai sitä pitää pystyä käsittelemään uudelleen, tallentakaa noudetut erät esimerkiksi Cloud Storageen NDJSON- tai Parquet-tiedostoina. BigQuery voi ladata tiedostoja erätyönä. Bucketin ja BigQuery-datasetin sijaintien on sovittava yhteen Googlen sijaintisääntöjen mukaan.
Anna jokaiselle ajolle oma tunniste ja tiedostopolku. Merkitse erä valmiiksi vasta, kun sen kaikki sivut ovat tallessa. Osittainen haku ei saa näyttää valmiilta aineistolta. Säilytysajat ja käyttöoikeudet määritellään myös raaka-aineistolle.
Raportointikerroksessa lähteen nimet, tyypit ja liiketoimintamääritelmät muutetaan yhteiseen muotoon esimerkiksi SQL:llä ja dbt:llä. Testaa avainten yksilöllisyys, pakolliset kentät, rivimäärän poikkeamat, aineiston tuoreus ja sovitut summat. Schema eli tietorakenne pidetään eksplisiittisenä; lähteen kenttämuutos arvioidaan ennen sen päästämistä raportille.
9. Milloin orkestrointi kannattaa ottaa mukaan?
Ajastus vastaa kysymykseen, milloin työ käynnistetään. Orkestrointi hallitsee myös sitä, mitä pitää valmistua ensin, mitä voidaan ajaa rinnakkain ja mitä tehdään virheen jälkeen. Yksi itsenäinen päivittäinen siirto voi toimia hyvin Schedulerilla ja Cloud Run Jobilla. Erillinen orkestrointityökalu ei ole lähtövaatimus.
Orkestrointi kannattaa ottaa mukaan, kun ajoketjun oikeellisuus tai ylläpito alkaa riippua useiden tehtävien yhteispelistä. Työkalun valintaa ei ratkaise yksin rivimäärä tai rajapintojen määrä.
Esimerkiksi CRM:n siirto ja laskutuksen siirto voidaan tehdä rinnakkain. Asiakaskohtainen kannattavuusmalli ajetaan vasta, kun molemmat saman raportointijakson aineistot ovat valmiit. Raportointiin tarkoitettu versio julkaistaan vasta laatutarkistusten jälkeen. Kellonajat 02.00, 02.15 ja 02.30 eivät takaa tätä riippuvuutta.
- Yhden vaiheen onnistuminen on seuraavan vaiheen edellytys.
- Usean lähteen pitää vastata samaa ajanjaksoa tai sovittua tuoreutta.
- Virheestä pitää jatkaa oikeasta vaiheesta ilman koko historian uudelleenhakua.
- Vanhoja jaksoja täydennetään säännöllisesti ja niiden etenemistä pitää seurata.
- Tehtävät jakautuvat usealle vastuuhenkilölle ja kokonaisuuden tilan pitää näkyä yhdessä paikassa.
Orkestraattori ei korjaa puuttuvia sivuja, väärää kohdistusta tai uusinnassa monistuvia rivejä. Siirtokoodin, tilan tallennuksen ja datan tarkistusten pitää toimia myös sen alla.
10. Workflows, Dagster vai Airflow?
Google Cloud Workflows sopii esimerkiksi rajattuun palvelukutsujen ketjuun: käynnistä Cloud Run Job, odota sen valmistumista ja käynnistä seuraava vaihe. Cloud Runin Workflows-liitin osaa odottaa pitkäkestoisen operaation päättymistä. Määrittele odotuksen aikaraja ja virhepolut; tavallisen HTTP-kutsun hyväksytty vastaus ei yksin riitä työn valmistumisen todisteeksi.
Dagster on harkinnan arvoinen, kun kokonaisuutta halutaan hallita tuotettavina aineistoina ja niiden riippuvuuksina. Esimerkiksi päivämäärän mukaan jaettujen aineistojen täydennykset ja valmistumistilanne voidaan tehdä näkyviksi. Osituksen ja ajossa käsiteltävän aikavälin on vastattava myös siirtokoodia.
Apache Airflow kuvaa tehtävien riippuvuuksia DAG-työnkulkuina. Se on luonteva vaihtoehto, jos organisaatio jo käyttää Airflow’ta tai tarvitsee yhteisen ympäristön monille ajoketjuille. Itse ylläpidetyn ja hallitun ympäristön välillä arvioidaan kustannuksia, päivityksiä, käyttöoikeuksia ja oman tiimin ylläpitokykyä.
Ensimmäistä Python-siirtoa ei tarvitse kirjoittaa uudelleen vain orkestraattorin vuoksi. Selkeästi rajattu Cloud Run Job voi säilyä työn suorittajana, ja uusi kerros vastaa käynnistämisestä, riippuvuuksista sekä seurannasta. Suuria aineistoja ei kuljeteta työnkulun välisissä viesteissä: välitä tiedostopolku, taulu tai erätunniste.
11. Milloin tarvitaan webhookeja tai viestijonoa?
Jos tieto tarvitaan pian muutoksen jälkeen ja lähde tukee webhookeja, jatkuvan kyselyn rinnalle voidaan rakentaa tapahtumapohjainen vastaanotto. Cloud Run -palvelu tarkistaa lähettäjän, tallentaa tapahtuman esimerkiksi Pub/Subiin ja kuittaa pyynnön vasta onnistuneen tallennuksen jälkeen. Varsinainen käsittely tehdään erikseen.
Pub/Sub toimittaa oletuksena viestin vähintään kerran, joten sama tapahtuma voi tulla käsittelyyn uudelleen. Käyttäkää lähteen tapahtumatunnistetta tai muuta sovittua avainta kaksoiskäsittelyn estämiseen. Huomioikaa myös tapahtumien järjestys, uudelleenyritykset ja viestit, joita ei saada käsiteltyä.
Webhook ei yksin korvaa historiatäydennystä tai täsmäytystä. Rajapinnan kautta tehtävä jaksottainen tarkistus voi edelleen olla tarpeen puuttuvien muutosten löytämiseksi. Reaaliaikaisuutta kannattaa rakentaa silloin, kun nopeudella on selvä käyttötarkoitus ja sen lisäylläpito on perusteltu.
12. Mitä tuotantovalmiista siirrosta pitää tietää?
- Omistaja: kuka vastaa rajapinnasta, tiedon laadusta ja hälytyksiin reagoinnista?
- Tuoreus: milloin viimeinen hyväksytty erä valmistui ja mitä lähteen ajanjaksoa se kattaa?
- Kattavuus: haettiinko kaikki sivut, miten poistot näkyvät ja voiko lähteen tietoja verrata kohteen lukuihin?
- Toipuminen: voiko ajon uusia tai tietyn jakson palauttaa turvallisesti?
- Skeema: miten lähteen kenttämuutos tunnistetaan ja julkaistaan hallitusti?
- Kustannukset: mitä aiheutuu API-lisensseistä, Cloud Runista, tallennuksesta, BigQuery-käsittelystä ja orkestroinnista?
- Ylläpito: missä ovat lähdekoodi, riippuvuudet, ympäristöasetukset ja ohjeet tunnusten vaihtoon?
Aloittakaa yhdellä lähteellä ja rajatulla aineistolla, mutta määritelkää heti, mitä onnistunut ajo tarkoittaa. Lisätkää muutoshaku, historiatäydennykset ja orkestrointi niiden ratkaisemien tarpeiden perusteella. Näin ensimmäinen toteutus antaa pohjan jatkokehitykselle.
Tarvitsetteko ylläpidettävän rajapintasiirron?
Vectura suunnittelee ja toteuttaa rajapintaintegraatioita, BigQuery-datamalleja ja ajoketjuja. Kokemusta raportoinnin ja dataputkien kehittämisestä on vuodesta 2018. Aloitamme lähteen mahdollisuuksista, tiedon käyttötarpeesta ja ylläpidon vastuista.
Tutustu data engineering -palveluihin