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

将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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.23 05:00:16