Pobieranie danych z API w potokach

Pobieranie danych z interfejsu API oznacza pobieranie ich za pośrednictwem protokołu HTTP z usługi sieciowej, zwykle w postaci stronicowanego JSON-a, zamiast odczytywania danych z pliku lub bazy danych. W przeciwieństwie do plików czy magistrali wiadomości, nie ma wbudowanego, uniwersalnego źródła API, więc sam zajmujesz się uwierzytelnianiem, paginacją i ograniczeniami prędkości. Potoki Lakeflow obsługują trzy wzorce pobierania z dowolnego API. To, który z nich będzie odpowiedni, zależy od Twoich potrzeb dotyczących wolumenu i częstotliwości odświeżania.

Important

Zanim napiszesz jakikolwiek własny kod do pobierania API, sprawdź, czy istnieje już zarządzany konektor dla Twojego źródła. Lakeflow Connect dostarcza wbudowane konektory dla wielu popularnych API oprogramowania jako usługi (SaaS), takich jak Salesforce, Workday, ServiceNow i Google Analytics, a także rośnie liczba partnerskich łączników. Jeśli złącze obsługuje dane źródło, zajmie się za ciebie uwierzytelnianiem, paginacją i ekstrakcją przyrostową, a to niemal zawsze wymaga mniej pracy niż własnoręcznie przygotowany mechanizm pozyskiwania danych. Zobacz Łączniki zarządzane w programie Lakeflow Connect. Używaj poniższych wzorów tylko wtedy, gdy nie pasuje żaden łącznik.

Wymagania wstępne

Wybieranie wzorca

W pipeline'ach nie ma natywnego źródła REST-API generycznego, więc gdy pobierasz z dowolnego API, wybierz jeden z trzech wzorców na podstawie objętości danych i częstotliwości pobierania:

Pattern Użyj, gdy
Okresowe przyciągania jako zmaterializowany widok Ładunki danych są małe lub średnie i są pobierane raz na każde uruchomienie potoku przetwarzania, na przykład dane referencyjne, dzienne kursy walut lub interfejs API z paginacją, ale o dającym się określić zakresie.
Python Data Source API Musisz przyrostowo odpytywać API o dużym wolumenie lub API strumieniowe, z zapisywaniem punktów kontrolnych postępu, aby po restarcie nie trzeba było odczytywać wszystkiego od początku.
Rozdzielone pozyskiwanie danych za pomocą Auto Loadera Chcesz odizolować specyficzne dla API niedociągnięcia od logiki transformacji i uzyskać darmowe śledzenie plików dokładnie raz.

Wzorzec 1: Okresowe pobieranie w formie widoku zmaterializowanego

Dla małych i średnich ładunków danych pobieranych raz na jedno uruchomienie potoku należy napisać funkcję w języku Python, która wywołuje interfejs API i zwraca ramkę danych Spark. Ponieważ zbiór danych jest widokiem materializowanym, potok ponownie uruchamia funkcję w całości i w sposób idempotentny za każdym razem, gdy potok jest aktualizowany.

Poniższe kroki pokazują, jak utworzyć widok materializowany z okresowym pobieraniem danych:

  1. Zapisz token API w sekretze, a następnie przypisz go do właściwości konfiguracji Spark w ustawieniach potoku, aby kod potoku mógł go odczytać. Dodaj tę właściwość do bloku spark_conf w konfiguracji klastra potoku:

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

    Kod w kolejnym kroku odczytuje tę wartość za pomocą spark.conf.get("api.token"). Aby dowiedzieć się więcej o konfigurowaniu sekretów w ustawieniach potoku, zobacz Bezpieczny dostęp do poświadczeń magazynu za pomocą sekretów w potoku.

  2. Zdefiniuj zmaterializowany widok, który wywołuje API i zwraca odpowiedź jako 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. Zajmij się paginacją w funkcji, zapętlając strony i łącząc wyniki przed zwróceniem 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)
    

    Dodaj logikę powtórek i wycofania się wokół żądania odporności.

Ten wzorzec ponownie odczytuje pełną odpowiedź API przy każdej aktualizacji potoku, więc używaj go tylko, gdy payload jest ograniczony. Do odczytów przyrostowych użyj wzorca 2.

Wzorzec 2: wysokowolumenowe lub strumieniowe interfejsy API z interfejsem API źródła danych Python

W przypadku interfejsów API, jeśli chcesz odpytywać je przyrostowo ze śledzeniem przesunięcia, zaimplementuj niestandardowe źródło danych przy użyciu interfejsu Spark Python Data Source API. Zapewnia to prawidłową semantykę strumieniowania, w tym śledzenie postępu za pomocą punktów kontrolnych i odczyty przyrostowe, dzięki czemu po ponownym uruchomieniu proces jest wznawiany od ostatniego offsetu zamiast ponownie pobierać wszystkie dane z API.

Poniższe kroki pokazują, jak pobierać dane z niestandardowego źródła danych:

  1. Zaimplementuj DataSource i DataSourceStreamReader, które wywołują interfejs API i śledzą przesunięcie odczytu. Aby uzyskać szczegółowe informacje na temat tworzenia niestandardowego źródła danych, zobacz niestandardowe źródła danych PySpark.

  2. Zarejestruj źródło danych, aby potok mógł się do niego odwoływać pod nazwą formatu:

    spark.dataSource.register(MyApiDataSource)
    
  3. Odczyt z zarejestrowanego źródła w tabeli strumieniowej:

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

Wzór 3: Odłącz spożycie od zaplanowanej pracy i Auto Loadera

Powszechnym wzorcem produkcyjnym jest oddzielenie wywołania API od potoku przetwarzania. Zaplanowane zadanie wysyła surowe odpowiedzi API jako pliki do woluminu Unity Catalog, a potok odbiera je za pomocą Auto Loadera. Izoluje to specyficzne dla API niuanse, takie jak paginacja i limity liczby żądań, od deklaratywnej logiki transformacji, a także zapewnia w Auto Loader dokładnie jednokrotne śledzenie plików bez dodatkowego wysiłku.

Poniższe kroki pokazują, jak oddzielić spożycie od zaplanowanej pracy:

  1. Napisz notatnik lub skrypt, który wywołuje interfejs API i zapisuje surowe odpowiedzi JSON w woluminie Unity Catalog. Odczytaj dane uwierzytelniające API z sekretu. Zobacz Zarządzanie tajemnicami.

    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. Zaplanuj notatnik lub skrypt do samodzielnego uruchamiania z Lakeflow Jobs. Zobacz Zadania lakeflow.

  3. W swoim pipeline zdefiniuj tabelę streamingową, która odczytuje pliki lądujące za pomocą Auto Loadera:

    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")
        )
    

Więcej informacji o niezawodnym pobieraniu plików za pomocą Auto Loadera znajdziesz w artykule Load files from cloud object storage oraz What is Auto Loader?.

Najlepsze praktyki dotyczące pobierania API

  • Trzymaj sekrety z dala od kodu źródłowego. Przechowuj tokeny i klucze API w tajnych zakresach Azure Databricks i czytaj je w czasie działania. Zobacz Zarządzanie tajemnicami.
  • Potwierdzaj odpowiedzi na początku. Dodaj oczekiwania na wprowadzonych wierszach, aby wychwycić nieprawidłowe odpowiedzi API zanim trafią dalej.
  • Zajmij się paginacją i limitami szybkości. Iterować po stronach i dodać mechanizm ponawiania z wydłużanym odstępem, aby przejściowy błąd nie powodował niepowodzenia całej aktualizacji.

Dodatkowe zasoby