【问题标题】:Kafka - problems with TimestampExtractorKafka - TimestampExtractor 的问题
【发布时间】:2017-01-24 21:38:53
【问题描述】:

我用org.apache.kafka:kafka-streams:0.10.0.1

我正在尝试使用基于时间序列的流,该流似乎不会触发 KStream.Process() 来触发(“标点符号”)。 (参考here

KafkaStreams 配置中,我传入了这个参数(以及其他参数):

config.put(
  StreamsConfig.TIMESTAMP_EXTRACTOR_CLASS_CONFIG,
  EventTimeExtractor.class.getName());

这里,EventTimeExtractor 是一个自定义时间戳提取器(它实现了org.apache.kafka.streams.processor.TimestampExtractor),用于从 JSON 数据中提取时间戳信息。

当每条新记录被拉入时,我希望这会调用我的对象(派生自 TimestampExtractor)。有问题的流是 2 * 10^6 记录/分钟。我将punctuate() 设置为 60 秒,但它永远不会触发。我知道数据非常频繁地通过这个跨度,因为它会拉动旧值来迎头赶上。

事实上,它根本不会被调用。

  • 这是在 KStream 记录上设置时间戳的错误方法吗?
  • 这是声明此配置的错误方式吗?

【问题讨论】:

    标签: java apache-kafka apache-kafka-streams


    【解决方案1】:

    您的方法似乎是正确的。比较http://docs.confluent.io/3.0.1/streams/developer-guide.html#optional-configuration-parameters中的"Timestamp Extractor (timestamp.extractor):"段落

    不确定,为什么不使用您的自定义时间戳提取器。看看org.apache.kafka.streams.processor.internals.StreamTask。在构造函数中应该有类似

    TimestampExtractor timestampExtractor1 = (TimestampExtractor)config.getConfiguredInstance("timestamp.extractor", TimestampExtractor.class);
    

    检查您的自定义提取器是否在那里被拾取...

    【讨论】:

    • 我在.extract() 函数中添加了一些日志,但它永远不会被击中。类的构造函数是。
    • 您能分享您的数据和/或代码吗? matthias@confluent.io
    【解决方案2】:

    2017 年 11 月更新: Kafka 1.0 中的 Kafka Streams 现在支持 punctuate() 的流时间和处理时间(挂钟时间)行为。因此,您可以选择您喜欢的任何行为。

    您的设置对我来说似乎是正确的。

    您需要注意的事项:从 Kafka 0.10.0 开始,punctuate() 方法在 stream-time 上运行(默认情况下,即基于默认时间戳提取器,stream-time将意味着事件时间)。而stream-time只有在有新数据记录进来时才会提前,而stream-time提前多少是由这些新记录的相关时间戳决定的。

    例如:

    • 假设您已将 punctuate() 设置为每 1 分钟调用一次 = 60 * 1000(注意:流时间的 1 分钟)。现在,如果碰巧在接下来的 5 分钟内没有收到任何数据,则根本不会调用 punctuate() —— 即使您可能期望它会被调用 5 次。为什么?同样,因为punctuate() 依赖于流时间,流时间仅根据新接收的数据记录提前。

    这可能会导致您看到的行为吗?

    展望未来:Kafka 项目中已经在讨论如何使punctuate() 更加灵活,例如触发它不仅基于stream-time(默认为event-time),还基于processing-time

    【讨论】:

    • 您的应用程序正在读取的记录是否包含嵌入的时间戳?例如,您使用的是 0.9 版的 Kafka 集群还是旧版本?
    • 所有服务器组件都来自 confluent repo。在过去几周内安装。记录确实包含时间戳字段(记录是 JSON 编码的)。我不希望 kafka 选择这个领域,因此使用“提取器”(听起来很不祥)
    • 我认为所说的 kafka 集群出了点问题。 5 个服务器和 1 个 zk 都在 AWS 上运行。我试图从 3 个线程推送数据,但几乎没有达到 1M/分钟。添加第二个进程(单独的机器)降低总吞吐量。网络接口均
    • 注意一些遥远的观察者:确保将文件描述符计数设置为实用的值。在 centos/aws 上,默认值似乎是 1024,即使是小型设置也不够(增加到 65535)。这是我的集群的一个大问题。
    • 但是你的问题依然存在,即使增加nofile?
    【解决方案3】:

    我认为这是经纪人级别的另一个问题。我使用具有更多 CPU 和 RAM 的实例重建了集群。现在我得到了我预期的结果。

    注意远方的观察者:如果您的 KStream 应用程序行为异常,请查看您的代理并确保它们没有卡在 GC 中并且有足够的“空间”用于文件句柄、RAM 等。

    See also

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2018-07-03
      • 2017-10-15
      • 1970-01-01
      • 2020-01-25
      • 1970-01-01
      • 1970-01-01
      • 2019-07-06
      • 2017-07-22
      相关资源
      最近更新 更多