Ввод данных из API в конвейерах

Загрузка данных из API означает извлечение данных через HTTP из веб-сервиса, обычно в виде страницированного JSON, а не чтения из файла или базы данных. В отличие от файлов или шины сообщений, здесь нет встроенного универсального API-источника, поэтому вы сами занимаетесь аутентификацией, страницированием и ограничениями скорости. Конвейеры Lakeflow поддерживают три шаблона для приема данных из произвольного API. Какой вариант подходит — зависит от ваших потребностей в объёме и обновлении.

Important

Прежде чем писать какой-либо пользовательский код для загрузки API, проверьте, существует ли уже управляемый коннектор для вашего исходного кода. Lakeflow Connect поставляет встроенные коннекторы для многих распространённых API программного обеспечения как услуги (SaaS), таких как Salesforce, Workday, ServiceNow и Google Analytics, а также растёт набор партнёрских коннекторов. Если коннектор покрывает ваш источник, он занимается аутентификацией, пагинацией и постепенным извлечением за вас, и это почти всегда меньше работы, чем ручной загрузка. См. концепции разъёмов Lakeflow Connect. Используйте приведённые ниже схемы только тогда, когда соединитель не подходит.

Необходимые условия

  • Трубопровод. Чтобы создать его, см. руководства по Lakeflow pipelines.
  • Учётные данные API, такие как токен или ключ, хранятся как секрет Azure Databricks. Никогда не закодуйте учетные данные в исходном коде конвейера. Дополнительные сведения см. в разделе Управление секретами.
  • Доступ к сети от вашего конвейера вычисления к API-конечной точке.
  • Знакомство с потоковыми таблицами и материализованными представлениями, а также с типами наборов данных, которые создаются с помощью этих шаблонов. См. потоковые таблицы и материализованные представления.

Выбор шаблона

В конвейерах нет встроенного универсального источника данных REST API, поэтому, когда вы получаете данные из произвольного API, выберите одну из трёх схем в зависимости от объёма данных и частоты загрузки:

Расписание Используйте, если
Периодическое извлечение данных как материализованное представление Объёмы данных небольшие или средние и извлекаются один раз за запуск конвейера, например: справочные данные, ежедневные обменные курсы валют или API с пагинацией, но с ограниченным объёмом данных.
API для источников данных Python Вам нужно инкрементально опрашивать API с большим объёмом данных или потоковый API, с сохранением контрольных точек прогресса, чтобы после перезапуска не приходилось считывать всё заново.
Независимая загрузка данных с Auto Loader Вам нужно выделить специфические особенности API от логики трансформации и получить точное отслеживание файлов бесплатно.

Паттерн 1: Периодические вытяжения как материализованный вид

Для небольших и средних полезных нагрузок, загружаемых один раз за запуск конвейера, напишите Python-функцию, которая вызывает API и возвращает Spark DataFrame. Поскольку набор данных является материализованным представлением, конвейер каждый раз при обновлении повторно запускает функцию полностью, с соблюдением идемпотентности.

Следующие шаги показывают, как построить материализированное изображение с периодическими вытягиваниями:

  1. Храните токен API в секрете, затем сопоставьте его с конфигурационным свойством Spark в настройках конвейера, чтобы код конвейера мог его прочитать. Добавьте это свойство в блок spark_conf конфигурации кластера конвейера:

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

    Код на следующем шаге читает это значение с помощью spark.conf.get("api.token"). Подробнее о настройке секретов в настройках пайплайна см. в статье Безопасный доступ к учетным данным хранилища с помощью секретов в пайплайне.

  2. Определите материализованный вид, который вызывает API и возвращает ответ в виде 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. Обработайте пагинацию внутри функции, перебирая страницы в цикле и объединяя результаты перед тем, как вернуть 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)
    

    Добавьте повторные попытки и логику отступления вокруг запроса на устойчивость.

Этот шаблон пересчитывает полный ответ API при каждом обновлении конвейера, поэтому используйте его только при ограничении полезной нагрузки. Для инкрементальных чтений используйте шаблон 2.

Паттерн 2: Высоконагруженные или потоковые API с API источника данных Python

Для API нужно делать постепенные опросы с офсетным отслеживанием, реализуйте пользовательский источник данных с помощью API Python Data Source от Spark. Это обеспечивает корректную семантику потоковой обработки, включая сохранение прогресса через контрольные точки и инкрементальное чтение данных, так что после перезапуска обработка продолжается с последнего смещения вместо повторного чтения всех данных через API.

Следующие шаги показывают, как принимать данные из пользовательского источника:

  1. Реализуйте DataSource и DataSourceStreamReader, которые вызывают API и отслеживают смещение чтения. Для подробностей о создании пользовательского источника данных см. PySpark пользовательские источники данных.

  2. Зарегистрируйте источник данных, чтобы конвейер мог ссылаться на него по названию формата:

    spark.dataSource.register(MyApiDataSource)
    
  3. Чтение из зарегистрированного источника в потоковой таблице:

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

Паттерн 3: Разделение приёма данных с помощью запланированного задания и Auto Loader

Распространённый производственный шаблон — отделить вызов API от конвейера. Запланированное задание отправляет необработанные ответы API в виде файлов в том Unity Catalog, и конвейер получает их с помощью Auto Loader. Это отделяет особенности, характерные для API, такие как пагинация и ограничения частоты запросов, от вашей декларативной логики преобразования и позволяет без дополнительных усилий использовать отслеживание файлов Auto Loader с гарантией однократной обработки.

Следующие шаги показывают, как отделить употребление от запланированной работы:

  1. Напишите блокнот или скрипт, который вызывает API и записывает необработанные JSON-ответы в том Unity Catalog. Читайте учетные данные API из секрета. Дополнительные сведения см. в разделе Управление секретами.

    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. Запланируйте автоматический запуск блокнота или скрипта с помощью Lakeflow Jobs. Смотрите Задания Lakeflow.

  3. В вашем конвейере определите потоковую таблицу, которая считывает поступившие файлы с помощью 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")
        )
    

Для получения дополнительной информации о надёжном вводе файлов с помощью Auto Loader см. разделы «Загрузка файлов из облачного объектного хранилища » и «Что такое Auto Loader?».

Лучшие практики для загрузки API

  • Держите секреты вне исходного кода. Храните токены и ключи API в секретных областях Azure Databricks и читайте их во время выполнения. Дополнительные сведения см. в разделе Управление секретами.
  • Проверяйте ответы заранее. Добавьте ожидания к поглощённым строкам, чтобы поймать искажённые ответы API до того, как они пойдут вниз по потоку.
  • Занимайтесь пагинацией и ограничениями скорости. Перебирайте страницы и добавьте механизм повторных попыток с нарастающей задержкой, чтобы временный сбой не приводил к сбою всего обновления.

Дополнительные ресурсы