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

如何使用Spark读取并解析自定义多行日志

使用Spark结合正则解析自定义多行日志

嘿,我来帮你梳理下怎么搞定这个自定义多行日志的解析问题!首先得明确你的日志特点:每一条完整的日志是4行一组——第一行是核心动作(received/sending),后面三行是对应的sec、Msg-Type、Content属性。Spark默认是按单行读取的,所以第一步必须先把这4行合并成一个完整的日志条目,之后再用正则解析就好办了。

第一步:合并多行日志

我们可以通过mapPartitions来迭代处理每个分区的日志行,把属于同一组的4行拼接成一个字符串。核心逻辑是:识别出包含received或sending的主日志行,把后续的属性行都附加到这行后面,直到遇到下一个主日志行为止。

val logLines = sc.textFile("/directory/logfile")

// 定义分区内的多行合并逻辑
def mergeRelatedLogs(iter: Iterator[String]): Iterator[String] = {
  var currentEntry = new StringBuilder()
  iter.flatMap { line =>
    // 判断当前行是否是主日志行(核心动作行)
    if (line.contains("received") || line.contains("sending")) {
      // 如果之前已经在构建一个日志条目,先把它输出
      val previousEntry = if (currentEntry.nonEmpty) Some(currentEntry.toString()) else None
      // 开始构建新的日志条目
      currentEntry = new StringBuilder(line)
      previousEntry
    } else {
      // 属性行,直接追加到当前条目后面
      currentEntry.append(" ").append(line.trim)
      None
    }
  } ++ (if (currentEntry.nonEmpty) Iterator(currentEntry.toString()) else Iterator.empty)
}

// 得到合并后的单条完整日志
val mergedLogEntries = logLines.mapPartitions(mergeRelatedLogs)

第二步:编写正则并解析日志

现在每条日志都是一个完整的字符串了,我们可以针对两种动作(received/sending)分别编写正则,再用Scala的模式匹配把内容提取到你定义的case class里。我调整了你原来的正则,让它能匹配合并后的完整字符串:

// 定义统一的日志实体类(你也可以用原来的Rlog和Slog分开处理)
case class Log(
  dateTime: String,
  serverName: String,
  logType: String,
  serverAddr: String,
  action: String,
  targetAddr: String,
  cost: String,
  msgType: String,
  content: String
)

// 匹配received类型的完整日志正则
val receivePattern = """([\d\-:\s]{19})\s([a-z]+)\s(\w+):\s(\w+)\sreceived\s(\w+)\sfrom\s([\d\.]+)\s.*sec\s=\s([\d\.]+)\s.*Msg-Type\s=\s(\w+)\s.*Content\s=\s(.+)""".r
// 匹配sending类型的完整日志正则
val sendPattern = """([\d\-:\s]{19})\s([a-z]+)\s(\w+):\s(\w+)\ssending\s(\w+)\sto\s([\d\.]+)\s.*sec\s=\s([\d\.]+)\s.*Msg-Type\s=\s(\w+)\s.*Content\s=\s(.+)""".r

// 解析合并后的日志条目
val parsedLogs = mergedLogEntries.flatMap { entry =>
  entry match {
    case receivePattern(dt, server, logType, srvAddr, action, target, cost, msgType, content) =>
      Some(Log(dt, server, logType, srvAddr, action, target, cost, msgType, content))
    case sendPattern(dt, server, logType, srvAddr, action, target, cost, msgType, content) =>
      Some(Log(dt, server, logType, srvAddr, action, target, cost, msgType, content))
    case _ => None // 过滤掉无法匹配的异常行
  }
}

// 转成DataFrame方便后续分析(如果用Spark SQL的话)
import spark.implicits._
val logDF = parsedLogs.toDF()
logDF.show(truncate = false)

几个关键注意点

  • 正则调整:我把时间戳的匹配改成了[\d\-:\s]{19},比你原来的[\w|\s|\:|-]更精准,因为你的时间戳只有数字、横杠、冒号和空格。
  • Content匹配:如果你的Content包含空格,正则里的(.+)会匹配到行尾,刚好符合需求;如果是无空格的内容,用(\S+)也可以。
  • 合并逻辑:如果你的日志里同一秒可能有多个不同的主日志行(比如同时有received和sending),这个合并逻辑依然有效,因为它是按主行出现的顺序来分组的。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.28 06:30:32