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,需解决类加载优先级问题:
- 编译自定义JAR:将修改后的
JdbcUtils.scala编译为JAR包,确保依赖与集群Spark版本一致,仅打包org.apache.spark.sql.execution.datasources.jdbc.JdbcUtils类即可。 - 提交任务时指定优先加载:使用
--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
相关产品推荐
相关产品推荐

