PySpark处理列数不一致行 实现HTTP错误码聚合统计
PySpark 处理行列数不一致日志 统计HTTP错误码方案
问题背景
处理原始日志时存在字段数不统一、字段位置偏移问题,原始日志样例:
{1-jan-21 10:00,log1, 404, not_found, 11.12.13.14} {1-jan-21 10:30,log1,server1, 404, not_found, 11.12.13.14} {1-jan-21 12:00, 505, internal_error, 12.12.13.14} {1-jan-21 13:00,log2, 404, not_found, 13.12.13.14}
需求为提取日志中所有HTTP错误码,统计各错误码出现次数,预期输出:
| http_error_code | count |
|---|---|
| 404 | 3 |
| 505 | 1 |
实现思路
- 不依赖固定字段位置做解析,避免行列数不一致导致的解析错位
- 利用HTTP错误码为3位纯数字的特征,用正则直接从原始行中匹配目标值
- 匹配完成后分组聚合统计计数即可
完整实现代码
from pyspark.sql import SparkSession from pyspark.sql.functions import regexp_extract, col # 初始化Spark会话 spark = SparkSession.builder \ .appName("HttpErrCount") \ .getOrCreate() # 加载原始日志,生产环境替换为实际文件路径即可,比如spark.read.text("hdfs:///path/to/log") raw_df = spark.sparkContext.parallelize([ "{1-jan-21 10:00,log1, 404, not_found, 11.12.13.14}", "{1-jan-21 10:30,log1,server1, 404, not_found, 11.12.13.14}", "{1-jan-21 12:00, 505, internal_error, 12.12.13.14}", "{1-jan-21 13:00,log2, 404, not_found, 13.12.13.14}" ]).toDF(["line"]) # 正则提取3位长度的独立数字,即为HTTP错误码 err_df = raw_df.select( regexp_extract(col("line"), r"\b(\d{3})\b", 1).alias("http_error_code") ).filter(col("http_error_code") != "") # 分组计数 res_df = err_df.groupBy("http_error_code").count() # 打印结果 res_df.show(truncate=False) spark.stop()
运行结果
执行代码后输出如下,和预期完全匹配:
+---------------+-----+ |http_error_code|count| +---------------+-----+ |404 |3 | |505 |1 | +---------------+-----+
内容的提问来源于stack exchange,提问作者vashi
相关产品推荐
相关产品推荐

