Inmata data från ett API i pipelines

Att ingäta från ett API innebär att hämta data över HTTP från en webbtjänst, vanligtvis som paginerad JSON, istället för att läsa från en fil eller databas. Till skillnad från filer eller en meddelandebuss finns det ingen inbyggd generisk API-källa, så du hanterar autentisering, paginering och hastighetsgränser själv. Lakeflow-pipelines stöder tre mönster för inmatning från ett godtyckligt API. Vilken som passar beror på dina volym- och uppdateringsbehov.

Important

Innan du skriver någon anpassad API-insamlingskod, kontrollera om en hanterad connector redan finns för din källkod. Lakeflow Connect levererar inbyggda kopplingar för många vanliga SaaS-API:er, såsom Salesforce, Workday, ServiceNow och Google Analytics, och det finns också ett växande antal partnerkopplingar. Om en connector täcker din källa hanterar den autentisering, paginering och inkrementell extraktion åt dig, och det är nästan alltid mindre arbete än en handrullad inmatning. Se Hanterade kontakter i Lakeflow Connect. Använd mönstren nedan endast när ingen kontakt passar.

Förutsättningar

  • En processkedja. För att skapa en, se Lakeflow pipelines-handledningar.
  • API-uppgifter, såsom en token eller nyckel, lagras som en Azure Databricks-hemlighet. Hårdkoda aldrig inloggningsuppgifter i pipelinens källkod. Se Hemlig hantering.
  • Nätverksåtkomst från din pipeline-beräkning till API-endpointen.
  • Bekantskap med strömmande tabeller och materialiserade vyer, de datamängdstyper som dessa mönster producerar. Se Strömningstabeller och Materialiserade vyer.

Välj ett mönster

Det finns ingen inbyggd generisk REST-API källkod i pipelines, så när du hämtar från ett godtyckligt API, välj ett av tre mönster baserat på datavolym och hur ofta du tar in:

Pattern Använd när
Periodiska drag som en materialiserad syn Nyttolaster är små till medelstora och hämtas en gång per pipelinekörning, såsom referensdata, dagliga FX-kurser eller ett paginerat men avgränsbart API.
Python Data Source API Du behöver avfråga ett API med hög volym eller ett strömmande API inkrementellt, med sparade kontrollpunkter för förloppet så att en omstart inte behöver läsa in allt på nytt.
Frikopplad datainläsning med Auto Loader Du vill isolera API-specifika egenheter från din transformationslogik och få exakt en gång-filspårning gratis.

Mönster 1: Periodiska drag som en materialiserad vy

För små till medelstora payloads som hämtas en gång per pipelinekörning, skriv en Python-funktion som anropar API:et och returnerar en Spark DataFrame. Eftersom datasetet är en materialiserad vy kör pipelinen om funktionen i sin helhet och på ett idempotent sätt varje gång pipelinen uppdateras.

Följande steg visar hur du bygger en materialiserad vy med periodiska drag:

  1. Lagra API-token i en hemlighet och mappa den sedan till en Spark-konfigurationsegenskap i dina pipeline-inställningar så att pipelinekoden kan läsa den. Lägg till egenskapen i blocket spark_conf i pipelinens klusterkonfiguration:

    {
      "clusters": [
        {
          "spark_conf": {
            "api.token": "{{secrets/<scope-name>/<secret-name>}}"
          }
        }
      ]
    }
    

    Koden i nästa steg läser detta värde med spark.conf.get("api.token"). För mer om att konfigurera hemligheter i pipeline-inställningar, se Säker åtkomst till lagringsuppgifter med hemligheter i en pipeline.

  2. Definiera en materialiserad vy som anropar API:et och returnerar svaret som en DataFrame:

    import requests
    from pyspark import pipelines as dp
    from pyspark.sql import Row
    
    @dp.materialized_view(
        name="exchange_rates_bronze",
        comment="Daily FX rates pulled from a public REST API",
    )
    def exchange_rates_bronze():
        resp = requests.get(
            "https://api.example.com/v1/rates",
            params={"base": "USD"},
            headers={"Authorization": f"Bearer {spark.conf.get('api.token')}"},
            timeout=30,
        )
        resp.raise_for_status()
        rates = resp.json()["rates"]
        rows = [Row(currency=k, rate=float(v), as_of_date=resp.json()["date"]) for k, v in rates.items()]
        return spark.createDataFrame(rows)
    
  3. Hantera paginering i funktionen genom att loopa över sidor och sammanfoga resultaten innan du returnerar DataFrame:

    import requests
    from pyspark import pipelines as dp
    from pyspark.sql import Row
    
    @dp.materialized_view(
        name="customers_bronze",
        comment="Customers pulled from a paginated REST API",
    )
    def customers_bronze():
        token = spark.conf.get("api.token")
        rows = []
        url = "https://api.example.com/v1/customers"
        while url:  # follow the API's next-page cursor until exhausted
            resp = requests.get(
                url,
                headers={"Authorization": f"Bearer {token}"},
                timeout=30,
            )
            resp.raise_for_status()
            payload = resp.json()
            rows.extend(Row(**record) for record in payload["data"])
            url = payload.get("next")  # None on the last page
        return spark.createDataFrame(rows)
    

    Lägg till åter- och tillbakagångslogik kring begäran om motståndskraft.

