如何从Spark(Scala)向PostgreSQL的JSONB列插入JSON字符串
解决Spark(Scala)插入JSON字符串到PostgreSQL JSONB列的问题
错误原因分析
- 直接用
lit(st)添加的列是字符串类型(StringType),PostgreSQL JDBC驱动会将其识别为varchar,与目标jsonb列类型不匹配,触发类型转换错误。 - 使用
spray.json.JsObject作为字面量时,Spark SQL不支持该类型作为列值,因此抛出特性不支持的异常。
正确解法
方法一:将字符串转换为Spark JSON类型
Spark 2.2+原生支持JsonType类型,将字符串列转换为该类型后,JDBC驱动会自动映射到PostgreSQL的jsonb列。
import org.apache.spark.sql.types.DataTypes val st = """{"OFFLINE": {"CHEQUE": 5}, "ONLINE": {"CREDIT": 135, "DEBIT": 297}}""" val jsonData = Seq((1, "name1"), (2, "name2")) val jsonDF = spark.createDataFrame(jsonData).toDF("id", "sample_name") // 将字符串转换为Spark JSON类型,也可简写为lit(st).cast("json") val dfWithJSONB = jsonDF.withColumn("json_data", lit(st).cast(DataTypes.JsonType)) val jdbcUrl = "jdbc:postgresql://localhost:5432/db" dfWithJSONB.write .format("jdbc") .option("url", jdbcUrl) .option("user", "postgres") .option("password", "password") .option("dbtable", "jsonb_table") .option("driver", "org.postgresql.Driver") .mode("append") .save()
方法二:使用JDBC选项stringtype=unspecified
如果不想转换列类型,可通过该选项告诉PostgreSQL JDBC驱动不强制将字符串映射为varchar,由数据库自动判断类型,直接将JSON字符串插入jsonb列。
val st = """{"OFFLINE": {"CHEQUE": 5}, "ONLINE": {"CREDIT": 135, "DEBIT": 297}}""" val jsonData = Seq((1, "name1"), (2, "name2")) val jsonDF = spark.createDataFrame(jsonData).toDF("id", "sample_name") val dfWithJSONB = jsonDF.withColumn("json_data", lit(st)) val jdbcUrl = "jdbc:postgresql://localhost:5432/db" dfWithJSONB.write .format("jdbc") .option("url", jdbcUrl) .option("user", "postgres") .option("password", "password") .option("dbtable", "jsonb_table") .option("driver", "org.postgresql.Driver") .option("stringtype", "unspecified") // 关键配置 .mode("append") .save()
注意事项
- 确保PostgreSQL JDBC驱动版本不低于42.2.x,旧版本对JSONB的支持存在缺陷。
- 若JSON字符串是从DataFrame其他列动态生成的,可先通过
to_json函数将StructType数据转换为JSON字符串,再按上述方法处理。
内容的提问来源于stack exchange,提问作者SRN
相关产品推荐
相关产品推荐

