使用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
相关产品推荐
相关产品推荐

