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 :
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.
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.