Pular para o conteúdo principal

Ingerir dados de uma API em 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

Antes de escrever qualquer código de ingestão de API personalizado, verifique se já existe um conector gerenciado para sua fonte. O Lakeflow Connect disponibiliza conectores integrada para muitas APIs comuns de software como serviço (SaaS), como Salesforce, Workday, ServiceNow e Google analítica, e há também um conjunto crescente de conectores de parceiros. Se um conector cobrir sua fonte, ele cuidará da autenticação, paginação e extração incremental para você, e quase sempre é menos trabalhoso do que uma ingestão manual. Consulte Conectores gerenciados no Lakeflow Connect. Use os padrões abaixo apenas quando nenhum conector for adequado.

Pré-requisitos

Escolher um padrão

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.

Ingestão desacoplada com o 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.

Ingestão desacoplada com o 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.

Padrão 2: APIs de alto volume ou de transmissão com a API de fonte de dados Python

Para APIs que precisam de consultas incrementais com acompanhamento de deslocamento, implemente uma fonte de dados personalizada usando a API de fonte de dados Python do Spark. Isso oferece a você semântica de transmissão adequada, incluindo progresso com checkpoint e leituras incrementais, para que uma reinicialização seja retomada a partir do último deslocamento em vez de extrair toda a API novamente.

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

  1. Implemente um DataSource e DataSourceStreamReader que chamem a API e rastreiem o deslocamento de leitura. Para obter detalhes sobre como criar uma fonte de dados personalizada, consulte fontes de dados personalizadas do PySpark.

  2. Faça o registro da fonte de dados para que o pipeline possa referenciá-la pelo nome do formato:

    Python
    spark.dataSource.register(MyApiDataSource)
  3. Ler da fonte registrada em uma tabela de transmissão:

    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.

Os passos a seguir mostram como desacoplar a ingestão de um Job agendado:

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

Para saber mais sobre a ingestão confiável de arquivos com o Auto Loader, consulte Carregar arquivos do armazenamento de objetos na nuvem e O que é o Auto Loader?.

Práticas recomendadas para ingestão de API

  • Mantenha os segredos fora do código-fonte. Armazene tokens e chaves de API em Secret Scope do Databricks e leia-os em Runtime. Veja Gestão Secreta.
  • 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