Deprecated: imwpcache\f884414bce24ee67f\f73723ec7b1919fa5::__construct(): Implicitly marking parameter $YECBGYFECGEAFWHA as nullable is deprecated, the explicit nullable type must be used instead in /www/wwwroot/www.chuangxiangniao.com/wp-content/plugins/imwpcache-dist/build/f884414bce24ee67ff73723ec7b1919fa5.php on line 2

Deprecated: imwpcache\f884414bce24ee67f\f73723ec7b1919fa5::__construct(): Implicitly marking parameter $BBWFDDBHHYHDXXAB as nullable is deprecated, the explicit nullable type must be used instead in /www/wwwroot/www.chuangxiangniao.com/wp-content/plugins/imwpcache-dist/build/f884414bce24ee67ff73723ec7b1919fa5.php on line 2
PySpark中多层嵌套Array Struct的扁平化处理技巧_创想鸟

PySpark中多层嵌套Array Struct的扁平化处理技巧

PySpark中多层嵌套Array Struct的扁平化处理技巧

本文深入探讨了在PySpark中如何高效地将复杂的多层嵌套 array(struct(array(struct))) 结构扁平化为 array(struct)。通过结合使用Spark SQL的 transform 高阶函数和 flatten 函数,我们能够优雅地提取内层结构字段并与外层字段合并,最终实现目标模式的简化,避免了传统 explode 和 groupBy 组合的复杂性,提供了一种更具声明性和可扩展性的解决方案。

理解复杂嵌套结构与目标

在处理大数据时,我们经常会遇到包含复杂嵌套数据类型的dataframe。一个常见的场景是列中包含 array(struct(array(struct))) 类型的结构,例如:

root |-- a: integer (nullable = true) |-- list: array (nullable = true) |    |-- element: struct (containsNull = true) |    |    |-- b: integer (nullable = true) |    |    |-- sub_list: array (nullable = true) |    |    |    |-- element: struct (containsNull = true) |    |    |    |    |-- c: integer (nullable = true) |    |    |    |    |-- foo: string (nullable = true)

我们的目标是将这种多层嵌套结构简化为 array(struct) 形式,即把 sub_list 中的 c 和 foo 字段提升到 list 内部的 struct 中,并消除 sub_list 的嵌套层级:

root |-- a: integer (nullable = true) |-- list: array (nullable = true) |    |-- element: struct (containsNull = true) |    |    |-- b: integer (nullable = true) |    |    |-- c: integer (nullable = true) |    |    |-- foo: string (nullable true)

这种扁平化处理对于后续的数据分析和处理至关重要。

挑战与传统方法局限性

传统的扁平化方法通常涉及 explode 函数,它会将数组中的每个元素扩展为单独的行。对于上述结构,如果直接使用 explode,可能需要多次 explode 操作,然后通过 groupBy 和 collect_list 来重新聚合,这在面对更深层次的嵌套时会变得异常复杂和低效。例如,以下方法虽然有效,但在复杂场景下维护成本高昂:

from pyspark.sql import SparkSessionfrom pyspark.sql.functions import inline, expr, collect_list, struct# 假设df是您的DataFrame# df.select("a", inline("list")) # .select(expr("*"), inline("sub_list")) # .drop("sub_list") # .groupBy("a") # .agg(collect_list(struct("b", "c", "foo")).alias("list"))

这种方法要求我们将所有嵌套层级“提升”到行级别,然后再进行聚合,这与我们期望的“自底向上”或“原地”转换理念相悖。我们更倾向于一种能够直接在数组内部进行操作,而无需改变DataFrame行数的解决方案。

PySpark解决方案:Transform与Flatten的组合运用

PySpark 3.x 引入了 transform 等高阶函数,极大地增强了对复杂数据类型(特别是数组)的处理能力。结合 transform 和 flatten,我们可以优雅地解决上述问题。

transform 函数允许我们对数组中的每个元素应用一个自定义的转换逻辑,并返回一个新的数组。当涉及到多层嵌套时,我们可以使用嵌套的 transform 来逐层处理。

核心逻辑

