You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在Scala的Kafka Producer中读取外部API数据并发送至Kafka Consumer

解决Kafka Producer发送JSON数据到主题并消费的问题

首先,你的代码里有几个关键问题需要修正,同时我会给你两种清晰的方案来处理JSON数据的发送和消费,适合新手快速上手:


问题分析

当前你的Producer代码存在两个核心问题:

  1. KafkaProducer泛型与序列化器不匹配:你配置的是StringSerializer,但声明的Producer泛型是[Nothing, (String,JsValue)],这会导致序列化失败——因为StringSerializer只能处理String类型的Key/Value。
  2. 发送的数据类型错误:content.as[JsObject].fields(1)返回的是(String, JsValue)元组,StringSerializer无法处理这种复杂类型的数据。

方案一:使用String序列化器(最简单,适合快速验证)

这种方案不需要自定义序列化器,直接将JSON对象转为字符串发送,消费时再解析回JSON,非常适合新手入门。

修改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.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)

  // 将JSON对象转为字符串,适配StringSerializer
  val jsonString = Json.stringify(content)

  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")

  // 泛型改为[String, String],与配置的序列化器匹配
  val producer = new KafkaProducer[String, String](props)
  val topic = "quick-start"

  try {
    // 用MSFT作为消息key,也可以传入null;发送JSON字符串
    val record = new ProducerRecord(topic, "MSFT", jsonString)
    producer.send(record).get() // 调用get()等待发送完成,方便调试,生产环境可异步发送
    println("数据发送成功")
  } catch {
    case e: Exception => e.printStackTrace()
  } finally {
    producer.close()
  }
}

修改Consumer.scala(解析JSON并准备存库)

import java.util.{Collections, Properties}
import org.apache.kafka.clients.consumer.KafkaConsumer
import scala.collection.JavaConverters._
import play.api.libs.json.Json

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[String, String](props)
  val topics = List("quick-start")

  try {
    consumer.subscribe(topics.asJava)
    while (true) {
      val records = consumer.poll(10)
      for (record <- records.asScala) {
        println(s"Topic: ${record.topic()}, Key: ${record.key()}, Offset: ${record.offset()}")
        // 将字符串解析回JSON对象
        val jsonData = Json.parse(record.value())
        // 提取时间序列数据,这里可以替换为数据库存储逻辑
        val timeSeries = (jsonData \ "Time Series (5min)").as[Map[String, Map[String, String]]]
        println(s"获取到${timeSeries.size}条时间序列数据")
        // TODO: 编写数据库插入代码,比如用JDBC、Slick等框架
      }
    }
  } catch {
    case e: Exception => e.printStackTrace()
  } finally {
    consumer.close()
  }
}

方案二:自定义JSON序列化器(生产环境推荐)

如果不想每次手动转字符串,可以自定义基于Play JSON的序列化器,让Kafka直接处理JsValue类型的数据,代码更简洁规范。

1. 创建自定义序列化器和反序列化器

import org.apache.kafka.common.serialization.{Deserializer, Serializer}
import play.api.libs.json.{JsValue, Json}
import java.nio.charset.StandardCharsets

class PlayJsonSerializer extends Serializer[JsValue] {
  override def configure(configs: java.util.Map[String, _], isKey: Boolean): Unit = {}

  override def serialize(topic: String, data: JsValue): Array[Byte] = {
    Json.stringify(data).getBytes(StandardCharsets.UTF_8)
  }

  override def close(): Unit = {}
}

class PlayJsonDeserializer extends Deserializer[JsValue] {
  override def configure(configs: java.util.Map[String, _], isKey: Boolean): Unit = {}

  override def deserialize(topic: String, data: Array[Byte]): JsValue = {
    Json.parse(new String(data, StandardCharsets.UTF_8))
  }

  override def close(): Unit = {}
}

2. 修改Producer配置和代码

// ... 其他代码不变
val props: Properties = new Properties()
props.put("bootstrap.servers","localhost:9092")
props.put("key.serializer", "org.apache.kafka.common.serialization.StringSerializer")
// 替换为自定义的JSON序列化器,注意填写正确的包路径
props.put("value.serializer", "your.package.name.PlayJsonSerializer")
props.put("acks","all")

// 泛型改为[String, JsValue]
val producer = new KafkaProducer[String, JsValue](props)
val topic = "quick-start"

try {
  val record = new ProducerRecord(topic, "MSFT", content)
  producer.send(record).get()
  println("数据发送成功")
} catch {
  case e: Exception => e.printStackTrace()
} finally {
  producer.close()
}

3. 修改Consumer配置和代码

// ... 其他代码不变
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")
// 替换为自定义的JSON反序列化器,注意填写正确的包路径
props.put("value.deserializer", "your.package.name.PlayJsonDeserializer")
props.put("enable.auto.commit", "true")
props.put("auto.commit.interval.ms", "1000")

// 泛型改为[String, JsValue]
val consumer = new KafkaConsumer[String, JsValue](props)
val topics = List("quick-start")

try {
  consumer.subscribe(topics.asJava)
  while (true) {
    val records = consumer.poll(10)
    for (record <- records.asScala) {
      println(s"Topic: ${record.topic()}, Key: ${record.key()}, Offset: ${record.offset()}")
      val jsonData = record.value()
      // 直接使用JsValue处理数据,无需手动解析字符串
      val timeSeries = (jsonData \ "Time Series (5min)").as[Map[String, Map[String, String]]]
      println(s"获取到${timeSeries.size}条时间序列数据")
    }
  }
} catch {
  case e: Exception => e.printStackTrace()
} finally {
  consumer.close()
}

额外注意事项

  • 数据库存储:你可以使用JDBC、Slick或其他ORM框架实现数据存储,比如在Consumer中解析出每条时间序列数据后,插入到对应的数据库表中。
  • 错误处理:生产环境中建议添加更完善的错误处理,比如重试机制、死信队列等,避免数据丢失。
  • 依赖兼容性:你的build.sbt中scala 2.12.2搭配kafka 2.1.0和play-json 2.8.0版本是兼容的,无需调整。

内容的提问来源于stack exchange,提问作者Ruchir Dixit

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.04.29 05:02:28