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

