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 将 JSON 属性数据透视为表格列_创想鸟

使用 PySpark 将 JSON 属性数据透视为表格列

使用 PySpark 将 JSON 属性数据透视为表格列

本教程详细介绍了如何使用 PySpark 将 Oracle REST API 返回的 JSON 数组数据(其中属性名和属性值以键值对形式存在)转换为结构化的表格格式。通过 PySpark 读取 JSON 数据并结合 Spark SQL 的 MAX(CASE WHEN …) 语句,实现将动态属性名称(如 ‘LOG_ID’ 和 ‘BUSINESS_UNIT’)透视为独立的列,从而方便数据分析和处理。

在数据集成和处理过程中,我们经常会遇到来自 rest api 的响应数据,其结构可能并非传统的行列表格形式。例如,某些 api 会以键值对数组的形式返回数据,其中每个对象包含一个属性名(attributename)和对应的属性值(attributevalue)。当需要将这些动态的属性名转换为固定的列,并将其对应的属性值填充到这些列中时,传统的转换方法可能不够灵活。本教程将展示如何利用 pyspark 的强大能力,特别是结合 spark sql,高效地实现这种数据透视操作。

问题描述

假设我们从 Oracle REST API 获得以下 JSON 响应数据:

[    {        "attributeId": 300000000227671,        "attributeName": "BUSINESS_UNIT",        "attributeType": "Number",        "attributeValue": "300000207138371",        "timeBuildingBlockId": 300000300319699,        "timeBuildingBlockVersion": 1    },    {        "attributeId": 300000000226689,        "attributeName": "LOG_ID",        "attributeType": "Number",        "attributeValue": "300000001228038",        "timeBuildingBlockId": 300000300319699,        "timeBuildingBlockVersion": 1    }]

我们的目标是将 attributeName 为 ‘LOG_ID’ 和 ‘BUSINESS_UNIT’ 的 attributeValue 提取出来,并将其转换为以下表格形式:

LOG_ID BUSINESS_UNIT

300000001228038300000207138371

解决方案:使用 PySpark 和 Spark SQL

PySpark 提供了强大的数据处理能力,结合 Spark SQL,可以非常灵活地处理这种数据透视场景。核心思路是先将 JSON 数据加载到 DataFrame 中,然后利用 Spark SQL 的条件聚合函数(CASE WHEN 和 MAX)实现透视。

步骤一:加载 JSON 数据到 DataFrame

首先,我们需要将 JSON 响应数据加载到 PySpark DataFrame 中。假设 json_data 是包含上述 JSON 字符串的变量。

from pyspark.sql import SparkSession# 初始化 SparkSessionspark = SparkSession.builder.appName("JsonPivotTutorial").getOrCreate()sc = spark.sparkContext# 模拟 JSON 数据,实际应用中可能是从文件或API响应获取json_data = """[    {        "attributeId": 300000000227671,        "attributeName": "BUSINESS_UNIT",        "attributeType": "Number",        "attributeValue": "300000207138371",        "timeBuildingBlockId": 300000300319699,        "timeBuildingBlockVersion": 1    },    {        "attributeId": 300000000226689,        "attributeName": "LOG_ID",        "attributeType": "Number",        "attributeValue": "300000001228038",        "timeBuildingBlockId": 300000300319699,        "timeBuildingBlockVersion": 1    }]"""# 将 JSON 字符串转换为 RDD 并读取为 DataFrame# 注意:如果 json_data 是一个列表,可以直接使用 spark.createDataFrame()# 但如果是一个多行 JSON 字符串,或者需要更灵活地处理,spark.read.json(sc.parallelize([json_data])) 是一个有效方法df = spark.read.json(sc.parallelize([json_data]))# 查看原始 DataFrame 结构df.printSchema()df.show(truncate=False)

执行上述代码后,df 将包含解析后的 JSON 数据,每行对应 JSON 数组中的一个对象。

步骤二:创建临时视图

为了方便使用 Spark SQL 进行查询,我们将 DataFrame 注册为一个临时视图(Temporary View)。

df.createOrReplaceTempView("myTable")

现在,我们可以像操作传统数据库表一样,通过 SQL 语句查询 myTable。

步骤三:使用 Spark SQL 进行数据透视

透视的核心在于使用 CASE WHEN 语句根据 attributeName 的值选择对应的 attributeValue,并通过聚合函数(如 MAX)将每个组中的非空值提取出来。由于我们希望将所有相关属性(例如 LOG_ID 和 BUSINESS_UNIT,它们共享相同的 timeBuildingBlockId 和 timeBuildingBlockVersion)聚合到一行,因此需要对这些共享字段进行隐式分组。

