【发布时间】:2020-09-18 10:44:55
【问题描述】:
我正在为我的项目做一些使用 Spark 和 Spark 流的 POC。所以我所做的就是从 Topic 中读取文件名。从“src/main/sresource”下载文件并执行通常的“WordCount”频率应用程序。
@KafkaListener(topics = Constants.ABCWordTopic, groupId = Constants.ABC_WORD_COMSUMER_GROUP_ID)
public void processTask(@Payload String fileResourcePath) {
log.info("ABC Receiving task from WordProducer filepath {} at time {}", fileResourcePath,
LocalDateTime.now());
// Spark job
/*
* JavaRDD wordRDD =
* sparkContext.parallelize(Arrays.asList(extractFile(fileResourcePath).split(" ")));
* log.info("ABC Map Contents : {}", wordRDD.countByValue().toString());
* wordRDD.coalesce(1,
* true).saveAsTextFile("ResultSparklog_"+ System.currentTimeMillis());
*/
// Spark Streaming job
JavaPairDStream wordPairStream = streamingContext
.textFileStream(extractFile(fileResourcePath))
.flatMap(line -> Arrays.asList(SPACE.split(line)).iterator())
.mapToPair(s -> new Tuple2(s, 1)).reduceByKey((i1, i2) -> i1 + i2);
wordPairStream.foreachRDD(wordRDD -> {
// javaFunctions(wordTempRDD).writerBuilder("vocabulary", "words", mapToRow(String.class))
// .saveToCassandra();
log.info("ABC Map Contents : {}", wordRDD.keys().countByValue().toString());
wordRDD.coalesce(1, true)
.saveAsTextFile("SparkStreamResultlog_" + System.currentTimeMillis());
});
streamingContext.start();
try {
streamingContext.awaitTerminationOrTimeout(-1);
} catch (InterruptedException e) {
log.error("Terminated streaming context {}", e);
}
}
- 在上面的代码中,我正在收听 Kafka 主题(“ABCtopic”)和
处理它。 “Spark 作业”注释代码工作得非常好。
它计算单词并按预期给出结果,但是“火花
流作业”代码的行为与预期不符,它输出 null。
-
log.info("ABC Map Contents : {}", wordRDD.keys().countByValue().toString());行将“{}”作为输出。 写入文件是空的。作为 Spark 流的新手,从什么开始 鲜为人知的“Spark Streaming”是一个额外的库 持续实时处理来自任何来源的数据,例如 文件、主题等
- 上面的代码中缺少什么用于火花流输出 突出显示的日志行和输出数据文件中的“null” 正在写入磁盘,而 Spark 作业也是如此 工作非常好。
【问题讨论】:
标签: spring-boot apache-spark spark-streaming spark-streaming-kafka