内层转换:首先,对最内层的 sub_list 进行 transform 操作。对于 sub_list 中的每个元素(即包含 c 和 foo 的 struct),我们将其与外层 struct 中的 b 字段结合,创建一个新的扁平化 struct。外层转换:这一步的 transform 会生成一个 array(array(struct)) 的结构。扁平化:最后,使用 flatten 函数将 array(array(struct)) 结构合并成一个单一的 array(struct)。

示例代码

首先,我们创建一个模拟的DataFrame来演示:

from pyspark.sql import SparkSessionfrom pyspark.sql.functions import col, transform, flatten, structfrom pyspark.sql.types import StructType, StructField, ArrayType, IntegerType, StringType# 初始化SparkSessionspark = SparkSession.builder.appName("FlattenNestedArrayStruct").getOrCreate()# 定义初始schemainner_struct_schema = StructType([    StructField("c", IntegerType(), True),    StructField("foo", StringType(), True)])outer_struct_schema = StructType([    StructField("b", IntegerType(), True),    StructField("sub_list", ArrayType(inner_struct_schema), True)])df_schema = StructType([    StructField("a", IntegerType(), True),    StructField("list", ArrayType(outer_struct_schema), True)])# 创建示例数据data = [    (1, [        {"b": 10, "sub_list": [{"c": 100, "foo": "x"}, {"c": 101, "foo": "y"}]},        {"b": 20, "sub_list": [{"c": 200, "foo": "z"}]}    ]),    (2, [        {"b": 30, "sub_list": [{"c": 300, "foo": "w"}]}    ])]df = spark.createDataFrame(data, schema=df_schema)df.printSchema()df.show(truncate=False)# 应用扁平化逻辑df_flattened = df.withColumn(    "list",    flatten(        transform(            col("list"),  # 外层数组 (array of structs)            lambda x: transform(  # 对外层数组的每个struct x 进行操作                x.getField("sub_list"),  # 获取struct x 中的 sub_list (array of structs)                lambda y: struct(x.getField("b").alias("b"), y.getField("c").alias("c"), y.getField("foo").alias("foo")),            ),        )    ),)df_flattened.printSchema()df_flattened.show(truncate=False)# 停止SparkSessionspark.stop()

代码解析

df.withColumn(“list”, …): 我们选择修改 list 列,使其包含扁平化后的结果。transform(col(“list”), lambda x: …): 这是外层 transform。它遍历 list 列中的每一个 struct 元素,我们将其命名为 x。x 的类型是 struct(b: int, sub_list: array(struct(c: int, foo: string)))。transform(x.getField(“sub_list”), lambda y: …): 这是内层 transform。它作用于 x 中的 sub_list 字段。sub_list 是一个数组,它的每个元素(一个 struct(c: int, foo: string))被命名为 y。struct(x.getField(“b”).alias(“b”), y.getField(“c”).alias(“c”), y.getField(“foo”).alias(“foo”)): 在内层 transform 内部,我们构建一个新的 struct。x.getField(“b”): 从外层 struct x 中获取 b 字段。y.getField(“c”): 从内层 struct y 中获取 c 字段。y.getField(“foo”): 从内层 struct y 中获取 foo 字段。alias(“b”), alias(“c”), alias(“foo”) 用于确保新生成的 struct 字段名称正确。这个 struct 函数会为 sub_list 中的每个 y 元素生成一个扁平化的 struct。因此,内层 transform 的结果是一个 array(struct(b: int, c: int, foo: string))。中间结果:外层 transform 会收集所有这些 array(struct),因此它的最终输出是一个 array(array(struct(b: int, c: int, foo: string)))。flatten(…): 最后,flatten 函数将这个 array(array(struct)) 结构扁平化为一个单一的 array(struct(b: int, c: int, foo: string)),这正是我们期望的目标 schema。

注意事项与最佳实践

