Ingestion des données d’une API dans des pipelines
L'ingestion à partir d'une API signifie extraire des données via HTTP à partir d'un service web, généralement sous forme de JSON paginé, plutôt que de lire à partir d'un fichier ou d'une base de données. Contrairement aux fichiers ou à un bus de messages, il n'existe pas de source d'API générique intégrée ; vous devez donc gérer vous-même l'authentification, la pagination et les limites de débit. Les LakeFlow Pipelines prennent en charge trois modèles pour l'ingestion à partir d'une API arbitraire. Celui qui convient dépend de votre volume et de vos besoins en matière de refresh.
Avant d’écrire un code d’ingestion d’API personnalisé, vérifiez si un connecteur géré existe déjà pour votre source. Lakeflow Connect propose des connecteurs intégré pour de nombreuses APIs de logiciels en tant que service (SaaS) courants, tels que Salesforce, Workday, ServiceNow et Google analytique, ainsi qu’un ensemble croissant de connecteurs Partenaires. Si un connecteur couvre votre source, il gère l’authentification, la pagination et l’extraction incrémentielle pour vous, et cela représente presque toujours moins de travail qu’une ingestion développée en interne. Voir Connecteurs gérés dans Lakeflow Connect. Utilisez les modèles ci-dessous uniquement lorsqu’aucun connecteur ne convient.
Prérequis
- Un pipeline. Pour en créer un, consultez les tutoriels sur les LakeFlow Pipelines.
- Identifiants API, tels qu’un jeton ou une clé, stockés en tant que secret Databricks. Ne codez jamais d’identifiants en dur dans le code source du pipeline. Voir Gestion des secrets.
- Accès réseau depuis votre compute de pipeline vers l'Endpoint de l'API.
- Familiarité avec les tables de streaming et les vues matérialisées, les types de dataset produits par ces motifs. Voir Tables en streaming et Vues matérialisées.
Choisir un modèle
Il n'existe pas de source REST-API générique native dans les pipelines ; par conséquent, lorsque vous effectuez une extraction à partir d'une API arbitraire, choisissez l'un des trois modèles en fonction du volume de données et de la fréquence d'ingestion :
Modèle | À utiliser lorsque |
|---|---|
Les charges utiles sont de taille petite à moyenne et sont extraites une fois par exécution de pipeline, comme les données de référence, les taux de change quotidiens ou une API paginée mais limitable. | |
Vous devez interroger une API à haut volume ou en streaming de manière incrémentielle, avec une progression vérifiée par des points de contrôle afin qu'un redémarrage ne relise pas tout. | |
Vous souhaitez isoler les particularités spécifiques à l'API de votre logique de Transformations et bénéficier gratuitement d'un suivi de fichier « exactly-once ». |
Modèle 1 : extractions périodiques sous forme de vue matérialisée
Pour les charges utiles de petite à moyenne taille extraites une fois par exécution de pipeline, écrivez une fonction Python qui appelle l'API et renvoie un DataFrame Spark. Comme le dataset est une vue matérialisée, le pipeline réexécute la fonction entièrement et de manière idempotente à chaque mise à jour du pipeline.
Les étapes suivantes vous montrent comment créer une vue matérialisée avec des extractions périodiques :
-
Stockez le jeton API dans un secret, puis mappez-le à une propriété de configuration Spark dans les paramètres de votre pipeline afin que le code du pipeline puisse le lire. Ajoutez la propriété au bloc
spark_confde la configuration du cluster du pipeline :JSON{
"clusters": [
{
"spark_conf": {
"api.token": "{{secrets/<scope-name>/<secret-name>}}"
}
}
]
}Le code de l’étape suivante lit cette valeur avec
spark.conf.get("api.token"). Pour en savoir plus sur la configuration des secrets dans les paramètres du pipeline, consultez Accéder en toute sécurité aux identifiants de stockage avec des secrets dans un pipeline. -
Définissez une vue matérialisée qui appelle l'API et renvoie la réponse sous forme de DataFrame :
Pythonimport 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) -
Gérez la pagination à l’intérieur de la fonction en bouclant les pages et en concaténant les résultats avant de retourner le DataFrame :
Pythonimport 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)Ajoutez une logique de nouvelle tentative et d'interruption (backoff) autour de la requête pour plus de résilience.
Ce modèle relit l'intégralité de la réponse de l'API à chaque mise à jour du pipeline ; utilisez-le donc uniquement lorsque la charge utile est limitée. Pour les lectures incrémentielles, utilisez le modèle 2.
Modèle 2 : API à haut volume ou de streaming avec l'API Python Data Source
Pour les APIs nécessitant une interrogation incrémentale avec suivi du décalage, implémentez une source de données personnalisée à l'aide de l'API de source de données Python de Spark. Cela vous donne une sémantique de streaming correcte, incluant la progression par checkpoint et des lectures incrémentales, afin qu’un redémarrage reprenne à partir du dernier décalage au lieu de reprendre toute l’API.
Les étapes suivantes vous montrent comment ingérer des données depuis une source de données personnalisée :
-
Implémentez un
DataSourceet unDataSourceStreamReaderqui appellent l’API et suivent le décalage de lecture. Pour plus de détails sur la création d’une source de données personnalisée, voir PySpark sources de données personnalisées. -
Enregistrez la source de données afin que le pipeline puisse y faire référence par son nom de format :
Pythonspark.dataSource.register(MyApiDataSource) -
Lire depuis la source enregistrée dans une table de streaming :
Pythonfrom pyspark import pipelines as dp
@dp.table(name="events_bronze")
def events_bronze():
return spark.readStream.format("my_api_source").load()
Modèle 3 : découpler l'ingestion avec un Job planifié et Auto Loader
Un modèle de production courant consiste à séparer l'appel d'API du pipeline. Un Job planifié dépose les réponses brutes de l'API sous forme de fichiers dans un volume Unity Catalog, et le pipeline les récupère avec Auto Loader. Cela isole les particularités spécifiques à l'API, telles que la pagination et les limites de débit, de votre logique de transformation déclarative, et vous offre gratuitement le suivi de fichiers « exactly-once » d'Auto Loader.
Les étapes suivantes vous montrent comment dissocier l'ingestion avec un job planifié :
-
Écrivez un notebook ou un script qui appelle l’API et écrit les réponses JSON brutes dans un volume Unity Catalog. Lisez les identifiants de l’API à partir d’un secret. Voir Gestion des secrets.
Pythonimport 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) -
Planifiez l’exécution du notebook ou du script de manière autonome avec Lakeflow Jobs. See Lakeflow Jobs.
-
Dans votre pipeline, définissez une table de streaming qui lit les fichiers transférés avec Auto Loader :
Pythonfrom 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")
)
Pour en savoir plus sur l'ingestion fiable de fichiers avec Auto Loader, consultez Charger des fichiers à partir d'un stockage d'objets cloud et Qu'est-ce qu'Auto Loader ?.
Bonnes pratiques pour l'ingestion d'API
- Ne stockez pas de secrets dans le code source. Stockez les jetons et les clés d'API dans des Databricks Secret Scope et lisez-les au moment de l'exécution. Voir Gestion des secrets.
- Validez les réponses le plus tôt possible. Ajoutez des attentes sur les lignes ingérées afin de détecter les réponses d'API mal formées avant qu'elles ne soient traitées en aval.
- Gérez la pagination et les limites de débit. Parcourez les pages et ajoutez une nouvelle tentative avec interruption exponentielle afin qu'une défaillance temporaire ne fasse pas échouer toute la mise à jour.