Flink InvalidProgramException: Job was submitted in detached mode. Results of job execution, such as accumulators, runtime, etc. are not available

Viewed 345

I have a Flink job running on a Kinesis Data Analytics application, which uses Flink's DataSet API to read data into two DataSet objects. Since I need the number of tuples in each DataSet, I call the count() method on each DataSet, but I keep seeing this error when run through my AWS Console:

org.apache.flink.api.common.InvalidProgramException: The main method caused an error: Job was submitted in detached mode. Results of job execution, such as accumulators, runtime, etc. are not available. Please make sure your program doesn't call an eager execution function [collect, print, printToErr, count].

For context, this is roughly the code that is causing the exception:

DataSet<Tuple> dataset = executionEnvironment.readTextFile(file);
log.info("Number of records: " + dataset.count());

Is there any way to change the execution mode from detached mode to another mode that would allow the call to the count and other accumulator functions?

0 Answers
Related