result = spark.sql("""    SELECT        MAX(CASE WHEN attributeName = 'LOG_ID' THEN attributeValue END) AS LOG_ID,        MAX(CASE WHEN attributeName = 'BUSINESS_UNIT' THEN attributeValue END) AS BUSINESS_UNIT    FROM myTable    GROUP BY timeBuildingBlockId, timeBuildingBlockVersion -- 根据业务逻辑分组,确保同一逻辑实体的数据聚合到一行""")result.show()

SQL 逻辑解释:

CASE WHEN attributeName = ‘LOG_ID’ THEN attributeValue END: 这部分逻辑会检查 attributeName 是否为 ‘LOG_ID’。如果是,则返回对应的 attributeValue;否则返回 NULL。MAX(…) AS LOG_ID: 由于每个 attributeName 对应的 attributeValue 在原始数据中只出现一次(对于特定的逻辑实体),所以 MAX 函数会从 CASE WHEN 表达式生成的多个 NULL 值和一个非 NULL 值中选择那个非 NULL 的 attributeValue。这有效地将特定属性的 attributeValue 提升为新的列。GROUP BY timeBuildingBlockId, timeBuildingBlockVersion: 这一步至关重要。原始 JSON 数据中,LOG_ID 和 BUSINESS_UNIT 属于同一个逻辑实体,它们共享相同的 timeBuildingBlockId 和 timeBuildingBlockVersion。通过对这些字段进行分组,我们可以确保属于同一逻辑实体(即同一组)的所有属性值被聚合到同一行中。如果没有 GROUP BY,或者分组字段选择不当,可能会导致结果不正确(例如,所有属性聚合到一行,或者数据被错误地分割)。

输出结果:

+---------------+-------------------+|LOG_ID         |BUSINESS_UNIT      |+---------------+-------------------+|300000001228038|300000207138371|+---------------+-------------------+

这正是我们期望的透视结果。

注意事项与总结

动态列处理: 上述方法适用于列名(LOG_ID, BUSINESS_UNIT)已知的情况。如果 attributeName 的种类是动态变化的,并且需要在运行时确定列名,则需要结合 PySpark 的 DataFrame API 中的 pivot 函数,或者在 Spark SQL 中使用动态 SQL 生成技术。然而,对于固定的少量列,CASE WHEN 语句更直接和高效。聚合函数选择: 除了 MAX,也可以根据实际需求选择其他聚合函数,如 MIN、SUM、AVG 等。但对于这种将单个值提升为列的场景,MAX(或 MIN)是最常见的选择,因为它会忽略 NULL 值并返回唯一的非 NULL 值。分组键的重要性: GROUP BY 子句的选择至关重要。它决定了哪些原始行的数据会被聚合成新的一行。在上述示例中,timeBuildingBlockId 和 timeBuildingBlockVersion 共同标识了一个唯一的业务实体,因此它们是理想的分组键。务必根据您的数据模型和业务需求来确定正确的分组键。性能考量: 对于非常大的数据集,Spark SQL 能够有效地并行处理数据。然而,过多的 CASE WHEN 表达式或过于复杂的分组逻辑可能会影响性能。在实际应用中,应根据数据量和集群资源进行调优。

通过 PySpark 和 Spark SQL 的结合,我们可以灵活高效地处理各种复杂的数据转换需求,将非结构化或半结构化的 JSON 数据转换为易于分析的表格格式。

以上就是使用 PySpark 将 JSON 属性数据透视为表格列的详细内容,更多请关注创想鸟其它相关文章!

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

(0)
打赏 微信扫一扫 微信扫一扫 支付宝扫一扫 支付宝扫一扫
Python源码实现在线视频转字幕 利用ASR模型的Python源码对接流程
上一篇 2025年12月14日 08:28:05
Python函数如何用闭包保存函数内部状态 Python函数闭包基础用法的入门操作指南​
下一篇 2025年12月14日 08:28:19

