将PostgreSQL的JSONB数据读取为PySpark的MapType并提取字段问题
问题
尝试将PostgreSQL中的JSONB类型数据读取为PySpark的MapType(字典格式),并从attributes列提取cost和size到单独列,但PySpark始终将JSONB识别为字符串,无法直接按Map方式提取字段。
环境信息
PostgreSQL表结构与数据:
create table products (product_id varchar, description varchar, attributes jsonb, tax_rate decimal); insert into products values ('P1', 'Detergent', '{"cost": 45.50, "size": "10g"}', 5.0 ); insert into products values ('P2', 'Bread', '{"cost": 45.5, "size": "200g"}',3.5);
原问题代码:
import pyspark from pyspark.sql import SparkSession from pyspark.sql.functions import col from pyspark.sql.types import StructType, StructField, StringType, IntegerType, MapType, DecimalType spark = SparkSession \ .builder \ .appName("Python Spark SQL basic example") \ .config("spark.jars", "C:\Users\nupsingh\Documents\Jars\postgresql-42.7.3.jar") \ .getOrCreate() schema = StructType([ StructField('product_id', StringType(), True), StructField('description', StringType(), True), StructField('attributes', MapType(StringType(),IntegerType()),False), StructField('tax_rate', DecimalType(), True) ]) df = spark.read \ .format("jdbc") \ .option("url", "jdbc:postgresql://localhost:5432/postgres") \ .option("dbtable", "products") \ .option("user", "user1") \ .option("password", "password") \ .option("driver", "org.postgresql.Driver") \ .option("schema", schema) \ .load() df.show() df.printSchema() attributes_col = df.select("attributes") attributes_col.show() products_df = attributes_col.withColumn("cost", col("attributes")["cost"]).withColumn("size", col("attributes")["size"]) products_df.show()
原因
PostgreSQL的JDBC驱动默认将JSONB映射为字符串,直接指定PySpark的MapType schema不会生效——JDBC无法自动完成JSON字符串到Map的转换。
解决方案
方法1:在PostgreSQL层面预处理字段
直接通过SQL查询提前提取JSONB中的字段,避免后续PySpark转换:
df = spark.read \ .format("jdbc") \ .option("url", "jdbc:postgresql://localhost:5432/postgres") \ .option("dbtable", """ SELECT product_id, description, (attributes->>'cost')::decimal AS cost, attributes->>'size' AS size, tax_rate FROM products """) \ .option("user", "user1") \ .option("password", "password") \ .option("driver", "org.postgresql.Driver") \ .load() df.show() df.printSchema()
方法2:读取后用PySpark转成Map/Struct类型
先读取原始数据(attributes为字符串),再用from_json转换为结构化类型:
from pyspark.sql.functions import from_json # 读取原始数据 df = spark.read \ .format("jdbc") \ .option("url", "jdbc:postgresql://localhost:5432/postgres") \ .option("dbtable", "products") \ .option("user", "user1") \ .option("password", "password") \ .option("driver", "org.postgresql.Driver") \ .load() # 定义精准的StructSchema(比MapType更适合固定结构的JSON) struct_schema = StructType([ StructField("cost", DecimalType(), True), StructField("size", StringType(), True) ]) # 将JSON字符串转成StructType df = df.withColumn("attributes_struct", from_json(col("attributes"), struct_schema)) # 提取字段 products_df = df.withColumn("cost", col("attributes_struct.cost")) \ .withColumn("size", col("attributes_struct.size")) products_df.select("product_id", "description", "cost", "size", "tax_rate").show()
方法3:直接用get_json_object提取单个字段
若仅需少量字段,可直接用该函数提取:
from pyspark.sql.functions import get_json_object df = spark.read \ .format("jdbc") \ .option("url", "jdbc:postgresql://localhost:5432/postgres") \ .option("dbtable", "products") \ .option("user", "user1") \ .option("password", "password") \ .option("driver", "org.postgresql.Driver") \ .load() products_df = df.withColumn("cost", get_json_object(col("attributes"), "$.cost").cast(DecimalType())) \ .withColumn("size", get_json_object(col("attributes"), "$.size")) products_df.select("product_id", "description", "cost", "size", "tax_rate").show()
注意事项
- 原代码中
MapType(StringType(), IntegerType())不合理:cost是小数类型,应使用DecimalType或DoubleType,否则会导致转换失败。 - PostgreSQL的JSONB转字符串后为标准JSON格式,PySpark的
from_json可正常解析。
内容的提问来源于stack exchange,提问作者Nupur
相关产品推荐
相关产品推荐

