diff --git a/storm-core/src/jvm/backtype/storm/spout/ShellSpout.java b/storm-core/src/jvm/backtype/storm/spout/ShellSpout.java index ba7bb63a8b..8391383e4e 100644 --- a/storm-core/src/jvm/backtype/storm/spout/ShellSpout.java +++ b/storm-core/src/jvm/backtype/storm/spout/ShellSpout.java @@ -88,6 +88,10 @@ private void querySubprocess(Object query) { } else if (command.equals("log")) { String msg = (String) action.get("msg"); LOG.info("Shell msg: " + msg); + } else if (command.equals("error")) { + String msg = (String) action.get("msg"); + _collector.reportError(new Exception("Shell Process Exception: " + msg)); + return; } else if (command.equals("emit")) { String stream = (String) action.get("stream"); if (stream == null) stream = Utils.DEFAULT_STREAM_ID;