KafkaBolt no longer tries to map/process/send Tick Tuples to Kafka.
Project: http://git-wip-us.apache.org/repos/asf/storm/repo Commit: http://git-wip-us.apache.org/repos/asf/storm/commit/73e54a8f Tree: http://git-wip-us.apache.org/repos/asf/storm/tree/73e54a8f Diff: http://git-wip-us.apache.org/repos/asf/storm/diff/73e54a8f Branch: refs/heads/master Commit: 73e54a8f8b0ef9a62955c7c1cee20925d359772b Parents: 6c6bec9 Author: Niels Basjes <[email protected]> Authored: Wed Oct 1 11:51:42 2014 +0200 Committer: Niels Basjes <[email protected]> Committed: Wed Oct 1 11:51:42 2014 +0200 ---------------------------------------------------------------------- external/storm-kafka/src/jvm/storm/kafka/bolt/KafkaBolt.java | 4 ++++ 1 file changed, 4 insertions(+) ---------------------------------------------------------------------- http://git-wip-us.apache.org/repos/asf/storm/blob/73e54a8f/external/storm-kafka/src/jvm/storm/kafka/bolt/KafkaBolt.java ---------------------------------------------------------------------- diff --git a/external/storm-kafka/src/jvm/storm/kafka/bolt/KafkaBolt.java b/external/storm-kafka/src/jvm/storm/kafka/bolt/KafkaBolt.java index b6c3de4..2025766 100644 --- a/external/storm-kafka/src/jvm/storm/kafka/bolt/KafkaBolt.java +++ b/external/storm-kafka/src/jvm/storm/kafka/bolt/KafkaBolt.java @@ -89,6 +89,10 @@ public class KafkaBolt<K, V> extends BaseRichBolt { @Override public void execute(Tuple input) { + if (input.isTick()) { + return; // Do not try to send ticks to Kafka + } + K key = null; V message = null; String topic = null;
