Utilisez un programme Spark pour importer par lot des données CSV dans ApsaraDB for ClickHouse via une connexion JDBC. Cette approche convient lorsque vous disposez déjà d'un pipeline Spark et souhaitez écrire directement des DataFrames dans ClickHouse à l'aide du pilote JDBC ClickHouse.
Si vous avez besoin d'une prise en charge complète des types complexes (MAP, ARRAY, STRUCT) ou si vous préférez une intégration plus native, utilisez le connecteur natif Spark-ClickHouse . L'approche JDBC décrite ici ne prend pas en charge ces types.
Prérequis
Avant de commencer, assurez-vous d'avoir :
Ajouté l'adresse IP de votre machine locale à la liste d'autorisation du cluster ApsaraDB for ClickHouse. Consultez la section Configurer une liste d'autorisation.
Créé une table ApsaraDB for ClickHouse dont les types de données des colonnes correspondent aux données que vous souhaitez importer. Consultez la section Créer une table.
Fonctionnement
Le programme Spark lit un fichier CSV dans un DataFrame et l'écrit dans une table ApsaraDB for ClickHouse à l'aide du pilote JDBC ClickHouse. Les lignes sont insérées par lots via une connexion JDBC sur le port 8123.
Paramètres de connexion clés :
| Paramètre | Valeur | Description |
|---|---|---|
batchsize |
100000 |
Nombre de lignes par insertion par lot |
socket_timeout |
300000 |
Délai d'expiration du socket en millisecondes |
numPartitions |
8 |
Nombre de partitions d'écriture parallèles. Augmentez cette valeur pour les jeux de données plus volumineux ; réduisez-la pour diminuer la charge sur le cluster. |
rewriteBatchedStatements |
true |
Réécrit les insertions par lots en une seule instruction multi-lignes pour améliorer le débit |
Importer des données depuis un fichier CSV
Étape 1 : Configurer la structure du projet
Créez la structure de répertoires suivante pour votre projet Spark :
find .
.
./build.sbt
./src
./src/main
./src/main/scala
./src/main/scala/com
./src/main/scala/com/spark
./src/main/scala/com/spark/test
./src/main/scala/com/spark/test/WriteToCk.scala
Étape 2 : Ajouter les dépendances
Ajoutez le contenu suivant au fichier build.sbt :
name := "Simple Project"
version := "1.0"
scalaVersion := "2.12.10"
libraryDependencies += "org.apache.spark" %% "spark-sql" % "3.0.0"
libraryDependencies += "ru.yandex.clickhouse" % "clickhouse-jdbc" % "0.2.4"
Étape 3 : Écrire le programme Spark
Créez le fichier src/main/scala/com/spark/test/WriteToCk.scala avec le contenu suivant. Remplacez les espaces réservés indiqués dans le tableau ci-dessous avant l'exécution.
package com.spark.test
import java.util
import java.util.Properties
import org.apache.spark.sql.execution.datasources.jdbc.JDBCOptions
import org.apache.spark.SparkConf
import org.apache.spark.sql.{SaveMode, SparkSession}
import org.apache.spark.storage.StorageLevel
object WriteToCk {
val properties = new Properties()
properties.put("driver", "ru.yandex.clickhouse.ClickHouseDriver")
properties.put("user", "<yourUserName>")
properties.put("password", "<yourPassword>")
properties.put("batchsize","100000")
properties.put("socket_timeout","300000")
properties.put("numPartitions","8")
properties.put("rewriteBatchedStatements","true")
val url = "jdbc:clickhouse://<yourUrl>:8123/default"
val table = "<yourTableName>"
def main(args: Array[String]): Unit = {
val sc = new SparkConf()
sc.set("spark.driver.memory", "1G")
sc.set("spark.driver.cores", "4")
sc.set("spark.executor.memory", "1G")
sc.set("spark.executor.cores", "2")
val session = SparkSession.builder().master("local[*]").config(sc).appName("write-to-ck").getOrCreate()
val df = session.read.format("csv")
.option("header", "true")
.option("sep", ",")
.option("inferSchema", "true")
.load("<yourFilePath>")
.selectExpr(
"colName1",
"colName2",
"colName3",
...
)
.persist(StorageLevel.MEMORY_ONLY_SER_2)
println(s"read done")
df.write.mode(SaveMode.Append).option(JDBCOptions.JDBC_BATCH_INSERT_SIZE, 100000).jdbc(url, table, properties)
println(s"write done")
df.unpersist(true)
}
}
Remplacez les espaces réservés suivants par vos valeurs réelles :
| Espace réservé | Description | Obligatoire |
|---|---|---|
<yourUserName> |
Nom d'utilisateur du compte de base de données dans ApsaraDB for ClickHouse | Oui |
<yourPassword> |
Mot de passe du compte de base de données | Oui |
<yourUrl> |
Endpoint du cluster ApsaraDB for ClickHouse | Oui |
<yourTableName> |
Nom de la table de destination dans ApsaraDB for ClickHouse | Oui |
<yourFilePath> |
Chemin d'accès au fichier CSV à importer, y compris le nom du fichier | Oui |
colName1,colName2,colName3 |
Noms des colonnes de la table ApsaraDB for ClickHouse à sélectionner depuis le DataFrame | Oui |
Étape 4 : Générer le package
Exécutez la commande suivante pour compiler et empaqueter le programme :
sbt package
Le fichier JAR de sortie est généré à l'emplacement target/scala-2.12/simple-project_2.12-1.0.jar.
Étape 5 : Soumettre la tâche Spark
Exécutez la commande suivante pour soumettre la tâche. Cette commande ajoute le pilote JDBC ClickHouse au classpath pour les processus driver et executor.
${SPARK_HOME}/bin/spark-submit --class "com.spark.test.WriteToCk" --master local[4] --conf "spark.driver.extraClassPath=${HOME}/.m2/repository/ru/yandex/clickhouse/clickhouse-jdbc/0.2.4/clickhouse-jdbc-0.2.4.jar" --conf "spark.executor.extraClassPath=${HOME}/.m2/repository/ru/yandex/clickhouse/clickhouse-jdbc/0.2.4/clickhouse-jdbc-0.2.4.jar" target/scala-2.12/simple-project_2.12-1.0.jar
Limites
Absence de prise en charge des types de données complexes : Le pilote JDBC ClickHouse ne prend pas en charge les types complexes Spark tels que MAP, ARRAY ou STRUCT. Pour utiliser ces types, basculez vers le connecteur natif Spark-ClickHouse.
Les tables doivent exister avant l'importation : JDBC ne peut pas créer automatiquement la table de destination. Créez la table dans ApsaraDB for ClickHouse avant d'exécuter la tâche Spark. Consultez la section Créer une table.