Remarque
L’accès à cette page nécessite une autorisation. Vous pouvez essayer de vous connecter ou de modifier des répertoires.
L’accès à cette page nécessite une autorisation. Vous pouvez essayer de modifier des répertoires.
Ingérer depuis une API signifie extraire des données via HTTP depuis 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 au bus de messages, il n’y a pas de source API générique intégrée, donc vous gérez vous-même l’authentification, la pagination et les limites de débit. Les pipelines Lakeflow prennent en charge trois modèles pour ingérer des données depuis une API quelconque. Lequel convient dépend de votre volume et de vos besoins de rafraîchissement.
Important
Avant d’écrire un code d’ingestion API personnalisé, vérifiez si un connecteur géré existe déjà pour votre source. Lakeflow Connect fournit des connecteurs intégrés pour de nombreuses API SaaS courantes, telles que Salesforce, Workday, ServiceNow et Google Analytics, et propose également un nombre croissant de connecteurs partenaires. Si un connecteur couvre votre source, il assure pour vous l’authentification, la pagination et l’extraction incrémentielle, ce qui représente presque toujours moins de travail qu’une ingestion manuelle. Consultez Connecteurs gérés dans Lakeflow Connect. Utilisez les patrons ci-dessous uniquement quand aucun connecteur ne s’adapte.
Prerequisites
- Un pipeline. Pour en créer un, consultez les tutoriels Lakeflow pipelines.
- Les identifiants API, tels qu’un jeton ou une clé, sont stockés comme un secret Azure Databricks. Ne codez jamais les identifiants en dur dans le code source du pipeline. Consultez Gestion des secrets.
- Accès réseau depuis les ressources de calcul de votre pipeline vers le point de terminaison de l’API.
- Connaissance des tables en continu et des vues matérialisées, ainsi que des types de jeux de données que ces modèles produisent. 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, donc lorsque vous puisez dans une API arbitraire, choisissez l’un des trois modèles selon le volume de données et la fréquence d’ingestion :
| Modèle | À utiliser quand |
|---|---|
| Extractions périodiques en tant que vue matérialisée | Les charges utiles sont faibles à moyennes 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 dont le volume peut être limité. |
| API Python Data Source | Vous devez interroger de manière incrémentale une API à fort volume ou en streaming, avec une progression jalonnée de points de contrôle afin qu’un redémarrage ne relise pas l’ensemble des données. |
| Ingestion découplée avec Auto Loader | Vous voulez isoler les particularités spécifiques à l’API de votre logique de transformation et obtenir un suivi de fichiers exactement une fois gratuitement. |
Modèle 1 : Extractions périodiques en tant que vue matérialisée
Pour les charges utiles de petite à moyenne taille tirées une fois par exécution de pipeline, écrire une fonction Python qui appelle l’API et retourne un DataFrame Spark. Comme le jeu de données est une vue matérialisée, le pipeline réexécute la fonction entièrement et de façon idempotente à chaque mise à jour du pipeline.
Les étapes suivantes vous montrent comment construire une vue matérialisée avec des tirages périodiques :
Stockez le jeton API en secret, puis associez-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 de cluster du pipeline :{ "clusters": [ { "spark_conf": { "api.token": "{{secrets/<scope-name>/<secret-name>}}" } } ] }Le code à 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, voir Accès sécurisé aux identifiants de stockage avec des secrets dans un pipeline.Définissons une vue matérialisée qui appelle l’API et renvoie la réponse sous forme de 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)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 :
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)Ajoutez une logique de réessaie et de recul autour de la demande de résilience.
Ce pattern relit la réponse complète de l’API à chaque mise à jour du pipeline, donc utilisez-le uniquement lorsque la charge utile est bornée. Pour les lectures incrémentales, utilisez le motif 2.
Modèle 2 : API à haut volume ou en streaming avec l’API Python Data Source
Pour les API, il faut faire des sondages incrémentaux avec le suivi décalé, implémenter une source de données personnalisée avec l'API Python Data Source de Spark. Cela vous donne une sémantique de streaming appropriée, notamment une progression avec des points de contrôle et des lectures incrémentielles, si bien qu’un redémarrage reprend au dernier décalage au lieu de réinterroger toute l’API.
Les étapes suivantes vous montrent comment ingérer à partir d’une source de données personnalisée :
Implémentez un
DataSourceetDataSourceStreamReaderqui appellent l’API et suivez 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 la référencer par le nom du format :
spark.dataSource.register(MyApiDataSource)Lisez à partir de la source enregistrée dans une table de streaming :
from 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 à l’aide d’un travail planifié et du chargeur automatique
Un schéma de production courant consiste à séparer l’appel API du pipeline. Un job planifié envoie les réponses API brutes sous forme de fichiers dans un volume Unity Catalog, et le pipeline les récupère avec Auto Loader. Cela isole de la logique de transformation déclarative les particularités propres à l’API, telles que la pagination et les limites de débit, et vous permet de bénéficier gratuitement du suivi de fichiers exactement une fois du chargeur automatique.
Les étapes suivantes vous montrent comment découpler l’ingestion d’un travail planifié :
Écrivez un notebook ou un script qui appelle l’API et écrit les réponses JSON brutes d’un volume du catalogue Unity. Lisez les informations d’identification de l’API à partir d’un secret. Consultez Gestion des secrets.
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)Programmez le notebook ou le script pour qu’il s’exécute seul avec Lakeflow Jobs. Consultez les offres d'emploi Lakeflow.
Dans votre pipeline, définissez une table de streaming qui lit les fichiers atterris avec 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") )
Pour en savoir plus sur l’ingestion fiable de fichiers avec Auto Loader, voir Charger des fichiers depuis le stockage d’objets cloud et Qu’est-ce qu’Auto Loader ?.
Meilleures pratiques pour l’ingestion d’API
- Garde les secrets hors du code source. Stockez les jetons et clés API dans les scopes secrets Azure Databricks et lisez-les à l’exécution. Consultez Gestion des secrets.
- Validez les réponses dès le début. Ajoutez des règles de validation sur les lignes ingérées pour détecter les réponses d’API mal formées avant qu’elles ne soient transmises aux étapes suivantes.
- Gérer la pagination et les limites de fréquence. Parcourez les pages et ajoutez un mécanisme de nouvelle tentative avec temporisation exponentielle afin qu’une erreur transitoire ne fasse pas échouer toute la mise à jour.