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

ZIO技术问询:如何在ZIO-Http中用Schema替代Case Class处理JSON?

技术求助:基于Schema映射JSON并转换为Avro写入Kafka Topic

我需要实现从HTTP请求中获取JSON体,将其转换为Avro格式后写入Kafka Topic。目前已经通过Case Class完成了基础代码,但希望改用Schema来映射处理JSON,尝试使用Zio-Json实现时遇到了问题,特此寻求帮助。

现有代码

import zhttp.http._
import zio._
import zhttp.http.{Http, Method, Request, Response, Status}
import zhttp.service.Server
import zio.json._
import zio.kafka._
import zio.kafka.serde.Serde
import zio.schema._


case class Experiments(experimentId: String,
                       variantId: String,
                       accountId: String,
                       deviceId: String,
                       date: Int)

//case class RootInterface (events: Seq[Experiments])


object Experiments {
  implicit val encoder: JsonEncoder[Experiments] = DeriveJsonEncoder.gen[Experiments]
  implicit val decoder: JsonDecoder[Experiments] = DeriveJsonDecoder.gen[Experiments]
  implicit val codec: JsonCodec[Experiments] = DeriveJsonCodec.gen[Experiments]
  implicit val schema: Schema[Experiments] = DeriveSchema.gen

}

object HttpService {
  def apply(): Http[ExpEnvironment, Throwable, Request, Response] =
    Http.collectZIO[Request] {

      case req@(Method.POST -> !! / "zioCollector") =>
        val c = req.body.asString.map(_.fromJson[Experiments])
        for {
          u <- req.body.asString.map(_.fromJson[Experiments])
          r <- u match {
            case Left(e) =>
              ZIO.debug(s"Failed to parse the input: $e").as(
                Response.text(e).setStatus(Status.BadRequest)
              )
            case Right(u) =>
              println(s"$u +       =====")
              ExpEnvironment.register(u)
                .map(id => Response.text(id))
          }
        }
        yield r
    }
}

//  val experimentsSerde: Serde[Any, Experiments] = Serde.string.inmapM { string =>
//    //desericalization
//    ZIO.fromEither(string.fromJson[Experiments].left.map(errorMessage => new RuntimeException(errorMessage)))
//  } { theMatch =>
//    ZIO.effect(theMatch.toJson)

//  }

object ZioCollectorMain extends ZIOAppDefault {
  def run: ZIO[Environment with ZIOAppArgs with Scope, Any, Any] = {
    Server.start(
      port = 9001,
      http = HttpService()).provide(ZLayerExp.layer)
  }
}

示例请求JSON

{
"experimentId": "abc",
"variantId": "123",
"accountId": "123",
"deviceId": "123",
"date": 1664544365
}

内容的提问来源于stack exchange,提问作者Mohammed Mukhtar Ali

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.15 22:20:33