Aller au contenu principal

Modèles Structured Streaming sur Databricks

Ceci contient des Notebooks et des exemples de code pour les modèles courants de travail avec Structured Streaming sur Databricks.

Prise en main de Structured Streaming

Si vous débutez avec Structured Streaming, consultez Exécuter votre première charge de travail Structured Streaming.

Écrire dans Cassandra en tant que récepteur pour Structured Streaming en Python

Apache Cassandra est une base de données OLTP distribuée, à faible latence, évolutive et hautement disponible.

Structured Streaming fonctionne avec Cassandra via le connecteur Spark Cassandra. Ce connecteur prend en charge les API RDD et DataFrame, et il offre un support natif pour l'écriture de données en streaming. Important Vous devez utiliser la version correspondante du spark-cassandra-connector-assembly.

L'exemple suivant se connecte à un ou plusieurs hôtes d'un cluster de base de données Cassandra. Il spécifie également les configurations de connexion telles que l'emplacement du point de contrôle et les noms spécifiques des keyspace et des tables :

Python
spark.conf.set("spark.cassandra.connection.host", "host1,host2")

df.writeStream \
.format("org.apache.spark.sql.cassandra") \
.outputMode("append") \
.option("checkpointLocation", "/path/to/checkpoint") \
.option("keyspace", "keyspace_name") \
.option("table", "table_name") \
.start()

Écrire dans Azure Synapse Analytics avec foreachBatch() en Python

streamingDF.writeStream.foreachBatch() vous permet de réutiliser les enregistreurs de données batch existants pour écrire la sortie d'une query streaming dans Azure Synapse Analytics. Consultez la documentation foreachBatch pour plus de détails.

Pour exécuter cet exemple, vous avez besoin du connecteur Azure Synapse Analytics. Pour plus de détails sur le connecteur Azure Synapse Analytics, consultez Interroger les données dans Azure Synapse Analytics.

Python
from pyspark.sql.functions import *
from pyspark.sql import *

def writeToSQLWarehouse(df, epochId):
df.write \
.format("com.databricks.spark.sqldw") \
.mode('overwrite') \
.option("url", "jdbc:sqlserver://<the-rest-of-the-connection-string>") \
.option("forward_spark_azure_storage_credentials", "true") \
.option("dbtable", "my_table_in_dw_copy") \
.option("tempdir", "wasbs://<your-container-name>@<your-storage-account-name>.blob.core.windows.net/<your-directory-name>") \
.save()

spark.conf.set("spark.sql.shuffle.partitions", "1")

query = (
spark.readStream.format("rate").load()
.selectExpr("value % 10 as key")
.groupBy("key")
.count()
.toDF("key", "count")
.writeStream
.foreachBatch(writeToSQLWarehouse)
.outputMode("update")
.start()
)

Jointures Stream-Stream

Ces deux Notebooks montrent comment utiliser les jointures de Stream à Stream en Python et Scala.

Jointures Stream-Stream Notebook Python

Jointures Stream-Stream du Notebook Scala