如何在Scala的Kafka Producer中读取外部API数据并发送至Kafka Consumer
解决Kafka Producer发送JSON数据到主题并消费的问题
首先,你的代码里有几个关键问题需要修正,同时我会给你两种清晰的方案来处理JSON数据的发送和消费,适合新手快速上手:
问题分析
当前你的Producer代码存在两个核心问题:
- KafkaProducer泛型与序列化器不匹配:你配置的是
StringSerializer,但声明的Producer泛型是[Nothing, (String,JsValue)],这会导致序列化失败——因为StringSerializer只能处理String类型的Key/Value。 - 发送的数据类型错误:
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
相关产品推荐
相关产品推荐

