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

Spark写入含MapType列的DataFrame到ClickHouse报错求助

问题描述

使用clickhouse-native-jdbc驱动将包含MapType列的Spark DataFrame写入对应Schema含Map类型列的ClickHouse时,抛出如下错误:

Caused by: java.lang.IllegalArgumentException: Can't translate non-null value for field 74
        at org.apache.spark.sql.execution.datasources.jdbc.JdbcUtils$.$anonfun$makeSetter$16(JdbcUtils.scala:593)
        at org.apache.spark.sql.execution.datasources.jdbc.JdbcUtils$.$anonfun$makeSetter$16$adapted(JdbcUtils.scala:591)

原因是Spark原生JDBC工具类JdbcUtils的makeSetter方法未适配MapType类型,无匹配逻辑时抛出该异常。

尝试修改org.apache.spark.sql.execution.datasources.jdbc.JdbcUtils添加MapType处理逻辑:

case MapType(_, _, _) =>
    (stmt: PreparedStatement, row: Row, pos: Int) =>
        val map = row.getMap[AnyRef, AnyRef](pos)
        stmt.setObject(pos + 1, mapAsJavaMap(map))

该修改在本地环境生效,但集群模式下执行器仍加载原生JdbcUtils,未使用修改后的代码。


解决方案

一、无需修改Spark源码的替代方案

1. 将MapType列序列化为JSON字符串写入

把Spark的Map列转换为JSON字符串,ClickHouse侧用String或JSON类型(ClickHouse 21.8+支持)接收,查询时再解析:

import org.apache.spark.sql.functions.to_json

// 将Map列转为JSON字符串
val dfWithJsonMap = originalDF.withColumn("map_col", to_json($"map_col"))

// 写入ClickHouse
dfWithJsonMap.write
  .format("jdbc")
  .option("url", "jdbc:clickhouse://host:port/db")
  .option("dbtable", "target_table")
  .option("driver", "com.clickhouse.jdbc.ClickHouseDriver")
  .save()

ClickHouse侧列定义示例:

CREATE TABLE target_table (
  id Int64,
  map_col String -- 或 JSON 类型
) ENGINE = MergeTree() ORDER BY id;

查询时解析JSON:

SELECT JSONExtractKeysAndValues(map_col, 'String', 'String') FROM target_table;

2. 使用ClickHouse官方Spark Connector

官方Connector原生支持Map类型读写,无需通过JDBC层适配:

originalDF.write
  .format("clickhouse")
  .option("host", "your_clickhouse_host")
  .option("port", "8123") -- 9000为原生协议端口
  .option("database", "your_db")
  .option("table", "your_table")
  .option("user", "your_user")
  .option("password", "your_password")
  .save()

需在Spark依赖中加入Connector包(根据Spark版本调整),Maven坐标示例:

<dependency>
  <groupId>com.clickhouse</groupId>
  <artifactId>clickhouse-spark-connector_2.12</artifactId>
  <version>0.5.0</version>
</dependency>

二、让集群加载修改后的JdbcUtils

若必须使用修改后的JdbcUtils,需解决类加载优先级问题:

  1. 编译自定义JAR:将修改后的JdbcUtils.scala编译为JAR包,确保依赖与集群Spark版本一致,仅打包org.apache.spark.sql.execution.datasources.jdbc.JdbcUtils类即可。
  2. 提交任务时指定优先加载:使用--jars传入自定义JAR,同时设置类加载优先级参数,让执行器优先加载自定义类:
spark-submit \
  --class your.main.application.class \
  --jars /path/to/custom-jdbc-utils.jar \
  --conf spark.driver.userClassPathFirst=true \
  --conf spark.executor.userClassPathFirst=true \
  your-application.jar

注意:userClassPathFirst可能引发类冲突,需测试验证;YARN模式下需将JAR上传至HDFS或通过--files分发。


内容的提问来源于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.25 05:55:12