Skip to main content

lean_ctx/core/patterns/
spark.rs

1//! Apache Spark (`spark-submit`) log compression.
2//!
3//! Spark drivers emit hundreds of `YY/MM/DD HH:MM:SS INFO Component: ...`
4//! lines. We drop INFO noise, keep finished-job lines, deduplicate WARNs and
5//! preserve ERRORs / exceptions, then prefix a one-line job/warn/error count.
6//!
7//! Crucially we also keep lines that are NOT framework log records — these are
8//! the application's own stdout (e.g. `Result: total words = 184273`), which is
9//! the actual point of the run and must never be dropped.
10
11use crate::core::compressor::strip_ansi;
12use std::collections::HashSet;
13
14pub fn compress(_cmd: &str, output: &str) -> Option<String> {
15    let trimmed = output.trim();
16    if trimmed.is_empty() {
17        return Some("spark: ok".to_string());
18    }
19
20    let mut jobs: Vec<String> = Vec::new();
21    let mut warnings: Vec<String> = Vec::new();
22    let mut warn_seen: HashSet<String> = HashSet::new();
23    let mut errors: Vec<String> = Vec::new();
24    let mut app_output: Vec<String> = Vec::new();
25    let mut saw_log = false;
26
27    for raw in trimmed.lines() {
28        let stripped = strip_ansi(raw);
29        let line = stripped.trim();
30        if line.is_empty() {
31            continue;
32        }
33
34        if let Some((level, rest)) = parse_log(line) {
35            saw_log = true;
36            match level {
37                "INFO" => {
38                    if rest.contains("Job ") && rest.contains("finished") {
39                        jobs.push(rest.to_string());
40                    }
41                }
42                "WARN" => {
43                    if warn_seen.insert(normalize(rest)) {
44                        warnings.push(rest.to_string());
45                    }
46                }
47                "ERROR" => errors.push(rest.to_string()),
48                _ => {}
49            }
50        } else if is_exception(line) {
51            errors.push(line.to_string());
52        } else {
53            // Not a framework log line → application stdout. Preserve it.
54            app_output.push(line.to_string());
55        }
56    }
57
58    if !saw_log && jobs.is_empty() && errors.is_empty() {
59        return Some(fallback(trimmed));
60    }
61
62    let mut parts = vec![format!(
63        "spark: {} job(s), {} warning(s), {} error(s)",
64        jobs.len(),
65        warnings.len(),
66        errors.len()
67    )];
68    push_capped(&mut parts, &jobs, 10, "more jobs");
69    push_capped(&mut parts, &warnings, 5, "more warnings");
70    push_capped(&mut parts, &errors, 10, "more errors");
71    push_capped(&mut parts, &app_output, 20, "more output lines");
72    Some(parts.join("\n"))
73}
74
75/// Parse a `DATE TIME LEVEL rest` Spark log line into `(level, rest)`.
76fn parse_log(line: &str) -> Option<(&str, &str)> {
77    let mut it = line.splitn(4, ' ');
78    let date = it.next()?;
79    let time = it.next()?;
80    let level = it.next()?;
81    let rest = it.next().unwrap_or("");
82    if looks_like_date(date) && looks_like_time(time) && is_level(level) {
83        Some((level, rest))
84    } else {
85        None
86    }
87}
88
89fn looks_like_date(s: &str) -> bool {
90    let parts: Vec<&str> = s.split('/').collect();
91    parts.len() == 3 && parts.iter().all(|p| p.chars().all(|c| c.is_ascii_digit()))
92}
93
94fn looks_like_time(s: &str) -> bool {
95    let parts: Vec<&str> = s.split(':').collect();
96    parts.len() == 3 && parts.iter().all(|p| p.chars().all(|c| c.is_ascii_digit()))
97}
98
99fn is_level(s: &str) -> bool {
100    matches!(s, "INFO" | "WARN" | "ERROR" | "DEBUG" | "TRACE")
101}
102
103fn is_exception(line: &str) -> bool {
104    line.contains("Exception") || line.starts_with("Caused by:") || line.contains("Error:")
105}
106
107/// Collapse digits so "took 5.1 s" / "took 9.2 s" warnings dedupe together.
108fn normalize(s: &str) -> String {
109    s.chars().filter(|c| !c.is_ascii_digit()).collect()
110}
111
112fn push_capped(parts: &mut Vec<String>, items: &[String], cap: usize, label: &str) {
113    for item in items.iter().take(cap) {
114        parts.push(format!("  {item}"));
115    }
116    if items.len() > cap {
117        parts.push(format!("  ... +{} {label}", items.len() - cap));
118    }
119}
120
121fn fallback(text: &str) -> String {
122    let lines: Vec<&str> = text.lines().filter(|l| !l.trim().is_empty()).collect();
123    let n = lines.len().min(8);
124    let mut s = lines[..n].join("\n");
125    if lines.len() > n {
126        s.push_str(&format!("\n... (+{} lines)", lines.len() - n));
127    }
128    s
129}
130
131#[cfg(test)]
132mod tests {
133    use super::*;
134
135    const LOG: &str = "23/01/01 12:00:00 INFO SparkContext: Running Spark version 3.4.0\n23/01/01 12:00:01 INFO ResourceUtils: No custom resources configured\n23/01/01 12:00:02 INFO Utils: Successfully started service\n23/01/01 12:00:03 WARN NativeCodeLoader: Unable to load native-hadoop\n23/01/01 12:00:10 INFO DAGScheduler: Job 0 finished: collect at Main.scala:20, took 5.123 s\n23/01/01 12:00:15 ERROR Executor: Exception in task 0.0 in stage 1.0\n";
136
137    #[test]
138    fn drops_info_keeps_job_warn_error() {
139        let r = compress("spark-submit app.py", LOG).unwrap();
140        assert!(r.contains("1 job(s), 1 warning(s), 1 error(s)"), "{r}");
141        assert!(r.contains("Job 0 finished"), "{r}");
142        assert!(r.contains("Executor: Exception"), "{r}");
143        assert!(!r.contains("ResourceUtils"), "drops info noise: {r}");
144        assert!(!r.contains("23/01/01"), "drops timestamps: {r}");
145    }
146
147    #[test]
148    fn keeps_application_stdout() {
149        let log = "23/01/01 12:00:00 INFO SparkContext: Running Spark version 3.4.0\n23/01/01 12:00:01 INFO ResourceUtils: No custom resources\n23/01/01 12:00:10 INFO DAGScheduler: Job 0 finished: collect, took 5.1 s\nResult: total words = 184273\n23/01/01 12:00:11 INFO SparkContext: Successfully stopped";
150        let r = compress("spark-submit app.py", log).unwrap();
151        assert!(
152            r.contains("Result: total words = 184273"),
153            "keeps app output: {r}"
154        );
155        assert!(!r.contains("ResourceUtils"), "still drops info noise: {r}");
156    }
157
158    #[test]
159    fn shorter_than_input() {
160        let r = compress("spark-submit app.py", LOG).unwrap();
161        assert!(r.len() < LOG.len());
162    }
163
164    #[test]
165    fn empty_is_ok() {
166        assert_eq!(compress("spark-submit app.py", "").unwrap(), "spark: ok");
167    }
168}