Akka Stream持续消费WebSocket并推送到Kafka的实现相关疑问
问题解答
1. 是否可以通过GraphDSL实现WebSocket的持续消费?如果可以麻烦给出示例。
可以实现,以下是符合Akka规范的示例,同时修正了你原有JSON字符串拆分的不合理实现,改用Circe做正规JSON解析:
import akka.NotUsed import akka.actor.ActorSystem import akka.http.scaladsl.Http import akka.http.scaladsl.model.ws.{Message, TextMessage, WebSocketRequest} import akka.stream.{ClosedShape, Materializer} import akka.stream.scaladsl.{Flow, GraphDSL, RunnableGraph, Sink, Source} import akka.kafka.scaladsl.Producer import akka.kafka.ProducerSettings import org.apache.kafka.common.serialization.StringSerializer import io.circe._ import io.circe.generic.semiauto._ import io.circe.parser._ // 定义单币种价格样例类 case class CoinPrice(coin: String, price: String) object CoinPrice { implicit val encoder: Encoder[CoinPrice] = deriveEncoder[CoinPrice] } implicit val system: ActorSystem = ActorSystem("CoinPriceStream") implicit val mat: Materializer = Materializer(system) implicit val ec = system.dispatcher // Kafka生产者配置 val kafkaProducerSettings = ProducerSettings(system, new StringSerializer, new StringSerializer) .withBootstrapServers("localhost:9092") // WebSocket消息转换Flow:提取文本、解析JSON、拆分为单币种消息 val messageProcessFlow: Flow[Message, String, NotUsed] = Flow[Message] .collect { case TextMessage.Strict(text) => text } // 仅处理严格文本消息,流式消息可按需扩展 .map(parse(_).flatMap(_.as[Map[String, String]])) // 解析原始JSON为KV结构 .collect { case Right(priceMap) => priceMap } // 忽略解析失败的消息,可按需加死信队列 .mapConcat(priceMap => priceMap.map { case (coin, price) => CoinPrice(coin, price) }) .map(coinPrice => coinPrice.encoder.noSpaces) // 序列化为单币种JSON字符串 // 用GraphDSL构建完整流 val streamGraph = RunnableGraph.fromGraph(GraphDSL.create() { implicit builder => import GraphDSL.Implicits._ // WebSocket流作为输入源 val webSocketSource = Http().webSocketClientFlow(WebSocketRequest("wss://ws.coincap.io/prices?assets=ALL")) .map(_._1) // 仅取接收的消息,忽略发送通道 // Kafka Sink val kafkaSink = Producer.plainSink(kafkaProducerSettings.withTopic("coin-price-topic")) // 流拓扑 webSocketSource ~> messageProcessFlow ~> kafkaSink ClosedShape }) // 启动流 streamGraph.run()
如果不需要主动给WebSocket服务端发消息,用webSocketClientFlow的写法比singleWebSocketRequest更简洁。
2. 使用GraphDSL消费WebSocket是否是合理方案?
- 如果你的流逻辑只有单一链路(接收消息→转换→发Kafka),完全不需要用GraphDSL,直接用链式写法更简洁易读,GraphDSL是为多输入输出、分流/广播/合并等复杂拓扑设计的,简单场景用属于过度设计。
- 如果你后续需要加分支逻辑(比如一部分消息发Kafka、一部分消息落本地日志、一部分消息做实时告警),用GraphDSL是合理的。
3. 我当前在推送到Kafka前对JSON消息做拆分处理的方式是否合适,还是直接推送原消息能获得更低延迟?
首先你的现有拆分方式完全不合适:
用字符串替换、分割的方式处理JSON容错性极差,只要JSON格式有微小变化(比如字段值里出现逗号、冒号,或者空格格式变化)就会生成错误数据,必须用正规JSON解析库处理。
关于拆分时机的延迟对比:
- 上游拆分:仅做一次解析和拆分,下游所有消费者都可以直接消费结构化的单币种数据,避免重复计算,总链路资源消耗更低,单条消息的延迟增加可以忽略不计(微秒级)。
- 直接推原消息:仅适合你有多个下游需要完整原始数据的场景,否则下游每个消费者都要重复解析大JSON,总延迟更高,资源浪费更严重。
推荐在推送Kafka前完成拆分,优先保证下游消费的易用性和总资源效率。
4. 我计划后续使用Apache Storm消费Kafka中的消息,该方案是否合理,还是应该继续使用Akka生态组件实现?
两个方案都可行,选择标准看你的现有技术栈:
- 如果你的团队已经有成熟的Storm运维、开发经验,现有业务也大量跑在Storm上,直接用Storm完全合理。
- 如果是新做的项目,没有历史技术栈包袱,更推荐用Akka生态组件(Alpakka Kafka、Akka Projections、Akka Streams)实现,整个从数据采集到消费的链路都用同一套技术栈,API统一,不用额外学习Storm的开发模型,运维成本也更低。
内容的提问来源于stack exchange,提问作者Ares
相关产品推荐
相关产品推荐

