You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.29 20:21:34