【发布时间】:2016-10-13 18:13:21
【问题描述】:
我将 Play Framework (Scala) 用于微服务,并使用 Kafka 作为事件总线。我有一个事件消费者,它映射到一个事件类,如下所示:
case class MovieEvent[T] (
mediaId: String,
config: T
)
object MovieEvent {
implicit def movieEventFormat[T: Format]: Format[MovieEvent[T]] =
((__ \ "mediaId").format[String] ~
(__ \ "config").format[T]
)(MovieEvent.apply _, unlift(MovieEvent.unapply))
}
object MovieProvider extends SerializableEnumeration {
implicit val providerReads: Reads[MovieProvider.Value] = SerializableEnumeration.jsonReader(MovieProvider)
implicit val providerWrites: Writes[MovieProvider.Value] = SerializableEnumeration.jsonWrites
val Dreamworks, Disney, Paramount = Value
}
消费者看起来像:
class MovieEventConsumer @Inject()(movieService: MovieService
) extends ConsumerRecordProcessor with LazyLogging {
override def process(record: IncomingRecord): Unit = {
val movieEventJson = Json.parse(record.valueString).validate[MovieEvent[DreamworksConfiguration]]
movieEventJson match {
case event: JsSuccess[MovieEvent[DreamworksJobOptions]] => processMovieEvent(event.get)
case er: JsError =>
logger.error("Unrecognized MovieEvent, attempting to parse as MovieUploadEvent: " + JsError.toJson(er).toString())
try {
val data = (Json.parse(record.valueString) \ "upload").as[MovieUploadEvent]
processUploadEvent(data)
} catch {
case er: Exception => logger.error("Unrecognized kafka event", er)
}
}
}
def processMovieEvent[T](event: MovieEvent[T]): Unit = {
logger.debug(s"Received movie event: ${event}")
movieService.createMovieJob(event)
}
def processUploadEvent(event: MovieUploadEvent): Unit = {
logger.debug(s"Received upload event: ${event}")
movieService.addToCollection(event)
}
}
目前,我只能验证三种不同的 MovieEvent 配置(Dreamwork、迪士尼和派拉蒙)中的一种。我可以换掉我通过代码验证的那个,但这不是重点。但是,我想验证这三个中的任何一个,而不必增加额外的消费者。我尝试过一些不同的想法,但没有一个可以编译。我对 Play 和 Kafka 还很陌生,想知道是否有一种好方法可以做到这一点。
提前致谢!
【问题讨论】:
标签: json scala playframework playframework-2.0 apache-kafka