Exemples de code pour Databricks Connect pour Scala
Cet article couvre Databricks Connect pour Databricks Runtime 13.3 LTS et les versions ultérieures.
Cet article fournit des exemples de code qui utilisent Databricks Connect pour Scala. Databricks Connect vous permet de connecter des IDEs, des serveurs Notebook et des applications personnalisées aux clusters Databricks. See Databricks Connect. Pour la version Python de cet article, consultez Exemples de code pour Databricks Connect pour Python.
Avant de commencer à utiliser Databricks Connect, vous devez configurer le client Databricks Connect.
Les exemples suivants supposent que vous utilisez l'authentification default pour la configuration du client Databricks Connect.
Exemple : Lire une table
Cet exemple de code simple interroge la table spécifiée, puis affiche les 5 premières lignes de la table spécifiée.
import com.databricks.connect.DatabricksSession
import org.apache.spark.sql.SparkSession
object Main {
def main(args: Array[String]): Unit = {
val spark = DatabricksSession.builder().getOrCreate()
val df = spark.read.table("samples.nyctaxi.trips")
df.limit(5).show()
}
}
Créer un DataFrame
L'exemple de code ci-dessous :
- Crée un DataFrame en mémoire.
- Crée une table portant le nom
zzz_demo_temps_tabledans le schémadefault. Si la table avec ce nom existe déjà, la table est d’abord supprimée. Pour utiliser un schéma ou une table différente, ajustez les appels àspark.sql,temps.write.saveAsTableou les deux. - Enregistre le contenu du DataFrame dans la table.
- Exécute une query
SELECTsur le contenu de la table. - Affiche le résultat de la query.
- Supprime la table.
import com.databricks.connect.DatabricksSession
import org.apache.spark.sql.SparkSession
import org.apache.spark.sql.types._
import java.time.LocalDate
object Main {
def main(args: Array[String]): Unit = {
val spark = DatabricksSession.builder().getOrCreate()
// Create a Spark DataFrame consisting of high and low temperatures
// by airport code and date.
val schema = StructType(
Seq(
StructField("AirportCode", StringType, false),
StructField("Date", DateType, false),
StructField("TempHighF", IntegerType, false),
StructField("TempLowF", IntegerType, false)
)
)
val data = Seq(
( "BLI", LocalDate.of(2021, 4, 3), 52, 43 ),
( "BLI", LocalDate.of(2021, 4, 2), 50, 38),
( "BLI", LocalDate.of(2021, 4, 1), 52, 41),
( "PDX", LocalDate.of(2021, 4, 3), 64, 45),
( "PDX", LocalDate.of(2021, 4, 2), 61, 41),
( "PDX", LocalDate.of(2021, 4, 1), 66, 39),
( "SEA", LocalDate.of(2021, 4, 3), 57, 43),
( "SEA", LocalDate.of(2021, 4, 2), 54, 39),
( "SEA", LocalDate.of(2021, 4, 1), 56, 41)
)
val temps = spark.createDataFrame(data).toDF(schema.fieldNames: _*)
// Create a table on the Databricks cluster and then fill
// the table with the DataFrame 's contents.
// If the table already exists from a previous run,
// delete it first.
spark.sql("USE default")
spark.sql("DROP TABLE IF EXISTS zzz_demo_temps_table")
temps.write.saveAsTable("zzz_demo_temps_table")
// Query the table on the Databricks cluster, returning rows
// where the airport code is not BLI and the date is later
// than 2021-04-01.Group the results and order by high
// temperature in descending order.
val df_temps = spark.sql("SELECT * FROM zzz_demo_temps_table " +
"WHERE AirportCode != 'BLI' AND Date > '2021-04-01' " +
"GROUP BY AirportCode, Date, TempHighF, TempLowF " +
"ORDER BY TempHighF DESC")
df_temps.show()
// Results:
// +------------+-----------+---------+--------+
// | AirportCode| Date|TempHighF|TempLowF|
// +------------+-----------+---------+--------+
// | PDX | 2021-04-03| 64 | 45 |
// | PDX | 2021-04-02| 61 | 41 |
// | SEA | 2021-04-03| 57 | 43 |
// | SEA | 2021-04-02| 54 | 39 |
// +------------+-----------+---------+--------+
// Clean up by deleting the table from the Databricks cluster.
spark.sql("DROP TABLE zzz_demo_temps_table")
}
}
Exemple : Utilisation de DatabricksSesssion ou de SparkSession
L'exemple suivant décrit comment utiliser la classe SparkSession dans les cas où la classe DatabricksSession dans Databricks Connect n'est pas disponible.
Cet exemple query la table spécifiée et renvoie les 5 premières lignes. Cet exemple utilise la variable d'environnement SPARK_REMOTE pour l'authentification.
import org.apache.spark.sql.{DataFrame, SparkSession}
object Main {
def main(args: Array[String]): Unit = {
getTaxis(getSpark()).show(5)
}
private def getSpark(): SparkSession = {
SparkSession.builder().getOrCreate()
}
private def getTaxis(spark: SparkSession): DataFrame = {
spark.read.table("samples.nyctaxi.trips")
}
}
Ressources supplémentaires
Databricks fournit des exemples d'applications supplémentaires qui montrent comment utiliser Databricks Connect dans le repository GitHub Databricks Connect, notamment :