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

如何将Databricks中的PySpark DataFrame转换为可重构代码?

如何将含嵌套结构的PySpark DataFrame转换为可复用的创建代码?

问题背景

假设在Databricks中有如下PySpark DataFrame:

some_other_columnprice_history
test1[{"date":"2021-03-21T01:20:33Z","price_tag":"N","price":"9.23","price_promotion":"1.34","AKT":false,"my_column":null,"supplier":"some_supplier"}]
test2[{"date":"2021-03-23T01:20:33Z","price_tag":null,"price":"10.40","price_promotion":null,"AKT":true,"my_column":"something","supplier":null}]

该DataFrame的price_history列包含嵌套的数组-结构体结构。目前仅能通过df.schema获取表结构,但生成的schema无法直接指导手动创建含复杂类型(如Struct、ArrayType等)的DataFrame。希望将现有DataFrame转换为可重新创建它的完整代码,方便编写测试时快速生成示例数据。

解决方案

以下是可直接复用的代码,完全匹配目标DataFrame的结构与数据:

from datetime import datetime
from decimal import Decimal

from pyspark.sql.types import (StructType, StructField, StringType,
                               ArrayType, TimestampType, BooleanType,
                               DecimalType)

# 定义包含嵌套结构的Schema
schema = StructType([
    StructField("some_other_column", StringType()),
    StructField('price_history', ArrayType(StructType([
        StructField('date', TimestampType()),
        StructField('price_tag', StringType()),
        StructField('price', DecimalType(12, 2)),
        StructField('price_promotion', DecimalType(12, 2)),
        StructField('AKT', BooleanType()),
        StructField("my_column", StringType()),
        StructField("supplier", StringType()),
    ]))),
])

# 构造与Schema匹配的数据集
data = [
    ('test1', [(datetime(2021, 3, 21, 1, 20, 33), 'N',
                Decimal('9.23'), Decimal('1.34'), False, None, 'some_supplier')]),
    ('test2', [(datetime(2021, 3, 23, 1, 20, 33), None,
                Decimal('10.40'), None, True, 'something', None)]),
]

# 创建目标DataFrame
df = spark.createDataFrame(schema=schema, data=data)

关键说明

  1. Schema定义:严格匹配原DataFrame的字段类型,包括嵌套数组内的结构体字段,Decimal类型需指定精度和小数位数以保证数据一致性。
  2. 数据转换:将原DataFrame中的JSON格式数据转换为Python原生类型,如时间字符串转datetime对象、价格字符串转Decimal对象,空值保持为None。
  3. 自动推导补充:若需从现有DataFrame自动生成代码,可编写辅助函数遍历df.schema生成类型定义,同时通过df.collect()获取行数据并转换为原生格式,但手动编写更可控,能避免自动推导可能出现的类型精度偏差。

内容的提问来源于stack exchange,提问作者the_economist

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.02 16:23:26