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

如何通过Spark JDBC将数组列写入ClickHouse?

解决Spark写入ClickHouse Array类型字段的问题

问题重现

你创建的Spark DataFrame:

val df = Seq(
    List("aa","bb","cc"),
    List("aa","bb","cc"),
    List("aa","bb","cc")
).toDF("str")

ClickHouse表结构:

CREATE TABLE IF NOT EXISTS testtable (
    str Array(String)
) ENGINE = Memory

使用Spark JDBC写入时触发错误:

java.lang.IllegalArgumentException: Can't get JDBC type for array

尝试用concat_ws转换数组为字符串后,仍报错且数据格式不符合预期:

java.sql.SQLException: Unsupported type ARRAY
写入后数据变成["aa,bb,cc"]而非预期的["aa", "bb", "cc"]。

解决方案

方法1:使用ClickHouse官方Spark Connector(推荐)

官方连接器对ClickHouse的类型支持更完善,无需手动处理类型转换。

  1. 添加Maven依赖(Scala 2.12为例):
<dependency>
    <groupId>com.clickhouse</groupId>
    <artifactId>clickhouse-spark-runtime-3.3_2.12</artifactId>
    <version>0.5.0</version>
</dependency>
  1. 写入代码示例:
df.write
  .format("clickhouse")
  .option("host", "your-clickhouse-host")
  .option("port", "8123")
  .option("user", "user")
  .option("password", "password")
  .option("database", "your-db")
  .option("table", "testtable")
  .mode("append")
  .save()

方法2:使用JDBC并手动处理数组类型

如果必须使用JDBC方式,可以通过foreachPartition手动创建PreparedStatement,利用ClickHouse JDBC的setArray方法处理数组:

import java.sql.{Connection, DriverManager, PreparedStatement}
import org.apache.spark.sql.Row

val url = "jdbc:clickhouse://your-host:8123/your-db"
val user = "user"
val password = "password"

df.foreachPartition { rows =>
  var conn: Connection = null
  var stmt: PreparedStatement = null
  try {
    Class.forName("com.clickhouse.jdbc.ClickHouseDriver")
    conn = DriverManager.getConnection(url, user, password)
    val sql = "INSERT INTO testtable (str) VALUES (?)"
    stmt = conn.prepareStatement(sql)
    
    rows.foreach { row =>
      val arr = row.getAs[Seq[String]](0).toArray
      stmt.setArray(1, conn.createArrayOf("String", arr))
      stmt.addBatch()
    }
    stmt.executeBatch()
  } catch {
    case e: Exception => e.printStackTrace()
  } finally {
    if (stmt != null) stmt.close()
    if (conn != null) conn.close()
  }
}

方法3:将数组序列化为ClickHouse兼容的字符串格式

ClickHouse支持用[]包裹、元素用单引号包裹并以逗号分隔的字符串表示数组(需确保元素不含单引号/逗号,或提前转义),可以先将数组转换为这种格式再写入:

import org.apache.spark.sql.functions.{concat, lit, array_join}

// 将数组转为ClickHouse兼容的字符串格式:"['aa','bb','cc']"
val formattedDf = df.withColumn(
  "str",
  concat(lit("["), array_join($"str", "','", "'"), lit("]"))
)

formattedDf.write
  .format("jdbc")
  .mode("append")
  .option("driver", "com.clickhouse.jdbc.ClickHouseDriver")
  .option("url", "jdbcurl")
  .option("user", "user")
  .option("password", "password")
  .option("dbtable", "testtable")
  .save()

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.28 03:13:13