在ADF或PySpark中实现Oracle REST API返回JSON的属性选择与透视
用Azure Data Factory(ADF)或PySpark处理Oracle REST API响应数据并生成透视表
一、Azure Data Factory(ADF)实现步骤
1. 获取API响应数据
- 使用Lookup活动调用Oracle REST API,配置正确的API端点、认证信息,将返回的JSON数组作为数据源读取。
2. 数据透视转换
- 添加Data Flow活动,将Lookup活动的输出作为数据源输入;
- 配置Pivot转换:
- 分组依据选择
timeBuildingBlockId和timeBuildingBlockVersion(确保同批次数据被正确聚合); - 透视列指定为
attributeName; - 值列选择
attributeValue,聚合方式选First(同分组下每个属性值唯一)。
- 分组依据选择
3. 筛选与输出
- 添加Select转换,仅保留
LDG_ID(注:原API响应中无LOG_ID字段,推测为笔误)和BUSINESS_UNIT两列; - 将转换后的数据输出到目标存储或数据库,即可得到所需的透视表结构。
二、PySpark代码实现
from pyspark.sql import SparkSession from pyspark.sql.functions import first # 初始化SparkSession spark = SparkSession.builder.appName("OracleAPIResponsePivot").getOrCreate() # 模拟API响应数据(实际可通过spark.read.json读取API返回文件或直接请求API获取) data = [ (300000000227671, "BUSINESS_UNIT", "Number", "300000207138371", 300000300319699, 1), (300000000227689, "LDG_ID", "Number", "300000001228038", 300000300319699, 1) ] columns = ["attributeId", "attributeName", "attributeType", "attributeValue", "timeBuildingBlockId", "timeBuildingBlockVersion"] df = spark.createDataFrame(data, columns) # 执行透视操作 pivoted_df = df.groupBy("timeBuildingBlockId", "timeBuildingBlockVersion") \ .pivot("attributeName") \ .agg(first("attributeValue")) # 选择所需列(若需LOG_ID,需确认API是否存在该字段,此处按现有LDG_ID处理) result_df = pivoted_df.select("LDG_ID", "BUSINESS_UNIT") # 显示结果 result_df.show() # 输出到目标位置(如Parquet、数据库等) # result_df.write.format("parquet").save("/path/to/output")
执行代码后输出结果如下:
| LDG_ID | BUSINESS_UNIT |
|---|---|
| 300000001228038 | 300000207138371 |
内容的提问来源于stack exchange,提问作者Advaitha
相关产品推荐
相关产品推荐

