【问题标题】:How to pass records from Kafka to method?如何将记录从 Kafka 传递到方法?
【发布时间】: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


【解决方案1】:

您可以在 kafka 流应用程序的接收器部分使用 ForeachWriter[T] 来处理查询的每一行,如下所示:

   datasetOfString.write.foreach(new ForeachWriter[String] {

     def open(partitionId: Long, version: Long): Boolean = {
       // open connection
     }

     def process(record: String) = {
       // write string to connection
     }

     def close(errorOrNull: Throwable): Unit = {
       // close the connection
     }
   })

【讨论】:

  • 你是我的明星!
  • 对未来的读者说一句话:public abstract void process(T value) Called to process the data in the executor side. This method will be called only when open returns true.
【解决方案2】:

linesDataset&lt;String&gt;,其中 Kafka 中的值作为行。

如何将这些传递给处理它们的方法。

根据您的具体需求,您当然可以使用 foreach 运算符或使用任何其他可用于批处理数据集的运算符或函数。

您可以使用 withColumn(...)selectmap 运算符。

换句话说,将 Spark Structured Streaming 视为带有流数据集的 Spark SQL。

【讨论】:

  • 感谢 Jack,这对我是个很有见地的新手,能像大海一样找到文档。
猜你喜欢
  • 1970-01-01
  • 2020-10-23
  • 2016-04-23
  • 1970-01-01
  • 1970-01-01
  • 2021-11-27
  • 1970-01-01
  • 2016-05-13
  • 1970-01-01
相关资源
最近更新 更多