diff options
author | seanm <sean.mcnamara@webtrends.com> | 2013-05-10 17:34:28 -0600 |
---|---|---|
committer | seanm <sean.mcnamara@webtrends.com> | 2013-05-10 17:34:28 -0600 |
commit | f25282def5826fab6caabff28c82c57a7f3fdcb8 (patch) | |
tree | a35530c4fa845eed356026cc35a373bc8922530c /streaming/src/test | |
parent | 3632980b1b61dbb9ab9a3ab3d92fb415cb7173b9 (diff) | |
download | spark-f25282def5826fab6caabff28c82c57a7f3fdcb8.tar.gz spark-f25282def5826fab6caabff28c82c57a7f3fdcb8.tar.bz2 spark-f25282def5826fab6caabff28c82c57a7f3fdcb8.zip |
fixing kafkaStream Java API and adding test
Diffstat (limited to 'streaming/src/test')
-rw-r--r-- | streaming/src/test/java/spark/streaming/JavaAPISuite.java | 6 |
1 files changed, 6 insertions, 0 deletions
diff --git a/streaming/src/test/java/spark/streaming/JavaAPISuite.java b/streaming/src/test/java/spark/streaming/JavaAPISuite.java index 350d0888a3..e5fdbe1b7a 100644 --- a/streaming/src/test/java/spark/streaming/JavaAPISuite.java +++ b/streaming/src/test/java/spark/streaming/JavaAPISuite.java @@ -1206,6 +1206,12 @@ public class JavaAPISuite implements Serializable { JavaDStream test1 = ssc.kafkaStream("localhost:12345", "group", topics); JavaDStream test2 = ssc.kafkaStream("localhost:12345", "group", topics, StorageLevel.MEMORY_AND_DISK()); + + HashMap<String, String> kafkaParams = Maps.newHashMap(); + kafkaParams.put("zk.connect","localhost:12345"); + kafkaParams.put("groupid","consumer-group"); + JavaDStream test3 = ssc.kafkaStream(String.class, StringDecoder.class, kafkaParams, topics, + StorageLevel.MEMORY_AND_DISK()); } @Test |