【发布时间】:2018-02-12 20:15:57
【问题描述】:
我有一个 Kafka 队列,从中读取数据如下:
private static void startKafkaConsumerStream() {
try {
System.out.println("Print method: startKafkaConsumerStream");
Dataset<String> lines = (Dataset<String>) _spark
.readStream()
.format("kafka")
.option("kafka.bootstrap.servers", getProperty("kafka.bootstrap.servers"))
.option("subscribe", HTTP_FED_VO_TOPIC)
.option("startingOffsets", "latest")
.load()
.selectExpr("CAST(value AS STRING)")
.as(Encoders.STRING());
StreamingQuery query = lines.writeStream()
.outputMode("append")
.format("console")
.start();
query.awaitTermination();
} catch (Exception e) {
e.printStackTrace();
}
}
要求:使用上面的代码,我可以将记录打印到控制台,但是,我受到威胁,因为我如何将这些传递给将处理它们的方法。
为此,我尝试查看文档,但找不到任何相关内容。由于我是新手,这听起来可能有点傻。但是我被卡住了,非常感谢任何提示。
应用程序的目标 应用程序的目标是接受请求并将其发送到 Kafka,然后在一个单独的线程中实现一个 Kafka 读取器,该读取器负责读取和处理请求并生成输出到另一个 Kafka 队列。我只是在实现这个,架构不是我的主意。
【问题讨论】:
-
显式
append输出模式是否有特定原因?这是默认输出模式,因此我在问。 -
之前我提到的 tut 没有设置完成我只是对 Kafka 中可用的最新数据感兴趣(实现请求响应管道)
-
“实现了一个kafka阅读器,负责读取和处理请求”
-
这是一个机器学习 - 人工智能组件,对我来说它是一个黑盒子。所以我只是得到一个方法的句柄,我必须将 Kafka 字符串作为方法参数传递给该方法。
-
这个组件是如何工作的?这是一个功能吗?一个Java类? REST 服务?记录如何传递到应用程序?看起来
foreach可能是目前最好的选择。
标签: java apache-spark apache-kafka kafka-consumer-api spark-structured-streaming