You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何将Spark DataFrame写入ClickHouse?附代码及遇到的问题

How to Write Spark DataFrame to ClickHouse

Got it, let's get your Spark DataFrame into ClickHouse smoothly. Here's a step-by-step solution that covers the most common approaches:

Prerequisite: Add Required Dependencies

First, make sure your Spark environment has the ClickHouse connectors. You can include them when submitting your Spark job:

spark-submit --packages ru.yandex.clickhouse:clickhouse-jdbc:0.3.2,com.github.housepower:clickhouse-spark-runtime-3.2_2.12:2.6.5 your-spark-app.jar

Note: Adjust the version numbers to match your Spark version (e.g., use clickhouse-spark-runtime-3.3_2.12 for Spark 3.3).

This dedicated connector offers better performance and native ClickHouse support. Build on your existing code like this:

import org.apache.spark.sql.SaveMode

// Your existing variables (fixed the typo in username for you)
val df = spark.read.parquet(path)
val IP ="190.176.35.145"
val port = "9000"
val table = "table1"
val user = "default"
val password = "default"
val database = "your_target_database" // Replace with your actual database name

// Write DataFrame to ClickHouse
df.write
  .format("clickhouse")
  .option("host", IP)
  .option("port", port)
  .option("user", user)
  .option("password", password)
  .option("database", database)
  .option("table", table)
  .mode(SaveMode.Append) // Choose the right mode based on your need
  .save()

Key Notes for This Method:

  • Table Structure Alignment: Ensure your DataFrame's column names and data types exactly match the ClickHouse table. Mismatches will cause write failures.
  • Write Modes:
    • SaveMode.Append: Add new data to the existing table
    • SaveMode.Overwrite: Replace the entire table with your DataFrame
    • SaveMode.Ignore: Do nothing if the table already exists
    • SaveMode.ErrorIfExists: Throw an error if the table exists (default behavior)
  • Performance Tuning: Add options like option("batch_size", "10000") to adjust batch size, or option("socket_timeout", "300000") to extend timeout for large datasets.

Method 2: Use JDBC Connection (Fallback Option)

If the dedicated connector isn't working for your setup, you can use the standard JDBC approach:

import org.apache.spark.sql.SaveMode

// Your existing variables (same as above)

df.write
  .format("jdbc")
  .option("url", s"jdbc:clickhouse://$IP:$port/$database")
  .option("dbtable", table)
  .option("user", user)
  .option("password", password)
  .option("driver", "ru.yandex.clickhouse.ClickHouseDriver")
  .mode(SaveMode.Append)
  .save()

When to Use JDBC:

  • You need a more universal approach that works with older Spark versions
  • You're facing compatibility issues with the dedicated connector

内容的提问来源于stack exchange,提问作者sparkFish

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.06 20:09:05