Use um programa Spark para importar dados CSV em lote para o ApsaraDB for ClickHouse por meio de uma conexão JDBC. Essa abordagem é ideal quando você já tem um pipeline Spark e quer gravar DataFrames diretamente no ClickHouse com o driver JDBC do ClickHouse.
Se precisar de suporte completo a tipos complexos (MAP, ARRAY, STRUCT) ou preferir uma integração mais nativa, use o conector nativo Spark-ClickHouse . O método JDBC descrito aqui não oferece suporte a esses tipos.
Pré-requisitos
Antes de começar, verifique se você:
Adicionou o endereço IP da máquina local à lista de permissões do cluster ApsaraDB for ClickHouse. Consulte Configure uma lista de permissões.
Crie uma tabela no ApsaraDB for ClickHouse com tipos de dados de coluna correspondentes aos dados a importar. Consulte Crie uma tabela.
Como funciona
O programa Spark lê um arquivo CSV em um DataFrame e o grava em uma tabela do ApsaraDB for ClickHouse com o driver JDBC do ClickHouse. As linhas são inseridas em lotes via conexão JDBC na porta 8123.
Principais parâmetros de conexão:
|
Parâmetro |
Valor |
Descrição |
|
|
|
Quantidade de linhas por inserção em lote |
|
|
|
Tempo limite do socket em milissegundos |
|
|
|
Número de partições de gravação paralelas. Aumente para conjuntos de dados maiores; diminua para reduzir a carga no cluster. |
|
|
|
Reescreve inserções em lote como uma única instrução de múltiplas linhas para melhorar o throughput |
Importar dados de CSV
Etapa 1: Configure a estrutura do projeto
Crie a seguinte estrutura de diretórios para o projeto 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
Etapa 2: Adicionar dependências
Adicione o conteúdo abaixo ao arquivo 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"
Etapa 3: Escrever o programa Spark
Crie o arquivo src/main/scala/com/spark/test/WriteToCk.scala com o conteúdo a seguir. Substitua os placeholders listados na tabela abaixo antes da execução.
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)
}
}
Substitua os seguintes placeholders pelos valores reais:
|
Placeholder |
Descrição |
Obrigatório |
|
|
Nome de usuário da conta do banco de dados no ApsaraDB for ClickHouse |
Sim |
|
|
Senha da conta do banco de dados |
Sim |
|
|
Endpoint do cluster ApsaraDB for ClickHouse |
Sim |
|
|
Nome da tabela de destino no ApsaraDB for ClickHouse |
Sim |
|
|
Caminho do arquivo CSV a importar, incluindo o nome do arquivo |
Sim |
|
|
Nomes das colunas na tabela do ApsaraDB for ClickHouse a selecione do DataFrame |
Sim |
Etapa 4: Compilar o pacote
Execute o comando a seguir para compilar e empacotar o programa:
sbt package
O JAR de saída será gerado em target/scala-2.12/simple-project_2.12-1.0.jar.
Etapa 5: Envie o job Spark
Execute o comando abaixo para enviar o job. Esse comando adiciona o driver JDBC do ClickHouse ao classpath dos processos do driver e do 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
Limitações
Sem suporte para tipos de dados complexos: O driver JDBC do ClickHouse não suporta tipos complexos do Spark, como MAP, ARRAY ou STRUCT. Para usar esses tipos, mude para o conector nativo Spark-ClickHouse.
As tabelas devem existir antes da importação: O JDBC não cria automaticamente a tabela de destino. Crie a tabela no ApsaraDB for ClickHouse antes de executar o job Spark. Consulte Crie uma tabela.