字段名称:确保在 getField() 和 struct() 中使用的字段名称与实际 schema 中的名称完全匹配。空值处理:transform 函数会自然地处理数组中的空元素。如果 sub_list 为空,内层 transform 会返回一个空数组;如果 list 为空,外层 transform 也会返回空数组。flatten 对空数组的处理也是安全的。性能:transform 是Spark SQL的内置高阶函数,通常比自定义UDF(用户定义函数)具有更好的性能,因为它可以在Spark Catalyst优化器中进行优化。可读性:虽然嵌套 transform 非常强大,但过度嵌套可能会降低代码的可读性。对于更复杂的场景,可以考虑将转换逻辑拆分成多个步骤或添加详细注释。通用性:这种 transform 结合 flatten 的模式可以推广到更深层次的嵌套结构,只需增加 transform 的嵌套层级即可。

总结

通过巧妙地结合使用PySpark的 transform 和 flatten 函数,我们能够以一种声明式且高效的方式,将复杂的多层嵌套 array(struct(array(struct))) 结构扁平化为更易于处理的 array(struct) 结构。这种方法避免了传统 explode 和 groupBy 组合的复杂性,特别适用于需要对数组内部元素进行精细化转换的场景,是处理Spark中复杂半结构化数据时一个非常有用的技巧。

以上就是PySpark中多层嵌套Array Struct的扁平化处理技巧的详细内容,更多请关注创想鸟其它相关文章!

版权声明:本文内容由互联网用户自发贡献,该文观点仅代表作者本人。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。
如发现本站有涉嫌抄袭侵权/违法违规的内容, 请发送邮件至 chuangxiangniao@163.com 举报,一经查实,本站将立刻删除。
发布者:程序猿,转转请注明出处:https://www.chuangxiangniao.com/p/1374680.html

(0)
打赏 微信扫一扫 微信扫一扫 支付宝扫一扫 支付宝扫一扫
在IIS 10上部署FastAPI应用的完整教程
上一篇 2025年12月14日 14:21:29
利用Tshark和PDML实现网络数据包十六进制字节到字段的映射
下一篇 2025年12月14日 14:21:47

