PySpark过滤JSON数据:排除指定username并按时间戳筛选
日志过滤实现方法
你需要筛选的username和审计时间戳都嵌套在properties.log这个JSON字符串字段里,需要先解析该字段再做筛选,推荐在Spark计算阶段直接完成过滤,避免拉取无效数据到Driver端占用内存。
核心过滤逻辑
- 剔除规则:解析
properties.log得到的user.username字段值不等于system:serviceaccount:internal - 时间筛选规则:将解析得到的
requestReceivedTimestamp(请求接收时间)或stageTimestamp(审计阶段完成时间)转为时间类型后,按实际需要的时间范围匹配即可
修改后的完整代码
import re import json # 补全缺失的glob规则转换依赖 from fnmatch import translate as glob2re %pip install azure import azure from azure.storage.blob import AppendBlobService from pyspark.sql import functions as F abs = AppendBlobService(account_name="azdevstoreforlogs", account_key="mykey") base_path = "resourceId=/SUBSCRIPTIONS/subscriptionid/RESOURCEGROUPS/AZURE-DEV/PROVIDERS/MICROSOFT.CONTAINERSERVICE/MANAGEDCLUSTERS/AZURE-DEV" pattern = base_path + "/*/*/*/*/m=00/*.json" filter_re = glob2re(pattern) df = ( spark.sparkContext.parallelize( [ blob.name for blob in abs.list_blobs("insights-logs-kube-audit", prefix=base_path) if re.match(filter_re, blob.name) ] ) .map( lambda blob_name: abs.get_blob_to_bytes("insights-logs-kube-audit", blob_name) .content.decode("utf-8") .splitlines() ) .flatMap(lambda lines: [json.loads(l) for l in lines]) .toDF() # 解析嵌套的log字段为结构化格式 .withColumn("audit_log", F.from_json("properties.log", "struct<user:struct<username:string>, requestReceivedTimestamp:string, stageTimestamp:string>")) # 过滤指定serviceaccount的记录 .filter(F.col("audit_log.user.username") != "system:serviceaccount:internal") # 时间戳筛选示例:筛选2022-05-23 13:45:00之后的记录,可根据实际时间范围修改 .withColumn("event_ts", F.to_timestamp("audit_log.requestReceivedTimestamp", "yyyy-MM-dd'T'HH:mm:ss.SSSSSS'Z'")) .filter(F.col("event_ts") >= "2022-05-23 13:45:00") # 取前10条结果 .take(10) )
小数据量本地快速过滤方法
如果你已经拿到了take(10)返回的Row列表,不想重新跑Spark任务,可以直接在本地遍历过滤,示例代码如下:
from datetime import datetime filtered_result = [] # 定义目标时间范围,示例为2022-05-23 13:45:13之后 target_start = datetime.fromisoformat("2022-05-23T13:45:13+00:00") for row in df1: # 解析嵌套的log字段 log_data = json.loads(row.properties["log"]) # 跳过指定username的记录 if log_data["user"]["username"] == "system:serviceaccount:internal": continue # 时间范围校验 event_ts = datetime.fromisoformat(log_data["requestReceivedTimestamp"].replace("Z", "+00:00")) if event_ts >= target_start: filtered_result.append(row)
注意:本地过滤仅适合数据量极小的场景,数据量较大时务必在Spark侧做过滤,避免Driver端内存溢出。
内容的提问来源于stack exchange,提问作者ZZZSharePoint
相关产品推荐
相关产品推荐

