【发布时间】:2019-01-09 05:21:48
【问题描述】:
我有一个简单的案例类
case class KafkaContainer(key: String, payload: AnyRef)
然后我想通过生产者将这个发送到 kafka 主题我这样做
val byteArrayStream = new ByteArrayOutputStream()
val output = AvroOutputStream.binary[KafkaContainer](byteArrayStream)
output.write(msg)
output.close()
val bytes = byteArrayStream.toByteArray
producer.send(new ProducerRecord("my_topic", msg.key, bytes))
而且效果很好
然后我尝试消费这个
Consumer.committableSource(consumerSettings, Subscriptions.topics("my_topic"))
.map { msg =>
val in: ByteArrayInputStream = new ByteArrayInputStream(msg.record.value())
val input: AvroBinaryInputStream[KafkaContainer] = AvroInputStream.binary[KafkaContainer](in)
val result: Option[KafkaContainer] = input.iterator.toSeq.headOption
input.close()
...
}.runWith(Sink.ignore)
这适用于有效载荷中的任何类。
但是!如果是任何参考。消费者代码失败
错误:(38, 96) 找不到证据参数的隐含值 类型 com.sksamuel.avro4s.FromRecord[test.messages.KafkaContainer] val 输入:AvroBinaryInputStream[KafkaContainer] = AvroInputStream.binaryKafkaContainer
错误:(38, 96) 不够 二进制方法的参数:(隐式证据$21: com.sksamuel.avro4s.SchemaFor[test.messages.KafkaContainer], 隐含证据$22: com.sksamuel.avro4s.FromRecord[test.messages.KafkaContainer]) com.sksamuel.avro4s.AvroBinaryInputStream[test.messages.KafkaContainer]。 未指定值参数证据 $22。 val 输入:AvroBinaryInputStream[KafkaContainer] = AvroInputStream.binaryKafkaContainer
如果我声明隐含
implicit val schemaFor: SchemaFor[KafkaContainer] = SchemaFor[KafkaContainer]
implicit val fromRecord: FromRecord[KafkaContainer] = FromRecord[KafkaContainer]
编译失败
错误:(58, 71) 找不到类型的惰性隐式值 com.sksamuel.avro4s.FromValue[对象] 隐式验证 fromRecord: FromRecord[KafkaContainer] = FromRecord[KafkaContainer]
错误:(58, 71) 没有足够的参数 方法lazyConverter:(隐式fromValue: shapeless.Lazy[com.sksamuel.avro4s.FromValue[Object]])shapeless.Lazy[com.sksamuel.avro4s.FromValue[Object]]。 未指定值参数 fromValue。 隐式验证 fromRecord: FromRecord[KafkaContainer] = FromRecord[KafkaContainer]
如果添加每个隐含的编译器是必需的
lazy implicit val fromValue: FromValue[Object] = FromValue[Object]
implicit val fromRecordObject: FromRecord[Object] = FromRecord[Object]
implicit val schemaFor: SchemaFor[KafkaContainer] = SchemaFor[KafkaContainer]
implicit val fromRecord: FromRecord[KafkaContainer] = FromRecord[KafkaContainer]
编译失败并出现错误
Error:(58, 69) 宏扩展时出现异常: java.lang.IllegalArgumentException:要求失败:需要一个案例 类但 Object 不在 scala.Predef$.require(Predef.scala:277) 在 com.sksamuel.avro4s.FromRecord$.applyImpl(FromRecord.scala:283) 隐式验证 fromRecordObject: FromRecord[Object] = FromRecord[Object]
但如果我将 AnyRef 替换为某个类 - 不需要隐式,一切都会再次正常运行
【问题讨论】:
-
编译器不知道如何编组 AnyRef,它实际上可以是任何东西。
标签: scala apache-kafka avro4s akka-kafka