Pular para o conteúdo principal

Ingest data from an API in pipelines

Ingerir dados de uma API significa extrair dados via HTTP de um serviço web, geralmente como JSON paginado, em vez de ler de um arquivo ou de um banco de dados. Ao contrário de arquivos ou de um barramento de mensagens, não há uma fonte de API genérica integrada, portanto, você mesmo gerencia a autenticação, a paginação e os limites de taxa. Os Lakeflow pipelines suportam três padrões para ingestão a partir de uma API arbitrária. Qual deles se ajusta depende do seu volume e das suas necessidades de refresh.

importante

Before you write any custom API-ingestion code, check whether a managed connector already exists for your source. Lakeflow Connect ships built-in connectors for many common software as a service (SaaS) APIs, such as Salesforce, Workday, ServiceNow, and Google Analytics, and there's a growing set of partner connectors as well. If a connector covers your source, it handles authentication, pagination, and incremental extraction for you, and it's almost always less work than a hand-rolled ingestion. See Managed connectors in Lakeflow Connect. Use the patterns below only when no connector fits.

Pré-requisitos

  • Um pipeline. Para criar um, consulte tutoriais de Lakeflow pipelines.
  • API credentials, such as a token or key, stored as a Databricks secret. Never hardcode credentials in pipeline source code. See Secret management.
  • Acesso de rede do seu compute de pipeline ao endpoint da API.
  • Familiarity with streaming tables and materialized views, the dataset types these patterns produce. See Streaming tables and Materialized views.

Choose a pattern

Não há uma fonte de API REST genérica nativa em pipelines, portanto, ao extrair de uma API arbitrária, escolha um dos três padrões com base no volume de dados e na frequência de ingestão:

Padrão

Usar quando

Extrações periódicas como uma view materializada

Os payloads são de pequeno a médio porte e extraídos uma vez por execução de pipeline, como dados de referência, taxas de câmbio diárias ou uma API paginada, mas limitável.

API de fonte de dados Python

É necessário realizar o polling de uma API de alto volume ou de transmissão de forma incremental, com progresso verificado por ponto de verificação para que uma reinicialização não releia tudo.

Decoupled ingestion with Auto Loader

Você deseja isolar peculiaridades específicas da API de sua lógica de transformação e obter acompanhamento de arquivos exatamente uma vez gratuitamente.

Padrão

Usar quando

Extrações periódicas como uma view materializada

Os payloads são de pequeno a médio porte e extraídos uma vez por execução de pipeline, como dados de referência, taxas de câmbio diárias ou uma API paginada, mas limitável.

API de fonte de dados Python

É necessário realizar o polling de uma API de alto volume ou de transmissão de forma incremental, com progresso verificado por ponto de verificação para que uma reinicialização não releia tudo.

Decoupled ingestion with Auto Loader

Você deseja isolar peculiaridades específicas da API de sua lógica de transformação e obter acompanhamento de arquivos exatamente uma vez gratuitamente.

Padrão 1: Extrações periódicas como uma visualização materializada

Para payloads de pequeno a médio porte extraídos uma vez por execução de pipeline, escreva uma função Python que chame a API e retorne um Spark DataFrame. Como o dataset é uma view materializada, o pipeline executa a função novamente de forma completa e idempotente toda vez que o pipeline é atualizado.

Os passos a seguir mostram como criar uma materialized view com extrações periódicas:

  1. Armazene o token da API em um segredo e, em seguida, mapeie-o para uma propriedade de configuração do Spark nas configurações do seu pipeline para que o código do pipeline possa lê-lo. Adicione a propriedade ao bloco spark_conf da configuração de cluster do pipeline:

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

    O código na próxima etapa lê esse valor com spark.conf.get("api.token"). Para obter mais informações sobre como configurar segredos nas configurações do pipeline, consulte Acessar credenciais de armazenamento com segurança usando segredos em um pipeline.

  2. Defina uma view materializada que chama a API e retorna a resposta como um DataFrame:

    Python
    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={&quot;base&quot;: &quot;USD&quot;},
    headers={&quot;Authorization&quot;: f&quot;Bearer {spark.conf.get('api.token')}&quot;},
    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. Trate a paginação dentro da função fazendo um loop pelas páginas e concatenando os resultados antes de retornar o DataFrame:

    Python
    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={&quot;Authorization&quot;: f&quot;Bearer {token}&quot;},
    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)

    Adicione lógica de repetição e backoff em torno da solicitação para resiliência.

Este padrão relê a resposta completa da API a cada atualização de pipeline, portanto, use-o apenas quando a carga útil for limitada. Para leituras incrementais, use o padrão 2.

Pattern 2: High-volume or streaming APIs with the Python Data Source API

For APIs you need to poll incrementally with offset tracking, implement a custom data source using Spark's Python Data Source API. This gives you proper streaming semantics, including checkpointed progress and incremental reads, so a restart resumes from the last offset instead of pulling the whole API again.

Os passos a seguir mostram como ingerir dados de uma fonte de dados personalizada:

  1. Implement a DataSource and DataSourceStreamReader that call the API and track the read offset. For details on authoring a custom data source, see PySpark custom data sources.

  2. Register the data source so the pipeline can reference it by format name:

    Python
    spark.dataSource.register(MyApiDataSource)
  3. Read from the registered source in a streaming table:

    Python
    from pyspark import pipelines as dp

    @dp.table(name="events_bronze")
    def events_bronze():
    return spark.readStream.format("my_api_source").load()

Padrão 3: desacoplar a ingestão com um job agendado e o Auto Loader

Um padrão de produção comum é separar a chamada de API do pipeline. Um Job agendado salva as respostas brutas da API como arquivos em um volume do Unity Catalog, e o pipeline os captura com o Auto Loader. Isso isola peculiaridades específicas da API, como paginação e limites de taxa, da sua lógica de transformação declarativa, e oferece o acompanhamento de arquivos exactly-once do Auto Loader gratuitamente.

The following steps show you how to decouple ingestion with a scheduled job:

  1. Escreva um notebook ou script que chame a API e grave as respostas JSON brutas em um volume do Unity Catalog. Leia as credenciais da API de um segredo. Consulte Gerenciamento de segredos.

    Python
    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={&quot;Authorization&quot;: f&quot;Bearer {token}&quot;},
    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. Programe o notebook ou script para ser executado por conta própria com o Lakeflow Jobs. See Lakeflow Jobs.

  3. Em seu pipeline, defina uma tabela de transmissão que leia os arquivos recebidos com o Auto Loader:

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

For more about reliable file ingestion with Auto Loader, see Load files from cloud object storage and What is Auto Loader?.

Práticas recomendadas para ingestão de API

  • Keep secrets out of source code. Store API tokens and keys in Databricks secret scopes and read them at runtime. See Secret management.
  • Valide as respostas antecipadamente. Adicione expectativas nas linhas ingeridas para detectar respostas de API malformadas antes que elas fluam para o downstream.
  • Lidar com paginação e limites de taxa. Faça um loop pelas páginas e adicione uma nova tentativa com backoff para que uma falha transitória não interrompa toda a atualização.

Outros recursos