相关推荐

  • MySQL最新版本如何下载?官方下载指南

    MySQL最新版本如何下载?官方下载指南MySQL最新版本如何下载?官方下载指南MySQL最新版本如何下载?官方下载指南MySQL最新版本如何下载?官方下载指南

    要下载mysql,推荐从官网直接下载;选择社区版或商业版取决于用途;下载时需选对操作系统和版本;安装遇到问题可查错误提示并搜索解决方案;验证安装成功可用命令行登录。下载步骤包括访问官网、选择版本与操作系统、使用installer、注册账号、开始下载安装。安装后配置root密码、字符集等。验证方式为命…

    2026年9月22日 用户投稿
    100
  • Flyway多数据库与CI/CD测试集成策略

    本文深入探讨了在CI/CD流程中,如何高效地配置Flyway以管理多数据库环境下的迁移,尤其关注集成测试场景。我们将比较使用真实数据库服务、Testcontainers以及Flyway自身多数据库配置的优劣,并提供关于分离生产与测试环境迁移脚本的实用策略,旨在确保开发、测试与生产环境的数据一致性与流…

    2026年9月22日
    100
  • Oracle定时任务

    oracle 介绍 oracle job 是应用在数据库层面,用来定时执行存储过程或者 SQL 语句的定时器。 测试数据代码语言:javascript代码运行次数:0运行复制 –建测试表CREATE TABLE TEST_JOB( DATETIME DATE, BEG_DATE VARCHAR2(…

    2026年9月22日
    500
  • MySQL复杂查询语句写作技巧_Sublime环境中编写多表关联逻辑

    MySQL复杂查询语句写作技巧_Sublime环境中编写多表关联逻辑MySQL复杂查询语句写作技巧_Sublime环境中编写多表关联逻辑MySQL复杂查询语句写作技巧_Sublime环境中编写多表关联逻辑MySQL复杂查询语句写作技巧_Sublime环境中编写多表关联逻辑

    提升mysql多表查询性能与可读性的方法包括:1. 优化索引,确保join和where字段有合适索引,理解复合索引左前缀原则;2. 使用cte分解逻辑,使结构清晰易维护;3. 利用sublime text插件如sqltools、sublimelinter提升编写效率;4. 拆解复杂逻辑,逐步构建查询…

    2026年9月22日 用户投稿
    200
  • 深入理解PHP数组中JSON字符串的解析与数据提取

    本文将详细讲解如何在PHP中处理包含JSON格式字符串的数组。通过使用json_decode函数,我们可以将这些JSON字符串转换为可操作的PHP数组,进而轻松提取所需的shortname和fullname等键值对。教程将提供清晰的示例代码,演示循环遍历和直接访问两种数据提取方式,帮助开发者高效地解…

    2026年9月22日
    300
  • PHP each() 函数的替代方案:自定义实现与常见错误修正

    本文探讨了PHP中已废弃的each()函数的替代方案。针对常见的自定义实现,如myEach(),文章详细指出了其在返回数组结构中常犯的错误,并提供了正确的代码示例,以确保替代函数能够模拟each()的预期行为,帮助开发者编写更健壮、兼容未来的PHP代码。 理解 each() 函数及其废弃背景 在PH…

    2026年9月22日
    000
  • MySQL数据备份自动化实施_MySQL定时任务与脚本管理

    MySQL数据备份自动化实施_MySQL定时任务与脚本管理MySQL数据备份自动化实施_MySQL定时任务与脚本管理MySQL数据备份自动化实施_MySQL定时任务与脚本管理MySQL数据备份自动化实施_MySQL定时任务与脚本管理

    mysql数据备份的自动化实施核心在于结合mysqldump等工具与操作系统的定时任务(如linux的cron或windows的task scheduler),通过编写和管理脚本实现定期执行备份。1. 使用mysqldump作为基础工具,编写包含数据库连接信息、时间戳文件名、日志记录、压缩清理等功能…

    2026年9月21日 用户投稿
    200
  • PHP/MySQL:高效合并订单商品并按日期分组显示

    本教程将指导如何在PHP/MySQL应用中,将同一日期的订单商品合并显示在同一行,以提高数据展示的清晰度。核心解决方案是利用MySQL的GROUP_CONCAT函数在数据库层面进行高效聚合,避免复杂的PHP逻辑处理,从而简化代码并优化性能。 订单数据展示的常见挑战 在开发在线购物平台时,通常需要向用…

    2026年9月21日
    200
  • Java ConcurrentSkipListMap在并发场景下应用

    ConcurrentSkipListMap是基于跳跃表实现的线程安全有序映射,支持高并发读写与高效范围查询,适用于需排序的并发场景,如排行榜系统;相比ConcurrentHashMap,它提供有序性与导航操作,但插入查找为O(log n),内存开销较大,适合读多写少或需区间扫描的业务。 在高并发场景…

    2026年9月21日
    100
  • 怎么全选VSCode多个光标_VSCode多光标操作与批量选择文本教程

    VSCode中高效创建多光标的方法包括:Alt+Click手动添加光标,适用于不规则位置;Ctrl+Alt+方向键垂直添加光标,适合连续多行操作;Ctrl+D逐个选择匹配项,精准控制选择范围;Ctrl+Shift+L一次性选择所有匹配项,实现全局批量修改。结合查找替换和列选择模式可进一步提升编辑效率…

    2026年9月21日
    100
  • avg计算平均值在mysql中如何使用

    AVG()是MySQL中计算列平均值的聚合函数,忽略NULL值。基本语法为SELECT AVG(列名) FROM 表名;可结合WHERE筛选条件,如SELECT AVG(score) FROM students WHERE subject = ‘math’ AND score…

    2026年9月21日
    100
  • MySQL缓存机制对性能提升的作用_MySQL缓存配置及调优方案

    MySQL缓存机制对性能提升的作用_MySQL缓存配置及调优方案MySQL缓存机制对性能提升的作用_MySQL缓存配置及调优方案MySQL缓存机制对性能提升的作用_MySQL缓存配置及调优方案MySQL缓存机制对性能提升的作用_MySQL缓存配置及调优方案

    mysql的缓存机制主要包括innodb缓冲池、查询缓存和操作系统文件系统缓存等,其中innodb缓冲池是性能优化的核心。1. innodb缓冲池缓存表数据和索引页,减少磁盘i/o,提升读写效率;2. 查询缓存因失效频繁及锁竞争问题,在高并发场景下易成瓶颈,已在mysql 8.0中移除;3. 操作系…

    2026年9月21日 用户投稿
    200
  • VSCode编写Java代码方法_VSCode搭建Java开发环境实战教程

    答案:在VSCode中配置Java开发环境需安装JDK并设置环境变量,再安装VSCode及Java扩展包,即可实现Java项目的创建、编写、运行与调试。它轻量、启动快,支持多语言和丰富扩展,集成Maven/Gradle,适合日常开发。 在VSCode里编写Java代码,说白了,就是把这个轻量级的代码…

    2026年9月21日
    100
  • Guava Multimap:高效获取并打印指定键的所有关联值

    guava multimap是处理一键多值映射关系的强大工具。要获取特定键的所有关联值,应直接使用其提供的`multimap#get(k)`方法。该方法会返回一个包含所有匹配值的`collection`,即使键不存在,也会返回一个空集合而非`null`,从而简化了值检索和空值处理逻辑,是比手动迭代键…

    2026年9月21日
    200
  • mysql如何在SQL中使用聚合函数

    聚合函数用于统计计算并返回单个值,常见函数有COUNT、SUM、AVG、MAX、MIN,通常与GROUP BY配合使用。1. COUNT统计非空值或总行数,SUM求和,AVG求平均,MAX和MIN分别取最大最小值。2. 对orders表整体统计可得总订单数、总额等信息。3. 按user_id分组后可…

    2026年9月21日
    400
  • mysql如何理解视图

    视图是基于SQL查询的虚拟表,不存储数据,每次查询时动态生成结果。1. 简化复杂查询,封装多表关联;2. 提高安全性,限制数据访问;3. 保持逻辑一致,避免重复定义;4. 兼容旧程序,表结构变更时减少修改;5. 更新受限,仅简单单表视图可写;6. 无性能提升,需依赖基础表索引优化。 视图在MySQL…

    2026年9月20日
    000
  • min和max在mysql中如何使用

    MIN()和MAX()用于查找列中的最小值和最大值,常用于数值、日期或字符串类型;基本语法为SELECT MIN(列名), MAX(列名) FROM 表名 [WHERE 条件];可单独或同时使用,如查询商品表中价格的最低与最高值;在日期字段中可找出最早和最晚时间;结合WHERE可按条件过滤,如统计某…

    2026年9月20日
    000
  • 在Linux中如何通过命令行安装Java JDK

    在Linux中安装Java JDK可通过包管理器或手动安装,推荐使用系统自带工具安装OpenJDK。对于Ubuntu/Debian系统,先更新软件包列表:sudo apt update,再搜索可用版本:apt search openjdk-*,然后安装指定版本如OpenJDK 17:sudo apt…

    2026年9月20日
    000
  • Java中如何高效地合并两个Map对象

    合并Map主要有三种方式:putAll()用于可变Map且性能高,Stream API适合不可变合并并支持冲突处理,Map.ofEntries()适用于小规模静态数据;选择依据是版本、是否需保持不可变及性能需求。 在Java中合并两个Map对象是常见操作,尤其在处理配置、缓存或数据聚合时。高效的方式…

    2026年9月20日
    200
  • group by分组在mysql中如何使用

    GROUP BY用于按列分组数据并配合聚合函数统计,如SELECT customer_id, SUM(amount) FROM orders GROUP BY customer_id计算每位客户总消费;可多字段分组如按客户和商品统计;结合WHERE过滤原始数据,HAVING筛选分组结果,常用函数有C…

    2026年9月20日
    100

发表回复

登录后才能评论
关注微信