【发布时间】:2018-02-09 11:32:08
【问题描述】:
我正在尝试集成 spark 和 Kafka 以使用来自 Kafka 的消息。我也有生产者代码来发送关于“临时”主题的消息。另外,我正在使用 Kafka 的控制台生产者来生产关于“临时”主题的消息。
我创建了下面的代码来使用来自同一“临时”主题的消息,但它也不会收到一条消息。
计划:
import java.util.Arrays;
import java.util.Map;
import java.util.HashMap;
import static org.apache.commons.lang3.StringUtils.SPACE;
import org.apache.spark.SparkConf;
import org.apache.spark.streaming.Duration;
import org.apache.spark.streaming.api.java.JavaDStream;
import org.apache.spark.streaming.api.java.JavaPairDStream;
import org.apache.spark.streaming.api.java.JavaPairReceiverInputDStream;
import org.apache.spark.streaming.api.java.JavaStreamingContext;
import org.apache.spark.streaming.kafka.KafkaUtils;
import scala.Tuple2;
import org.apache.log4j.Logger;
import org.apache.spark.api.java.JavaSparkContext;
import scala.collection.immutable.ListSet;
import scala.collection.immutable.Set;
public class ConsumerDemo {
public void main() {
String zkGroup = "localhost:2181";
String group = "test";
String[] topics = {"temp"};
int numThreads = 1;
SparkConf sparkConf = new SparkConf().setAppName("JavaKafkaWordCount").setMaster("local[4]").set("spark.ui.port", "7077").set("spark.executor.memory", "1g");
JavaStreamingContext jssc = new JavaStreamingContext(sparkConf, new Duration(2000));
Map<String, Integer> topicMap = new HashMap<>();
for (String topic : topics) {
topicMap.put(topic, numThreads);
}
System.out.println("topics : " + Arrays.toString(topics));
JavaPairReceiverInputDStream<String, String> messages
= KafkaUtils.createStream(jssc, zkGroup, group, topicMap);
messages.print();
JavaDStream<String> lines = messages.map(Tuple2::_2);
//lines.print();
JavaDStream<String> words = lines.flatMap(x -> Arrays.asList(SPACE.split(x)).iterator());
JavaPairDStream<String, Integer> wordCounts = words.mapToPair(s -> new Tuple2<>(s, 1))
.reduceByKey((i1, i2) -> i1 + i2);
//wordCounts.print();
jssc.start();
jssc.awaitTermination();
}
public static void main(String[] args) {
System.out.println("Started...");
new ConsumerDemo().main();
System.out.println("Ended...");
}
}
我在 pom.xml 文件中添加了以下依赖项:
<dependency>
<groupId>org.apache.kafka</groupId>
<artifactId>kafka_2.10</artifactId>
<version>0.9.0.0</version>
</dependency>
<dependency>
<groupId>org.apache.kafka</groupId>
<artifactId>kafka-clients</artifactId>
<version>0.11.0.0</version>
</dependency>
<dependency>
<groupId>org.apache.spark</groupId>
<artifactId>spark-core_2.11</artifactId>
<version>2.2.0</version>
<scope>provided</scope>
</dependency>
<dependency>
<groupId>org.apache.spark</groupId>
<artifactId>spark-streaming_2.10</artifactId>
<version>0.9.0-incubating</version>
<type>jar</type>
</dependency>
<dependency>
<groupId>org.apache.spark</groupId>
<artifactId>spark-streaming_2.10</artifactId>
<version>1.6.3</version>
</dependency>
<dependency>
<groupId>org.apache.spark</groupId>
<artifactId>spark-streaming-kafka_2.10</artifactId>
<version>1.6.3</version>
<type>jar</type>
</dependency>
<dependency>
<groupId>log4j</groupId>
<artifactId>log4j</artifactId>
<version>1.2.17</version>
</dependency>
<dependency>
<groupId>org.anarres.lzo</groupId>
<artifactId>lzo-core</artifactId>
<version>1.0.5</version>
<type>jar</type>
</dependency>
<dependency>
<groupId>com.fasterxml.jackson.core</groupId>
<artifactId>jackson-databind</artifactId>
<version>2.8.2</version>
</dependency>
<dependency>
<groupId>com.fasterxml.jackson.module</groupId>
<artifactId>jackson-module-scala_2.10</artifactId>
<version>2.8.2</version>
</dependency>
<dependency>
<groupId>com.msiops.footing</groupId>
<artifactId>footing-tuple</artifactId>
<version>0.2</version>
</dependency>
我是否缺少某些依赖项或代码中的问题?为什么这段代码不会收到任何消息?
【问题讨论】:
-
您是否能够使用基于控制台的消费者来消费消息?如果不是,那么生产者可能有问题。另外,请检查您的端口号是否正确。我认为 POM 中不应该有任何问题,如果有问题,它应该不允许您构建/编译项目。
-
@NileshPhrate- 是的,我可以使用 Kafka 的控制台消费者来消费消息,所以我们可以说这个问题与 kafka 或 zookeeper 无关,即我用于控制台方法的相同 ip 和端口。
标签: java apache-spark apache-kafka spark-streaming