Skip to content

Commit edf2370

Browse files
veloferenc-csaky
andauthored
fix: log uncaught exceptions from CliRunner main so JobManager pod logs capture them (#342)
Co-authored-by: Ferenc Csaky <ferenc@datasqrl.com>
1 parent 0a0cf92 commit edf2370

3 files changed

Lines changed: 38 additions & 4 deletions

File tree

flink-sql-runner/src/main/java/com/datasqrl/flinkrunner/CliRunner.java

Lines changed: 14 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -90,7 +90,7 @@ public CliRunner(
9090
}
9191

9292
public static void main(String[] args) throws Exception {
93-
log.info("Executing flink-sql-runner: {}", Arrays.toString(args));
93+
log.info("Starting Flink SQL runner: {}", Arrays.toString(args));
9494

9595
var cmd = new CommandLine(new SqlRunner());
9696
cmd.setUnmatchedArgumentsAllowed(true);
@@ -111,9 +111,19 @@ public static void main(String[] args) throws Exception {
111111
runner.udfPath = System.getenv("UDF_PATH");
112112
}
113113

114-
new CliRunner(runner.mode, runner.sqlFile, runner.planFile, runner.configDir, runner.udfPath)
115-
.run();
114+
try {
115+
var cliRunner =
116+
new CliRunner(
117+
runner.mode, runner.sqlFile, runner.planFile, runner.configDir, runner.udfPath);
116118

117-
log.info("Finished flink-sql-runner execution");
119+
cliRunner.run();
120+
121+
} catch (Throwable t) {
122+
// Make sure we log any error to be able to present it in a K8s env
123+
log.error("Flink SQL runner failed", t);
124+
throw t;
125+
}
126+
127+
log.info("Flink SQL runner finished");
118128
}
119129
}

flink-sql-runner/src/test/java/com/datasqrl/flinkrunner/CliRunnerIT.java

Lines changed: 23 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -15,6 +15,8 @@
1515
*/
1616
package com.datasqrl.flinkrunner;
1717

18+
import static org.assertj.core.api.Assertions.assertThat;
19+
1820
import java.util.ArrayList;
1921
import java.util.stream.Stream;
2022
import org.junit.jupiter.api.Test;
@@ -71,4 +73,25 @@ void givenKafkaPlanScript_whenExecuting_thenSuccess() throws Exception {
7173
String jobId = flinkRun("--planfile", "/it/planfile/kafka_plan.json");
7274
assertJobIsRunning(jobId);
7375
}
76+
77+
@Test
78+
void givenFailingSqlScript_whenExecuting_thenLogsFailure() throws Exception {
79+
var execRes =
80+
flinkContainer.execInContainer("sql-runner", "--sqlfile", "/it/sqlfile/dummy_error.sql");
81+
var commandOutput = execRes.getStdout() + execRes.getStderr();
82+
83+
assertThat(execRes.getExitCode()).isNotZero();
84+
85+
untilAssert(
86+
() -> {
87+
var logRes =
88+
flinkContainer.execInContainer(
89+
"bash",
90+
"-c",
91+
"grep -R 'Flink SQL runner failed\\|DUMMY_ERROR' /opt/flink/log || true");
92+
var loggedOutput = commandOutput + logRes.getStdout() + logRes.getStderr();
93+
94+
assertThat(loggedOutput).contains("Flink SQL runner failed").contains("DUMMY_ERROR");
95+
});
96+
}
7497
}
Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1 @@
1+
SELECT DUMMY_ERROR();

0 commit comments

Comments
 (0)