【发布时间】:2016-01-22 23:29:02
【问题描述】:
我正在使用以下代码来读取我的主题数据,即“sha-test2”,但它正在读取完全替代的代码行,即 20 行中的 10 行。 但是当我运行控制台时,它显示了所有 20 行。 IE 。 bin/kafka-console-consumer.sh --zookeeper localhost:2181 --topic sha-test2 --from-beginning
我哪里错了?非常感谢您的帮助。
public class KafkaTestConsumer extends Thread {
//final static String clientId = "SimpleConsumerDemoClient";
final static String TOPIC = "sha-test2";
ConsumerConnector consumerConnector;
public static void main(String[] argv) throws
UnsupportedEncodingException {
KafkaTestConsumer helloKafkaConsumer = new KafkaTestConsumer();
helloKafkaConsumer.start();
}
public KafkaTestConsumer(){
Properties properties = new Properties();
properties.put("zookeeper.connect","172.23.32.35:2181");
properties.put("group.id","test-group");
ConsumerConfig consumerConfig = new ConsumerConfig(properties);
consumerConnector =
Consumer.createJavaConsumerConnector(consumerConfig);
}
@Override
public void run() {
Map<String, Integer> topicCountMap = new HashMap<String, Integer>();
topicCountMap.put(TOPIC, new Integer(1));
Map<String, List<KafkaStream<byte[], byte[]>>> consumerMap =
consumerConnector.createMessageStreams(topicCountMap);
KafkaStream<byte[], byte[]> stream = consumerMap.get(TOPIC).get(0);
System.out.println("consumerMap : \n " + consumerMap.toString() );
ConsumerIterator<byte[], byte[]> it = stream.iterator();
System.out.println("run started");
while(it.hasNext()){
System.out.println(new String(it.next().message()));
}
}
Thank you.
~Shyam
【问题讨论】:
标签: regex apache-kafka hadoop-streaming kafka-consumer-api