【问题标题】:custom timestamp extractor in kafka streamskafka 流中的自定义时间戳提取器
【发布时间】:2023-03-27 11:02:01
【问题描述】:

我正在尝试使用 kafka 流来处理来自 kafka 主题的一些数据。数据来自 kafka 0.11.0 写入的 kafka 主题,该主题没有嵌入时间戳。在网上阅读了一番后,我了解到我可以通过在自定义类中扩展 TimestampExtractor 类并将其传递给 StreamsConfig 来解决这个问题。

我就是这样做的--

class MyEventTimestampExtractor extends TimestampExtractor {
  override def extract(record: ConsumerRecord[AnyRef, AnyRef], prev: Long) = {
    record.value() match {
        case w: String => 1000L
        case _ => throw new RuntimeException(s"Called for $record")
    }
  }
}

我基于this code on github

但是,当我执行sbt run 时出现此错误

[error] /home/someuser/app/blahblah/src/main/scala/main.scala:34: class MyEventTimestampExtractor needs to be abstract, since method extract in trait TimestampExtractor of type (x$1: org.apache.kafka.clients.consumer.ConsumerRecord[Object,Object], x$2: Long)Long is not defined
[error] (Note that Long does not match Long)
[error] class MyEventTimestampExtractor extends TimestampExtractor {
[error]       ^
[error] /home/someuser/app/blahblah/src/main/scala/main.scala:35: method extract overrides nothing.
[error] Note: the super classes of class MyEventTimestampExtractor contain the following, non final members named extract:
[error] def extract(x$1: org.apache.kafka.clients.consumer.ConsumerRecord[Object,Object],x$2: Long): Long
[error]   override def extract(record: ConsumerRecord[AnyRef, AnyRef], prev: Long): Long = {
[error]                ^
[error] two errors found
[error] (compile:compileIncremental) Compilation failed

我的 build.sbt 文件是这个 --

name := "kafka streams experiment"
version := "1.0"
scalaVersion := "2.12.4"

libraryDependencies ++= Seq(
  "org.apache.kafka" % "kafka-streams" % "1.0.0"
)

我真的不明白这个错误。特别是Note that Long does not match Long 周围的部分。我可能做错了什么?谢谢!

【问题讨论】:

  • 你试过没有override关键字吗?
  • 我没有提供问题中的所有上下文:(我导入了 scala long,它以 java long 的方式出现并导致编译失败:/

标签: apache-kafka apache-kafka-streams confluent-platform


【解决方案1】:

您需要 java.Long,因为此 API 是在 Java 中定义的,而您使用的是 Scala Long

【讨论】:

    【解决方案2】:

    尝试使用(观察prev函数参数的类型:

      override def extract(record: ConsumerRecord[AnyRef, AnyRef], prev: java.lang.Long) = {
    

    【讨论】:

      猜你喜欢
      • 2019-05-13
      • 1970-01-01
      • 2017-03-27
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2019-08-18
      相关资源
      最近更新 更多