【发布时间】:2020-07-29 12:58:04
【问题描述】:
假设我有两个 kafka 主题,request_topic 用于我的 Post 请求,response_topic 用于我的回复。
这是模型:
case class Request(requestId: String, body: String)
case class Response(responseId: String, body: String, requestId: String)
这是我的套接字处理程序
def socket = WebSocket.accept[String, String] { req =>
val requestId = ??? // Generate a unique requestId
val in: Sink[String, Future[Done]] = Sink.foreach[String]{ msg =>
val record = new ProducerRecord[String, Request]("request_topic", "key", Request(requestId, msg))
val producer: KafkaProducer[String, Request] = ???
Future { producer.send(record).get }
}
// Once produced, some stream processing apps will manage to process request and publish the reponse to response_topic
// The Request and Response object are linked by the requestId field.
val consumerSettings = ???
val out: Source[ConsumerRecord[String, Response], _] = Consumer
.plainSource(consumerSettings, Subscriptions.topics("response_topic"))
.filter(cr => cr.value.requestId == requestId)
.map(cr => someResponseString(cr.value))
Flow.formSinkAndSource(in, out)
}
def someResponseString(res: Response): String = ???
基本上,对于每条传入的消息,我都会向 Kafka 发布一个 Request 对象,然后该请求由一些流处理应用程序(此处未显示)进行处理,并希望将响应发布回 Kafka。
我有一些顾虑:
1 - Alpakka Kafka 连接器会为每条传入消息创建一个新的连接器实例,还是会在 Play 运行时使用相同的实例?
2 - 根据单个 requestId 过滤响应是个好主意,还是应该将整个流发送回每个客户端,让他们根据他们感兴趣的 requestId 过滤响应。
3 - 我错了吗? (我是 Websocket 的真正新手)
提前致谢。
【问题讨论】:
标签: scala websocket apache-kafka playframework alpakka