如何使用Pyspark提取Azure Application Insights事件并转为表格格式
问题原因
Application Insights查询API返回的结果为嵌套结构,实际的行列数据存储在响应根节点下的tables数组中,直接读取完整响应JSON只会得到顶层结构,无法生成预期的明细表格。
推荐解决方案(直接提取数据生成结构化DataFrame,效率更高)
import requests import json from pyspark.sql.types import TimestampType appId = "..." appKey = "..." query = """traces | where timestamp > ago(1d) | order by timestamp""" params = {"query": query} headers = {'X-Api-Key': appKey} url = f'https://api.applicationinsights.io/v1/apps/{appId}/query' response = requests.get(url, headers=headers, params=params) logs = json.loads(response.text) # 提取返回结果的第一份表(单query查询对应只有一张表) result_table = logs['tables'][0] # 提取列名列表 columns = [col['name'] for col in result_table['columns']] # 提取行数据列表 rows = result_table['rows'] # 直接基于行列数据生成Spark DataFrame df = spark.createDataFrame(rows, schema=columns) # 可选:调整字段类型,比如将timestamp从字符串转为时间类型 df = df.withColumn("timestamp", df["timestamp"].cast(TimestampType())) display(df)
保留原JSON读取流程的处理方案
如果你需要保留原有读取JSON的逻辑,可通过展开嵌套字段实现结构化转换:
import requests import json from pyspark.sql.functions import explode, col appId = "..." appKey = "..." query = """traces | where timestamp > ago(1d) | order by timestamp""" params = {"query": query} headers = {'X-Api-Key': appKey} url = f'https://api.applicationinsights.io/v1/apps/{appId}/query' response = requests.get(url, headers=headers, params=params) logs = json.loads(response.text) json_str = json.dumps(logs) jsonRDD = sc.parallelize([json_str]) df = spark.read.option('multiline', "true").json(jsonRDD) # 展开嵌套的tables数组 df = df.select(explode(col("tables")).alias("table")) # 展开行数据,关联列名定义 df = df.select(explode(col("table.rows")).alias("row"), col("table.columns.name").alias("col_names")) # 遍历列名,将行数组的元素映射为独立字段 col_list = df.select("col_names").first()[0] for idx, col_name in enumerate(col_list): df = df.withColumn(col_name, col("row")[idx]) # 清理中间辅助字段 df = df.drop("row", "col_names") display(df)
内容的提问来源于stack exchange,提问作者Dipanjan Mallick
相关产品推荐
相关产品推荐

