High-Performance CSV Processor
This project demonstrates deep understanding of Java by processing massive CSV files (gigabytes of logs) using highly efficient, parallel Streams and NIO lazy loading, completely avoiding OutOfMemoryErrors.
public class LogProcessor {
public Map<String, Long> analyzeErrorsByService(Path filePath) throws IOException {
try (Stream<String> lines = Files.lines(filePath, StandardCharsets.UTF_8)) {
return lines
.skip(1)
.parallel()
.map(LogEntry::parseCsvLine)
.filter(Objects::nonNull)
.filter(entry -> "ERROR".equalsIgnoreCase(entry.level()))
.collect(Collectors.groupingByConcurrent(
LogEntry::serviceName,
Collectors.counting()
));
}
}
}
public record LogEntry(
LocalDateTime timestamp,
String level,
String serviceName,
String message
) {
public static LogEntry parseCsvLine(String line) {
if (line == null || line.isBlank()) {
return null;
}
String[] parts = line.split(",", 4);
if (parts.length != 4) {
return null;
}
try {
return new LogEntry(
LocalDateTime.parse(parts[0].trim()),
parts[1].trim(),
parts[2].trim(),
parts[3].trim()
);
} catch (Exception e) {
return null;
}
}
}
public class Main {
public static void main(String[] args) {
Path logFilePath = Paths.get("data", "server_logs_large.csv");
LogProcessor processor = new LogProcessor();
try {
Instant start = Instant.now();
Map<String, Long> errorAggregates = processor.analyzeErrorsByService(logFilePath);
Instant finish = Instant.now();
long timeElapsed = Duration.between(start, finish).toMillis();
System.out.println("Processing completed in " + timeElapsed + " ms.");
errorAggregates.entrySet().stream()
.sorted(Map.Entry.<String, Long>comparingByValue().reversed())
.forEach(entry -> System.out.printf("%-20s : %d errors%n", entry.getKey(), entry.getValue()));
} catch (Exception e) {
System.err.println("Fatal error: " + e.getMessage());
}
}
}
Click "RUN" to simulate the execution of the DataStream Engine processing a 4.5 GB CSV log file in real time.