Detta mönster läser om hela API-svaret vid varje pipelineuppdatering, så använd det endast när payloaden är begränsad. Vid inkrementell läsning använder du mönster 2.

Mönster 2: Högvolyms- eller strömmande API:er med Python Data Source API

För API:er behöver du polla inkrementellt med offset tracking, implementera en anpassad datakälla med Sparks Python Data Source API. Detta ger dig korrekt streaming-semantik, inklusive checkpoint-framsteg och inkrementella läsningar, så att en omstart återupptas från senaste offset istället för att hämta hela API:et igen.

Följande steg visar hur du kan ta in från en anpassad datakälla:

  1. Implementera en DataSource och DataSourceStreamReader som anropar API:et och spårar läsoffsetet. För detaljer om hur man skapar en anpassad datakälla, se PySpark anpassade datakällor.

  2. Registrera datakällan så att pipelinen kan referera till den med formatnamn:

    spark.dataSource.register(MyApiDataSource)
    
  3. Läs från den registrerade källan i en streamingtabell:

    from pyspark import pipelines as dp
    
    @dp.table(name="events_bronze")
    def events_bronze():
        return spark.readStream.format("my_api_source").load()
    

Mönster 3: Koppla loss intagning från ett schemaläggt jobb och Auto Loader

Ett vanligt produktionsmönster är att separera API-anropet från pipelinen. Ett schemalagt jobb landar de råa API-svaren som filer i en Unity Catalog-volym, och pipelinen plockar upp dem med Auto Loader. Detta isolerar API-specifika egenheter som paginering och hastighetsgränser från din deklarativa transformationslogik, och ger dig Auto Loaders exakt en gång-filspårning gratis.

Följande steg visar hur du kan koppla bort intagning från ett schemalagt jobb:

  1. Skriv en anteckningsbok eller ett skript som anropar API:et och skriver de råa JSON-svaren till en Unity Catalog-volym. Läs API-uppgifterna från en hemlighet. Se Hemlig hantering.

    import requests, json, time
    
    token = dbutils.secrets.get(scope="<scope-name>", key="<secret-name>")
    volume_path = "/Volumes/main/raw/landing/api_events"
    
    resp = requests.get(
        "https://api.example.com/v1/events",
        headers={"Authorization": f"Bearer {token}"},
        timeout=30,
    )
    resp.raise_for_status()
    # One file per run; the pipeline's Auto Loader tracks which files it has ingested.
    with open(f"{volume_path}/events_{int(time.time())}.json", "w") as f:
        json.dump(resp.json()["data"], f)
    
  2. Schemalägg anteckningsboken eller skriptet så att det körs självständigt med Lakeflow Jobs. Se Lakeflow Jobs.

  3. I din pipeline, definiera en strömningstabell som läser de landade filerna med Auto Loader:

    from pyspark import pipelines as dp
    
    @dp.table(name="api_events_bronze")
    def api_events_bronze():
        return (
            spark.readStream.format("cloudFiles")
                .option("cloudFiles.format", "json")
                .load("/Volumes/main/raw/landing/api_events")
        )
    

För mer om pålitlig filinsamling med Auto Loader, se Ladda filer från molnobjektlagring och Vad är Auto Loader?.

Bästa praxis för API-intagning

  • Håll hemligheter utanför källkoden. Lagra API-token och nycklar i hemliga omfång i Azure Databricks och läs dem under körning. Se Hemlig hantering.
  • Bekräfta svar tidigt. Lägg till valideringsregler för de inlästa raderna för att fånga felaktiga API-svar innan de skickas vidare.
  • Hantera paginering och hastighetsgränser. Iterera över sidorna och lägg till återförsök med fördröjning mellan försöken så att ett tillfälligt fel inte gör att hela uppdateringen misslyckas.

Ytterligare resurser