Nota:
El acceso a esta página requiere autorización. Puede intentar iniciar sesión o cambiar directorios.
El acceso a esta página requiere autorización. Puede intentar cambiar los directorios.
Ingerir desde una API significa extraer datos por HTTP desde un servicio web, normalmente como JSON paginado, en lugar de leer de un archivo o una base de datos. A diferencia de los archivos o un bus de mensajes, no hay una API genérica incorporada, así que tú gestionas la autenticación, paginación y límites de tasa. Las pipelines Lakeflow soportan tres patrones para la ingesta desde una API arbitraria. La opción que mejor se adapta depende de tu volumen y de tus necesidades de actualización.
Importante
Antes de escribir cualquier código personalizado para la ingestión de API, comprueba si ya existe un conector gestionado para tu código fuente. Lakeflow Connect incluye conectores integrados para muchas APIs comunes de software como servicio (SaaS), como Salesforce, Workday, ServiceNow y Google Analytics, y también hay un conjunto creciente de conectores de socios. Si un conector cubre su origen, se encarga de la autenticación, paginación y extracción incremental por usted, y casi siempre es menos trabajo que una ingestión enrollada a mano. Consulte Conectores administrados en Lakeflow Connect. Usa los patrones de abajo solo cuando no encaje ningún conector.
Requisitos previos
- Un oleoducto. Para crear uno, consulta los tutoriales de Lakeflow pipelines.
- Credenciales de API, como un token o clave, almacenadas como un secreto de Azure Databricks. Nunca incruste credenciales de forma fija en el código fuente de la canalización. Consulte Administración de secretos.
- Acceso a la red desde su proceso de canalización hasta el punto de conexión de API.
- Familiaridad con las tablas de flujo y las vistas materializadas, y los tipos de conjuntos de datos que generan estos patrones. Consulte Tablas en streaming y vistas materializadas.
Elección de un patrón
No hay una fuente nativa de REST-API genérica en los pipelines, así que cuando extraes de una API arbitraria, elige uno de tres patrones según el volumen de datos y la frecuencia con la que ingieres:
| Pattern | Se utiliza cuando |
|---|---|
| Tirones periódicos como vista materializada | Las cargas son pequeñas o medianas y se obtienen una vez por ejecución de canalización, como datos de referencia, tipos de cambio diarios o una API paginada pero con un límite definible. |
| API de fuentes de datos en Python | Necesitas consultar una API de alto volumen o de streaming de forma incremental, con un progreso controlado para que un reinicio no lo relea todo. |
| Ingesta desacoplada con Auto Loader | Quiere aislar particularidades específicas de la API de su lógica de transformación y obtener el seguimiento exacto de un archivo gratis. |
Patrón 1: Extracciones periódicas como vista materializada
Para cargas útiles pequeñas o medianas extraídas una vez por ejecución de pipeline, escribe una función en Python que llame a la API y devuelva un DataFrame de Spark. Dado que el conjunto de datos es una vista materializada, la canalización vuelve a ejecutar la función de forma completa e idempotente cada vez que se actualiza.
Los siguientes pasos muestran cómo crear una vista materializada con extracciones periódicas:
Guarda el token de la API en secreto y luego mapearlo a una propiedad de configuración de Spark en la configuración de tu pipeline para que el código del pipeline pueda leerlo. Añada la propiedad al bloque
spark_confde la configuración del clúster de la canalización:{ "clusters": [ { "spark_conf": { "api.token": "{{secrets/<scope-name>/<secret-name>}}" } } ] }El código en el siguiente paso lee este valor con
spark.conf.get("api.token"). Para más información sobre cómo configurar secretos en la configuración de la pipeline, consulta Acceso seguro a credenciales de almacenamiento con secretos en una pipeline.Definamos una vista materializada que llama a la API y devuelve la respuesta como 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)Maneja la paginación dentro de la función iterando por las páginas y concatenando los resultados antes de devolver el 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)Añade lógica de reintentos y retrocesos alrededor de la petición de resiliencia.
Este patrón relee la respuesta completa de la API en cada actualización del pipeline, así que úsalo solo cuando la carga útil esté acotada. Para lecturas incrementales, utiliza el patrón 2.
Patrón 2: APIs de alto volumen o streaming con la API de Python Data Source
Si para las API necesitas realizar sondeos incrementales con seguimiento del offset, implementa una fuente de datos personalizada mediante la API de fuentes de datos de Python de Spark. Esto le proporciona una semántica de streaming adecuada, incluido el progreso registrado mediante puntos de control y las lecturas incrementales, de modo que, tras un reinicio, se reanuda desde el último offset en lugar de volver a consultar toda la API.
Los siguientes pasos te muestran cómo ingerir desde una fuente de datos personalizada:
Implemente un
DataSourceyDataSourceStreamReaderque llamen a la API y registren el offset de lectura. Para más detalles sobre la creación de una fuente de datos personalizada, consulte PySpark fuentes de datos personalizadas.Registrar la fuente de datos para que la tubería pueda referenciarla por nombre de formato:
spark.dataSource.register(MyApiDataSource)Leído del origen registrado en una tabla en streaming:
from pyspark import pipelines as dp @dp.table(name="events_bronze") def events_bronze(): return spark.readStream.format("my_api_source").load()
Patrón 3: Desacoplar la ingestión con un trabajo programado y Auto Loader
Un patrón de producción común es separar la llamada API de la pipeline. Un trabajo programado deposita las respuestas sin procesar de la API como archivos en un volumen de Unity Catalog y la canalización las ingiere con Auto Loader. Esto aísla particularidades específicas de la API como la paginación y los límites de velocidad de su lógica de transformación declarativa, y le da el seguimiento de archivos exactamente una vez de Auto Loader gratis.
Los siguientes pasos le muestran cómo desacoplar la ingestión de un trabajo programado:
Escribe un cuaderno o script que llame a la API y escriba las respuestas JSON en bruto a un volumen de Unity Catalog. Lea las credenciales de la API desde un secreto. Consulte Administración de secretos.
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)Programa el cuaderno o script para que se ejecute solo con Lakeflow Jobs. Consulte Trabajos de Lakeflow.
En su canalización, defina una tabla en streaming que lea los archivos de aterrizaje con 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") )
Para más información sobre la ingestión fiable de archivos con Auto Loader, consulta Cargar archivos desde almacenamiento de objetos en la nube y ¿Qué es Auto Loader?.
Mejores prácticas para la ingestión de APIs
- Mantén los secretos fuera del código fuente. Guarda tokens y claves API en los ámbitos secretos de Azure Databricks y léelos en tiempo de ejecución. Consulte Administración de secretos.
- Valida las respuestas con antelación. Añada expectativas sobre las filas ingeridas para detectar respuestas API con formato incorrecto antes de que fluyan aguas abajo.
- Gestiona la paginación y los límites de velocidad. Recorra las páginas y añada reintentos con retroceso progresivo para que un fallo transitorio no haga fallar toda la actualización.