【发布时间】:2018-01-18 05:42:45
【问题描述】:
我有以下案例类:
trait Event
object Event {
case class ProducerStreamActivated[T <: KafkaMessage](kafkaTopic: String, stream: SourceQueueWithComplete[T]) extends Event
}
trait KafkaMessage
object KafkaMessage {
case class DefaultMessage(message: String, timestamp: DateTime) extends KafkaMessage {
def this() = this("DEFAULT-EMPTY-MESSAGE", DateTime.now(DateTimeZone.UTC))
}
case class DefaultMessageBundle(messages: Seq[DefaultMessage], timeStamp: DateTime) extends KafkaMessage {
def this() = this(Seq.empty, DateTime.now(DateTimeZone.UTC))
}
}
在我的一个 Actor 中,我有以下方法来识别实际类型:
class KafkaPublisher[T <: KafkaMessage: TypeTag] extends Actor {
def paramInfo[T](x: T)(implicit tag: TypeTag[T]): Unit = {
val targs = typeOf[T] match { case TypeRef(_, _, args) => args }
println(s"type of $x has type arguments $targs")
}
implicit val system = context.system
val log = Logging(system, this.getClass.getName)
override final def receive = {
case ProducerStreamActivated(_, stream) =>
paramInfo(stream)
log.info(s"Activated stream for Kafka Producer with ActorName >> ${self.path.name} << ActorPath >> ${self.path} <<")
context.become(active(stream))
case other =>
log.warning("KafkaPublisher got some unknown message while producing: " + other)
}
def active(stream: SourceQueueWithComplete[KafkaMessage]): Receive = {
case msg: T =>
stream.offer(msg)
case other =>
log.warning("KafkaPublisher got the unknown message while producing: " + other)
}
}
object KafkaPublisher {
def props[T <: KafkaMessage: TypeTag] =
Props(new KafkaPublisher[T])
}
我在父 Actor 中创建了一个 ProducerStreamActivated(...) 的实例,如下所示:
val stream = producerStream[DefaultMessage](producerProperties)
def producerStream[T: Converter](producerProperties: Map[String, String]): SourceQueueWithComplete[T] = {
if (Try(producerProperties("isEnabled").toBoolean).getOrElse(false)) {
log.info(s"Kafka is enabled for topic ${producerProperties("publish-topic")}")
val streamFlow = flowToKafka[T](producerProperties)
val streamSink = sink(producerProperties)
source[T].via(streamFlow).to(streamSink).run()
} else {
// We just Log to the Console and by pass all Kafka communication
log.info(s"Kafka is disabled for topic ${producerProperties("publish-topic")}")
source[T].via(flowToLog[T](log)).to(Sink.ignore).run()
}
}
当我现在在我的子角色中打印 SourceQueueWithComplete[T] 流中包含的类型时,我会看到包含的基类 KafkaMessage,而不是预期的 DefaultMessage。有什么办法可以缓解这种情况吗?
【问题讨论】:
-
只是一个疯狂的猜测......但不是因为 JVM(类型擦除)而丢失信息的问题,而是因为您在 Kafka 上发布时对其进行了序列化,然后尝试匹配在将消息解码回原始类型之前打开消息?
-
我还没有将它发布到 Kafka,正如您在我的 producerStream 方法中看到的那样, producerProperties("isEnabled").toBoolean 设置为 false。因此它被禁用,我有一个写入日志文件的流程!可以在 flowToLog 函数中看到,然后应该写入日志文件!
-
试试 ClassTag 而不是 TypeTag?
-
ClassTag 也不起作用!我最初有 ClassTag,然后切换到 TypeTag,因为它具有比 ClassTag 更先进的功能
标签: scala apache-kafka akka akka-stream