【问题标题】:How to Read data from external API in Kafka producer and send it to Kafka consumer in Scala如何从 Kafka 生产者中的外部 API 读取数据并将其发送给 Scala 中的 Kafka 消费者
【发布时间】:2021-04-23 21:12:34
【问题描述】:

我是 Apache Kafka 的新手,我想从 https://www.alphavantage.co/query?function=TIME_SERIES_INTRADAY&symbol=MSFT&interval=5min&outputsize=full&apikey=demo API 读取生产者内部的数据,然后将其发送到主题并从消费者内部的主题中读取此数据以将其保存到数据库。

我无法弄清楚如何以 JSON 格式发送此数据。

我尝试了一个带有字符串值的 Kafka 消费者生产者示例:

在我的示例中,我的 Producer.scala 是:

import java.util.Properties

import org.apache.http.client.methods.HttpGet
import org.apache.http.impl.client.HttpClientBuilder
import org.apache.http.util.EntityUtils
import org.apache.kafka.clients.producer.{KafkaProducer, ProducerRecord}
import play.api.libs.json.{JsObject, JsValue, Json}
object Producer extends App {

  val url = "https://www.alphavantage.co/query?function=TIME_SERIES_INTRADAY&symbol=MSFT&interval=5min&outputsize=full&apikey=demo"
  val httpClient = HttpClientBuilder.create().build()
  val httpResponse = httpClient.execute(new HttpGet(url))
  val entity = httpResponse.getEntity
    val str = EntityUtils.toString(entity, "UTF-8")
    val content = Json.parse(str)
  val props:Properties = new Properties()
  props.put("bootstrap.servers","localhost:9092")
  props.put("key.serializer",
    "org.apache.kafka.common.serialization.StringSerializer")
  props.put("value.serializer",
    "org.apache.kafka.common.serialization.StringSerializer")
  props.put("acks","all")
  val producer = new KafkaProducer[Nothing, (String,JsValue)](props)
  val topic = "quick-start"
  try {
      val record = new ProducerRecord(topic, content.as[JsObject].fields(1))
      producer.send(record)
  }catch{
    case e:Exception => e.printStackTrace()
  }finally {
    producer.close()
  }
}

而我的 Consumer.scala 是:

import java.util.{Collections, Properties}
import java.util.regex.Pattern
import org.apache.kafka.clients.consumer.KafkaConsumer
import scala.collection.JavaConverters._
object Consumer extends App {

  val props:Properties = new Properties()
  props.put("group.id", "test")
  props.put("bootstrap.servers","localhost:9092")
  props.put("key.deserializer",
    "org.apache.kafka.common.serialization.StringDeserializer")
  props.put("value.deserializer",
    "org.apache.kafka.common.serialization.StringDeserializer")
  props.put("enable.auto.commit", "true")
  props.put("auto.commit.interval.ms", "1000")
  val consumer = new KafkaConsumer(props)
  val topics = List("quick-start")
  try {
    consumer.subscribe(topics.asJava)
    while (true) {
      val records = consumer.poll(10)
      for (record <- records.asScala) {
        println("Topic: " + record.topic() +
          ",Key: " + record.key() +
          ",Value: " + record.value() +
          ", Offset: " + record.offset() +
          ", Partition: " + record.partition())
      }
    }
  }catch{
    case e:Exception => e.printStackTrace()
  }finally {
    consumer.close()
  }
}

我的 built.sbt 是:

name := "Kafka-AkkaPractice"

version := "0.1"

scalaVersion := "2.12.2"

libraryDependencies ++= Seq(
  "org.apache.kafka" %% "kafka" % "2.1.0",
  "ch.qos.logback" % "logback-classic" % "1.1.3" % Runtime,
  "org.apache.httpcomponents" % "httpclient" % "4.5.2",
  "com.typesafe.play" %% "play-json" % "2.8.0"
)

我的理解是

props.put("key.serializer",
        "org.apache.kafka.common.serialization.StringSerializer")
props.put("value.serializer",
            "org.apache.kafka.common.serialization.StringSerializer")

适用于String类型的Key和Value。

谁能建议如何将这种类型的数据发送到我的 Kafka 主题,以便我可以从 Consumer 读取它并将其保存到数据库中?

【问题讨论】:

  • 假设 val content = Json.parse(str) 在生产者中实际上是 JSON,你需要在消费者中 解析它,那么你有什么问题呢?如果您想将数据放入数据库,这是 Kafka-Connect 的用例,而不是您自己的消费者

标签: scala apache-kafka kafka-consumer-api producer-consumer


【解决方案1】:

您无需在生产者中解析 JSON,只需验证 API 响应是否可以解析。

如果你想按原样发送数据,那么你需要val record = new ProducerRecord(topic, str)

在消费者中

 for (record <- records.asScala) {
    val content = Json.parse(record.value())

【讨论】:

  • 在消费者内部,它在val content = Json.parse(record.value()) 线上给了我错误。我使用import play.api.libs.json.Json 导入Json.
  • 我没有使用 PlayJSON 的经验,但是什么错误?
  • Error:(26, 28) ambiguous reference to overloaded definition, both method parse in object Json of type (input: Array[Byte])play.api.libs.json.JsValue and method parse in object Json of type (input: java.io.InputStream)play.api.libs.json.JsValue match argument types (Nothing) val content = Json.parse(record.value()) 这个错误
  • record.value() 应该已经是一个字符串,而不是一个字节数组,但你可以用 record.value().asInstanceOf[String] 强制它
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 2019-05-09
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2017-09-07
相关资源
最近更新 更多