diff --git a/wayang-platforms/wayang-jdbc-template/src/main/java/org/apache/wayang/jdbc/execution/JdbcExecutor.java b/wayang-platforms/wayang-jdbc-template/src/main/java/org/apache/wayang/jdbc/execution/JdbcExecutor.java index 4b816b2eb..344f89ea8 100644 --- a/wayang-platforms/wayang-jdbc-template/src/main/java/org/apache/wayang/jdbc/execution/JdbcExecutor.java +++ b/wayang-platforms/wayang-jdbc-template/src/main/java/org/apache/wayang/jdbc/execution/JdbcExecutor.java @@ -371,14 +371,29 @@ private static ExecutionTask findJdbcExecutionOperatorTaskInStage(final Executio final Channel outputChannel = task.getOutputChannel(0); - if (outputChannel.getConsumers().size() != 1) { + final Collection consumers = outputChannel.getConsumers(); + + // No consumer → end of pipeline + if (consumers.isEmpty()) { return null; } - final ExecutionTask consumer = outputChannel.getConsumers().iterator().next(); + // Multiple consumers are currently unsupported + if (consumers.size() != 1) { + throw new WayangException( + String.format("Expected a single consumer for task %s but found %d.", + task, + consumers.size() + ) + ); + } + + final ExecutionTask consumer = consumers.iterator().next(); - return consumer.getStage() == stage && consumer.getOperator() instanceof JdbcExecutionOperator ? consumer - : null; + return consumer.getStage() == stage + && consumer.getOperator() instanceof JdbcExecutionOperator + ? consumer + : null; } private static SqlQueryChannel.Instance instantiateOutboundChannel(final ExecutionTask task,