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

使用Apache Beam过滤指定关键词日志行并存储为Pandas DataFrame

问题描述

我有一个包含10000条日志行的文本文件,示例日志行如下:

2022-12-27T00:00:00+00:00 VM_DEV02 sshd[25690]: pam_unix(sshd:session): session closed for user USER7

需要完成两个核心任务:

  • 仅过滤出包含['unauthorized','error','kernel error','OS error','rejected','warning']这些关键词的日志行;
  • 将日志行拆分为不同字段,并通过Apache Beam将数据存储到Pandas DataFrame中。

我编写了以下代码,但未达到预期效果:

import apache_beam as beam
from apache_beam.io import ReadFromText
from apache_beam.pvalue import AsList
from apache_beam.transforms import Map, Filter
import pandas as pd

def extract_fields(line):
    timestamp = line.split(" ")[0]
    print(timestamp)
    hostname = line.split(" ")[1]
    print(hostname)
    process_name = line.split(" ")[2].split("[")[0]
    print(process_name)
    pid = line.split(" ")[2].split("[")[1].split("]")[0]
    print(pid)
    text = " ".join(line.split(" ")[3:])
    print(text)

    return [timestamp, hostname, process_name, pid, text]

with beam.Pipeline() as pipeline:
    log_lines = pipeline | 'Read log line' >> beam.io.ReadFromText("/Analytics/venv/Jup/CAPE_Apache_Beam/Sample_text_file")
    print(log_lines)

    fields = log_lines | 'Extract fields' >> Map(extract_fields)
    print(fields)

    df = (fields | 'Write to dataframe' >> Map(lambda fields: pd.DataFrame(fields, columns=['timestamp', 'hostname', 'process_name', 'pid', 'text'])))
    print(df)
解决方案

问题分析

原代码存在几个核心问题:

  • 缺少关键词过滤逻辑,未筛选目标日志行
  • 直接在Map中创建DataFrame会生成大量小数据框,无法得到完整的结果集
  • 字段分割方式不够健壮,若日志后续字段含空格会导致分割错误

修正后的代码

import apache_beam as beam
from apache_beam.io import ReadFromText
import pandas as pd

# 定义需要匹配的关键词集合
TARGET_KEYWORDS = {'unauthorized', 'error', 'kernel error', 'OS error', 'rejected', 'warning'}

def filter_target_logs(line):
    # 不区分大小写检查日志是否包含目标关键词
    line_lower = line.lower()
    return any(keyword.lower() in line_lower for keyword in TARGET_KEYWORDS)

def parse_log_fields(line):
    # 按空格分割,仅分割前3次,避免后续文本含空格导致错误
    split_parts = line.split(maxsplit=3)
    timestamp = split_parts[0]
    hostname = split_parts[1]
    
    # 解析进程名和PID
    process_info = split_parts[2].split('[')
    process_name = process_info[0]
    pid = process_info[1].rstrip(']:')  # 清理PID后的多余符号
    
    log_text = split_parts[3]
    # 返回字典格式,与DataFrame列名直接对应
    return {
        'timestamp': timestamp,
        'hostname': hostname,
        'process_name': process_name,
        'pid': pid,
        'text': log_text
    }

with beam.Pipeline() as pipeline:
    # 1. 读取日志文件
    raw_logs = pipeline | '读取日志文件' >> ReadFromText("/Analytics/venv/Jup/CAPE_Apache_Beam/Sample_text_file")
    
    # 2. 过滤符合关键词条件的日志
    filtered_logs = raw_logs | '过滤目标日志' >> beam.Filter(filter_target_logs)
    
    # 3. 解析日志字段为字典
    parsed_fields = filtered_logs | '解析日志字段' >> beam.Map(parse_log_fields)
    
    # 4. 收集所有结果并转换为DataFrame
    collected_data = parsed_fields | '收集结果列表' >> beam.combiners.ToList()
    final_df = collected_data | '转换为DataFrame' >> beam.Map(lambda data: pd.DataFrame(data))
    
    # 打印结果(可替换为保存到文件等操作)
    final_df | '输出DataFrame' >> beam.Map(lambda df: print(df.head()))

关键改进点

  • 新增关键词过滤逻辑,支持不区分大小写匹配,确保筛选准确
  • 使用split(maxsplit=3)分割日志,避免后续文本含空格导致的字段混乱
  • 字段解析返回字典格式,与DataFrame列名直接对应,逻辑更清晰
  • 通过ToList()收集所有结果后再创建DataFrame,确保生成完整的数据集
  • 优化PID字段的清理逻辑,适配日志格式细节

内容的提问来源于stack exchange,提问作者samh125679

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.31 12:46:15