相关推荐

  • Spring Boot自定义Kafka配置与动态Bean注册最佳实践

    本文探讨了在Spring Boot应用中通过自定义注解简化Kafka配置的挑战与解决方案。重点介绍了如何利用META-INF/spring.factories实现早期自动配置,并详细阐述了使用ImportBeanDefinitionRegistrar在应用上下文初始化早期动态注册Kafka生产者工厂…

    2026年9月22日
    000
  • mysql安装后怎么维护 mysql日常维护操作大全

    mysql安装后怎么维护 mysql日常维护操作大全mysql安装后怎么维护 mysql日常维护操作大全mysql安装后怎么维护 mysql日常维护操作大全mysql安装后怎么维护 mysql日常维护操作大全

    开启并分析慢查询日志以优化 sql 性能;2. 定期使用逻辑或物理方式备份数据并异地存储;3. 监控连接数和服务器资源,防止资源耗尽;4. 定期执行 analyze、optimize 和 check 表操作以维护表健康;5. 合理管理日志配置与清理策略。mysql 安装后的日常维护主要包括慢查询监控…

    2026年9月22日 用户投稿
    100
  • 深度解析蝴蝶号如何实现AI实景24小时无人直播

    深度解析蝴蝶号如何实现AI实景24小时无人直播深度解析蝴蝶号如何实现AI实景24小时无人直播深度解析蝴蝶号如何实现AI实景24小时无人直播深度解析蝴蝶号如何实现AI实景24小时无人直播

    蝴蝶号能实现ai实景24小时无人直播,主要靠智能中控系统+实景画面采集+自动化互动机制。一、ai中控系统作为“大脑”,自动控制画面切换、语音播报、商品推荐和评论区互动,具备一定判断能力,确保稳定性与持续性。二、实景画面采集作为“眼睛”,通过高清摄像头和云台控制,在门店、仓库等场景采集实时画面,保障真…

    2026年9月22日 用户投稿
    100
  • 在Java中如何开发简易问答社区

    答案是Java结合Spring Boot可快速构建问答社区,通过设计questions、answers、users三张表实现数据存储,使用JPA进行持久化,前端用HTML+JS调用后端API完成用户提问、回答、查看与互动功能。 开发一个简易问答社区,核心是实现用户提问、回答、查看问题和互动功能。Ja…

    2026年9月22日
    100
  • HitPawVideoEditor如何制作AI视频?教你快速创建AI内容的步骤

    答案是HitPaw Video Editor通过AI文本转视频、AI图片生成、智能抠图、自动字幕等功能,显著提升视频创作效率。它以“AI创作+人工精修”模式降低制作门槛,帮助用户快速生成初稿、丰富视觉素材、简化复杂操作,并支持快速迭代,但需避免过度依赖AI,仍需人工打磨以确保情感表达与叙事质量。 ☞…

    2026年9月22日
    000
  • linux系统下codeblocks控制台打印中文乱码[通俗易懂]

    linux系统下codeblocks控制台打印中文乱码[通俗易懂]linux系统下codeblocks控制台打印中文乱码[通俗易懂]linux系统下codeblocks控制台打印中文乱码[通俗易懂]linux系统下codeblocks控制台打印中文乱码[通俗易懂]

    大家好,很高兴再次和大家见面,我是你们的朋友全栈君。 在Linux系统下使用CodeBlocks时,如果在控制台中打印中文可能会遇到乱码问题。以下是解决这一问题的详细步骤: 首先,我们来看一下在Linux系统下安装CodeBlocks后,运行以下代码时出现的问题: #include #include…

    2026年9月22日 用户投稿
    600
  • 解决Android设备管理移除时的SecurityException

    本文将详细介绍如何解决在尝试从Android设备移除设备管理员时遇到的java.lang.SecurityException异常。该异常通常发生在尝试移除一个非测试用途的设备管理员应用时。通过修改应用的配置,将其临时标记为测试应用,可以绕过此安全限制,从而成功移除设备管理员。请务必注意,这种方法仅适…

    2026年9月22日
    100
  • 百度地图导航语音无法关闭怎么办 百度地图语音设置调整技巧

    首先要分清关闭的是导航语音播报还是语音唤醒功能。①关闭导航语音:打开百度地图→【我的】→【设置】→【导航设置】→【导航中语音】→选择【静音】。②关闭语音唤醒:进入【我的】→【设置】→【语音设置】→【智能语音】→关闭【说“小度小度”唤醒】开关。操作后仍无效可尝试更新App或重新登录账号。 百度地图导航…

    2026年9月22日
    000
  • 如何用Blender打造AI生成3D视频?免费软件制作AI视频的步骤

    如何用Blender打造AI生成3D视频?免费软件制作AI视频的步骤如何用Blender打造AI生成3D视频?免费软件制作AI视频的步骤如何用Blender打造AI生成3D视频?免费软件制作AI视频的步骤如何用Blender打造AI生成3D视频?免费软件制作AI视频的步骤

    答案是可行,通过Blender与免费AI工具结合,构建以AI辅助概念设计、纹理生成和动作参考,Blender主导建模、动画与渲染的混合工作流,实现高效3D视频创作。 ☞☞☞AI 智能聊天, 问答助手, AI 智能搜索, 免费无限量使用 DeepSeek R1 模型☜☜☜ 用Blender制作AI生成…

    2026年9月22日 用户投稿
    200
  • 悟空浏览器看小说章节错乱怎么办_悟空浏览器小说章节错乱解决方法

    1、清除悟空浏览器缓存:进入手机设置→悟空浏览器→存储→清除缓存,重启应用查看是否恢复。2、切换网络环境:断开当前Wi-Fi改用移动数据或更换其他Wi-Fi,关闭并重开浏览器,下拉刷新章节页面。3、更新或重装应用:前往应用商店更新悟空浏览器至最新版本;若无效则卸载后重新下载安装,登录账号后再次阅读,…

    2026年9月22日
    000
  • Linux基础必知必会(一)

    文章目录 前言 一、初识Linux操作系统 二、网络配置原理 三、虚拟机网络配置原理 四、虚拟机网络环境配置 五、远程工具Xshell 六、Linux目录结构讲解 七、Linux常用的命令讲解 八、用户和用户组的管理 结语 前言 为什么需要学习Linux系统? 许多人可能疑惑,为什么在当前可视化操作…

    2026年9月22日
    1200
  • 抖音号如何升级成企业号?升级成企业号需要多久?

    随着短视频平台的迅猛发展,抖音已成为企业进行品牌宣传与用户运营的核心渠道。将普通个人账号升级为企业号,不仅能够解锁更多营销工具,还能增强品牌的权威性与可信度。 一、抖音个人号怎样升级为企业号? 确认基本条件 在申请前,需确保账号已完成实名认证,且未有违反社区规范的行为。个人账号必须绑定手机号,并完善…

    2026年9月22日
    600
  • Java类中Jackson @JsonNaming策略的运行时内省

    本文介绍如何在运行时动态内省Java类上通过@JsonNaming注解配置的Jackson PropertyNamingStrategy。通过利用ObjectMapper的SerializationConfig和JacksonAnnotationIntrospector,开发者可以编程方式获取类的命…

    2026年9月22日
    600
  • VSCode安装C/C++文档查看 提升开发效率的VSCode技巧

    答案是利用C/C++扩展和cppreference插件实现高效文档查阅。首先安装微软官方C/C++扩展,启用智能感知与悬停提示;再安装cppreference扩展,通过命令面板直接搜索标准库函数,实现离线在线无缝查阅;结合Doxygen生成项目文档,使用“转到定义”功能快速跳转源码;同时借助Inte…

    2026年9月22日
    100
  • 高效利用 PriorityQueue 合并并排序多个列表

    本教程详细阐述了如何使用 Java 的 PriorityQueue 高效地合并并排序多个整数列表。文章首先指出将列表作为元素放入 PriorityQueue 的常见误区,进而纠正为应将单个整数元素放入队列。接着,它演示了如何正确声明、填充 PriorityQueue,并强调了通过循环调用 poll(…

    2026年9月22日
    400
  • mac怎么设置动态壁纸_mac设置动态壁纸方法

    可通过系统设置直接启用macOS自带的动态桌面;2. 利用iCloud同步功能实现Apple设备间的动态壁纸共享;3. 安装第三方应用扩展更多动态效果如视频或实时天气;4. 手动将.heic或.mov格式文件添加至指定目录以设置自定义动态壁纸。 如果您希望为Mac增添个性化视觉体验,可以将动态壁纸设…

    2026年9月22日
    300
  • 如何配置Android开发环境 Android Studio安装与JDK配置方法

    答案:配置Android开发环境需先安装JDK并设置环境变量,再下载安装Android Studio,配置SDK及虚拟设备,最后创建项目测试。具体步骤包括:1. 安装JDK 17并配置JAVA_HOME和Path;2. 从官网下载Android Studio并安装,自动集成SDK;3. 通过SDK …

    2026年9月22日
    200
  • ClipStudioPaintPro如何导出AI漫画图片?保存图像的详细指南

    导出AI漫画图片需通过Clip Studio Paint Pro的“文件”菜单选择“导出”,根据用途选单页、多页或Webtoon导出,推荐PNG用于高质量或透明背景需求,JPG用于网络分享以平衡文件大小与画质,设置300dpi以上分辨率确保清晰度,色彩配置选用sRGB保障跨平台一致性,批量导出时利用…

    2026年9月22日
    100
  • 拼多多商家版怎么赚钱?拼多多商家版怎么赚钱提现

    答案:通过拼多多商家版开设虚拟资料店铺,缴纳500元保证金注册个人店,选品如PPT模板、教程等高需求低风险商品,利用阿奇索等工具设置自动发货,结合低价引流和付费推广提升销量,订单完成后资金在T+1或T+7结算周期后提现至银行卡。 如果您想通过拼多多商家版实现盈利并完成资金提现,需要了解其核心的运营模…

    2026年9月22日
    500
  • 如何在mysql中实现热备份

    最推荐的MySQL热备份方案是结合Percona XtraBackup全量备份与binlog增量备份,并通过主从复制实现高可用。首先使用XtraBackup对InnoDB引擎进行在线全量备份,无需锁表;备份后执行–prepare确保数据一致性,恢复时用–copy-back还原…

    2026年9月22日
    100

发表回复

登录后才能评论
关注微信