001/** 002 * Licensed to the Apache Software Foundation (ASF) under one 003 * or more contributor license agreements. See the NOTICE file 004 * distributed with this work for additional information 005 * regarding copyright ownership. The ASF licenses this file 006 * to you under the Apache License, Version 2.0 (the 007 * "License"); you may not use this file except in compliance 008 * with the License. You may obtain a copy of the License at 009 * 010 * http://www.apache.org/licenses/LICENSE-2.0 011 * 012 * Unless required by applicable law or agreed to in writing, software 013 * distributed under the License is distributed on an "AS IS" BASIS, 014 * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. 015 * See the License for the specific language governing permissions and 016 * limitations under the License. 017 */ 018 019package org.apache.oozie.util; 020 021import java.io.BufferedReader; 022import java.io.IOException; 023import java.io.Writer; 024import java.util.ArrayList; 025import java.util.regex.Pattern; 026 027import org.apache.oozie.service.Services; 028import org.apache.oozie.service.XLogStreamingService; 029import org.apache.oozie.util.LogLine.MATCHED_PATTERN; 030 031 /** 032 * Encapsulates the parsing and filtering of the log messages from a BufferedReader. It makes sure not to read in the entire log 033 * into memory at the same time; at most, it will have two messages (which can be multi-line in the case of exception stack traces). 034 * <p> 035 * To use this class: Calling {@link TimestampedMessageParser#increment()} will tell the parser to read the next message from the 036 * Reader. It will return true if there are more messages and false if not. Calling 037 * {@link TimestampedMessageParser#getLastMessage()} and {@link TimestampedMessageParser#getLastTimestamp()} will return the last 038 * message and timestamp, respectively, that were parsed when {@link TimestampedMessageParser#increment()} was called. Calling 039 * {@link TimestampedMessageParser#processRemaining(java.io.Writer)} will write the remaining log messages to the given Writer. 040 */ 041public class TimestampedMessageParser { 042 043 static final String SYSTEM_LINE_SEPARATOR = System.getProperty("line.separator"); 044 protected BufferedReader reader; 045 private LogLine nextLine = null; 046 private String lastTimestamp = null; 047 private XLogFilter filter; 048 private boolean empty = false; 049 private String lastMessage = null; 050 private boolean patternMatched = false; 051 public int count = 0; 052 private Pattern splitPattern = null; 053 054 /** 055 * Creates a TimestampedMessageParser with the given BufferedReader and filter. 056 * 057 * @param reader The BufferedReader to get the log messages from 058 * @param filter The filter 059 */ 060 public TimestampedMessageParser(BufferedReader reader, XLogFilter filter) { 061 this.reader = reader; 062 this.filter = filter; 063 if (filter == null) { 064 filter = new XLogFilter(); 065 } 066 filter.constructPattern(); 067 String regEx = XLogFilter.PREFIX_REGEX + filter.getFilterPattern().pattern(); 068 this.splitPattern = Pattern.compile(regEx); 069 } 070 071 /** 072 * Causes the next message and timestamp to be parsed from the BufferedReader. 073 * 074 * @return true if there are more messages left; false if not 075 * @throws IOException If there was a problem reading from the Reader 076 */ 077 public boolean increment() throws IOException { 078 if (empty) { 079 return false; 080 } 081 082 StringBuilder message = new StringBuilder(); 083 if (nextLine == null) { // first time only 084 nextLine = parseNextLogLine(); 085 if (nextLine == null || nextLine.getLine() == null) { 086 // reader finished 087 empty = true; 088 return false; 089 } 090 } 091 lastTimestamp = parseTimestamp(nextLine); 092 String nextTimestamp = null; 093 while (nextTimestamp == null) { 094 message.append(nextLine.getLine()).append(SYSTEM_LINE_SEPARATOR); 095 nextLine = parseNextLogLine(); 096 if (nextLine != null && nextLine.getLine() != null) { 097 // exit loop if we have a timestamp, continue if not 098 nextTimestamp = parseTimestamp(nextLine); 099 } 100 else { // reader finished 101 empty = true; 102 nextTimestamp = ""; // exit loop 103 } 104 } 105 106 lastMessage = message.toString(); 107 return true; 108 } 109 110 /** 111 * Returns the timestamp from the last message that was parsed. 112 * 113 * @return the timestamp from the last message that was parsed 114 */ 115 public String getLastTimestamp() { 116 return lastTimestamp; 117 } 118 119 /** 120 * Returns the message that was last parsed. 121 * 122 * @return the message that was last parsed 123 */ 124 public String getLastMessage() { 125 return lastMessage; 126 } 127 128 /** 129 * Closes the Reader. 130 * 131 * @throws IOException 132 */ 133 public void closeReader() throws IOException { 134 reader.close(); 135 } 136 137 /** 138 * Reads the next line from the Reader and checks if it matches the filter. 139 * It can also handle multi-line messages (i.e. exception stack traces). If 140 * it returns null, then there are no lines left in the Reader. 141 * 142 * @return LogLine 143 * @throws IOException 144 */ 145 protected LogLine parseNextLogLine() throws IOException { 146 String line; 147 LogLine logLine = new LogLine(); 148 while ((line = reader.readLine()) != null) { 149 logLine.setLine(line); 150 logLine.setLogParts(null); 151 filter.splitLogMessage(logLine, splitPattern); 152 // check the splits if logLine matches with the splitPattern 153 // Otherwise, go with previous patternMatched value. This is needed 154 // in parsing stack trace 155 patternMatched = logLine.getMatchedPattern() == MATCHED_PATTERN.NONE ? patternMatched 156 : filter.splitsMatches(logLine); 157 if (patternMatched) { 158 if (filter.getLogLimit() != -1) { 159 if (logLine.getLogParts() != null) { 160 if (count >= filter.getLogLimit()) { 161 return null; 162 } 163 count++; 164 } 165 } 166 if (logLine.getLogParts() != null) { 167 if (filter.getEndDate() != null) { 168 // Ignore the milli second part 169 if (logLine.getLogParts().get(0).substring(0, 19).compareTo(filter.getFormattedEndDate()) > 0) 170 return null; 171 } 172 } 173 return logLine; 174 } 175 } 176 logLine.setLine(null); 177 return logLine; 178 } 179 180 /** 181 * Parses the timestamp out of the passed in line. If there isn't one, it 182 * returns null. 183 * 184 * @param logLine The LogLine to check 185 * @return the timestamp of the line, or null 186 */ 187 private String parseTimestamp(LogLine logLine) { 188 String timestamp = null; 189 if (logLine != null && logLine.getLogParts() != null && logLine.getLogParts().size() > 0) { 190 timestamp = logLine.getLogParts().get(0); 191 } 192 return timestamp; 193 } 194 195 /** 196 * Streams log messages to the passed in Writer. Flushes the log writing 197 * based on buffer len 198 * 199 * @param writer 200 * @param bufferLen maximum len of log buffer 201 * @param bytesWritten num bytes already written to writer 202 * @throws IOException 203 */ 204 public void processRemaining(Writer writer, int bufferLen, int bytesWritten) throws IOException { 205 while (increment()) { 206 writer.write(lastMessage); 207 bytesWritten += lastMessage.length(); 208 if (bytesWritten > bufferLen) { 209 writer.flush(); 210 bytesWritten = 0; 211 } 212 } 213 writer.flush(); 214 } 215 216 /** 217 * Streams log messages to the passed in Writer, with zero bytes already 218 * written 219 * 220 * @param writer 221 * @param bufferLen maximum len of log buffer 222 * @throws IOException 223 */ 224 public void processRemaining(Writer writer, int bufferLen) throws IOException { 225 processRemaining(writer, bufferLen, 0); 226 } 227 228 /** 229 * Streams log messages to the passed in Writer, with default buffer len 4K 230 * and zero bytes already written 231 * 232 * @param writer 233 * @throws IOException 234 */ 235 public void processRemaining(Writer writer) throws IOException { 236 processRemaining(writer, Services.get().get(XLogStreamingService.class).getBufferLen()); 237 } 238 239 /** 240 * Splits the log message into parts 241 * 242 * @param line 243 * @return List of log parts 244 */ 245 protected ArrayList<String> splitLogMessage(String line) { 246 return filter.splitLogMessage(line); 247 } 248 249 250}