Structured Streaming écrit dans Azure Synapse
Cette documentation a été retirée et pourrait ne pas être mise à jour.
Le connecteur Azure Synapse offre un support d'écriture Structured Streaming efficace et évolutif pour Azure Synapse qui offre une expérience utilisateur cohérente avec les écritures par batch et utilise COPY pour les transferts de données volumineux entre un cluster Databricks et une instance Azure Synapse.
La prise en charge de Structured Streaming entre Databricks et Synapse fournit une sémantique simple pour la configuration des Jobs ETL incrémentiels. Le modèle utilisé pour charger des données de Databricks vers Synapse introduit une latence qui pourrait ne pas satisfaire aux exigences SLA pour les charges de travail en quasi-temps réel. Voir query data in Azure Synapse Analytics.
Modes de sortie pris en charge pour les écritures en streaming vers Synapse
Le connecteur Azure Synapse prend en charge les modes de sortie Append et Complete pour les ajouts d'enregistrements et les agrégations. Pour plus de détails sur les modes de sortie et la matrice de compatibilité, consultez le guide Structured Streaming.
Sémantique de tolérance aux pannes Synapse
Par default, Azure Synapse Streaming offre une garantie de traitement *exactement une fois* de bout en bout pour l'écriture de données dans une table Azure Synapse, en suivant de manière fiable la progression de la query à l'aide d'une combinaison d'emplacement de point de contrôle dans DBFS, de table de point de contrôle dans Azure Synapse et de mécanisme de verrouillage pour garantir que le streaming peut gérer tout type de défaillance, de nouvelle tentative et de redémarrage de query.
En option, vous pouvez sélectionner une sémantique au moins une fois moins restrictive pour Azure Synapse Streaming en définissant l'option spark.databricks.sqldw.streaming.exactlyOnce.enabled sur false, auquel cas une duplication des données pourrait se produire en cas d'échecs de connexion intermittents vers Azure Synapse ou d'arrêt inattendu de la query.
Syntaxe Structured Streaming pour l'écriture dans Azure Synapse
Les exemples de code suivants illustrent les écritures en streaming vers Synapse à l'aide de Structured Streaming en Scala et Python :
- Scala
- Python
// Set up the Blob storage account access key in the notebook session conf.
spark.conf.set(
"fs.azure.account.key.<your-storage-account-name>.dfs.core.windows.net",
"<your-storage-account-access-key>")
// Prepare streaming source; this could be Kafka or a simple rate stream.
val df: DataFrame = spark.readStream
.format("rate")
.option("rowsPerSecond", "100000")
.option("numPartitions", "16")
.load()
// Apply some transformations to the data then use
// Structured Streaming API to continuously write the data to a table in Azure Synapse.
df.writeStream
.format("com.databricks.spark.sqldw")
.option("url", "jdbc:sqlserver://<the-rest-of-the-connection-string>")
.option("tempDir", "abfss://<your-container-name>@<your-storage-account-name>.dfs.core.windows.net/<your-directory-name>")
.option("forwardSparkAzureStorageCredentials", "true")
.option("dbTable", "<your-table-name>")
.option("checkpointLocation", "/tmp_checkpoint_location")
.start()
# Set up the Blob storage account access key in the notebook session conf.
spark.conf.set(
"fs.azure.account.key.<your-storage-account-name>.dfs.core.windows.net",
"<your-storage-account-access-key>")
# Prepare streaming source; this could be Kafka or a simple rate stream.
df = spark.readStream \
.format("rate") \
.option("rowsPerSecond", "100000") \
.option("numPartitions", "16") \
.load()
# Apply some transformations to the data then use
# Structured Streaming API to continuously write the data to a table in Azure Synapse.
df.writeStream \
.format("com.databricks.spark.sqldw") \
.option("url", "jdbc:sqlserver://<the-rest-of-the-connection-string>") \
.option("tempDir", "abfss://<your-container-name>@<your-storage-account-name>.dfs.core.windows.net/<your-directory-name>") \
.option("forwardSparkAzureStorageCredentials", "true") \
.option("dbTable", "<your-table-name>") \
.option("checkpointLocation", "/tmp_checkpoint_location") \
.start()
Pour la liste complète des configurations, consultez Query data in Azure Synapse Analytics.
Gestion des tables de point de contrôle de Synapse streaming
Le connecteur Azure Synapse ne supprime pas la table de points de contrôle de streaming qui est créée lorsqu'une nouvelle query de streaming est start. Ce comportement est conforme à checkpointLocation normalement spécifié pour le stockage d'objets. Databricks vous recommande de supprimer périodiquement les tables de points de contrôle pour les requêtes qui ne seront pas exécutées à l'avenir.
Par default, toutes les tables de point de contrôle ont le nom <prefix>_<query-id>, où <prefix> est un préfixe configurable avec la valeur default databricks_streaming_checkpoint et query_id est un ID de query en streaming avec _ caractères supprimés.
Pour trouver toutes les tables de point de contrôle pour les queries de streaming obsolètes ou supprimées, exécutez la query :
SELECT * FROM sys.tables WHERE name LIKE 'databricks_streaming_checkpoint%'
Vous pouvez configurer le préfixe avec l'option de configuration Spark SQL spark.databricks.sqldw.streaming.exactlyOnce.checkpointTableNamePrefix.
Référence des options de streaming du connecteur Databricks Synapse
Les OPTIONS fournis dans Spark SQL prennent en charge les options suivantes pour le streaming en plus des options batch:
parameter | Obligatoire | Par défaut | Notes |
|---|---|---|---|
| Oui | No default | Emplacement sur DBFS qui sera utilisé par Structured Streaming pour écrire les métadonnées et les informations de point de contrôle. Consultez le guide de programmation de Structured Streaming sur la récupération après des pannes avec des points de contrôle. |
| Non | 0 | Indique combien de répertoires temporaires (les plus récents) conserver pour le nettoyage périodique des micro-batch en streaming. Lorsqu'il est défini sur |
checkpointLocation et numStreamingTempDirsToKeep ne sont pertinentes que pour les écritures en streaming de Databricks vers une nouvelle table dans Azure Synapse.