【问题标题】:Streaming Data Processing and nano second time resolution流数据处理和纳秒时间分辨率
【发布时间】:2019-06-21 12:01:19
【问题描述】:

我刚刚开始研究实时流数据处理框架的主题,我有一个问题我至今找不到任何结论性的答案:

常见的嫌疑人(Apache 的 Spark、Kafka、Storm、Flink 等)是否支持以 纳秒(甚至皮秒)的事件时间分辨率处理数据?

大多数人和文档都在谈论毫秒或微秒分辨率,但如果可能有更多分辨率或问题,我无法找到明确的答案。我推断唯一具有该功能的框架是 influxData 的 Kapacitor 框架,因为他们的 TSDB influxDB 似乎以纳秒分辨率存储时间戳。

这里的任何人都可以就此提供一些见解,甚至提供一些知情的事实吗?提供此功能的替代解决方案/框架?

任何事情都将不胜感激!

感谢和问候,

西蒙


我的问题的背景:我正在一个环境中工作,该环境具有大量用于数据存储和处理的专有实现,并且目前正在考虑一些组织/优化。我们正在使用许多不同的诊断/测量系统以不同的采样率进行等离子体物理实验,现在高达“每秒超过千兆样本”。我们系统中的一个常见事实/假设是每个样本确实有一个以纳秒分辨率记录的事件时间。当尝试使用已建立的流(或批处理)处理框架时,我们必须保持这个时间戳分辨率。或者更进一步,因为我们最近在某些系统上突破了 1 Gsps 阈值。因此我的问题。

【问题讨论】:

    标签: apache-kafka apache-flink apache-storm apache-kafka-streams kapacitor


    【解决方案1】:

    如果不清楚,您应该注意事件时间和处理时间之间的差异:

    事件时间 - 事件源的生成时间

    处理时间 - 处理引擎内的事件执行时间

    src:Flink docs

    AFAIK Storm 不支持事件时间,Spark 的支持有限。这让 Kafka Streams 和 Flink 有待考虑。

    Flink 使用 long 类型作为时间戳。在docs 中提到,该值是自 1970-01-01T00:00:00Z 以来的毫秒数,但是 AFAIK,当您使用事件时间特征时,唯一衡量进度的是事件时间戳。因此,如果您可以将您的价值观纳入长期范围,那么它应该是可行的。

    编辑:

    一般来说,水印(基于时间戳)用于测量窗口、触发器等中事件时间的进度。所以,如果你使用:

    • AssignerWithPeriodicWatermarks 然后在处理时域的配置(自动水印间隔)中定义的间隔中发出新的水印 - 即使使用事件时间特征也是如此。有关详细信息,请参见例如org.apache.flink.streaming.runtime.operators.TimestampsAndPeriodicWatermarksOperator#open() 方法,其中注册了处理时间中的计时器。因此,如果自动水印设置为 500 毫秒,那么每 500 毫秒的处理时间(取自 System.currentTimeMillis())就会发出一个新的水印,但水印的时间戳是基于事件的时间戳。

    • AssignerWithPunctuatedWatermarks 那么最好的描述可以在org.apache.flink.streaming.api.datastream.DataStream#assignTimestampsAndWatermarks(org.apache.flink.streaming.api.functions.AssignerWithPunctuatedWatermarks<T>)的文档中找到:

    为数据流中的元素分配时间戳,并根据元素本身创建水印以指示事件时间进度。

    此方法完全基于流元素创建水印。 对于通过AssignerWithPunctuatedWatermarks#extractTimestamp(Object, long) 处理的每个元素,调用AssignerWithPunctuatedWatermarks#checkAndGetNextWatermark(Object, long) 方法,如果返回的水印,则会发出新的水印值是非负的并且大于前一个水印。

    当数据流中嵌入了水印元素,或者某些元素带有可以用来确定当前事件时间水印的标记时,这种方法很有用。 此操作使程序员可以完全控制水印的生成。用户应该注意,过于激进的水印生成(即每秒生成数百个水印)可能会降低性能。

    要了解水印的工作原理,强烈推荐阅读:Tyler Akidau on Streaming 102

    【讨论】:

    • 感谢您了解 Spark 和 Storm。关于 Flink:我是否理解您所说的正确,我可以使用 long 值,以纳秒为单位,并确保每个使用数据的人都知道它不是以毫秒为单位的? - 我想知道框架是否可能在内部某个地方使用毫秒时间戳的假设,正如他们特别提到的那样。你对此有什么确定的吗?
    • 认为更高分辨率的事件时间时间戳可能在 Flink 中工作似乎是合理的。这正在用户邮件列表中进行讨论——请关注此线程:apache-flink-user-mailing-list-archive.2336050.n4.nabble.com/…。
    【解决方案2】:

    虽然 Kafka Streams 使用毫秒级分辨率,但运行时实际上是不可知的。到头来就是长篇大论。

    话虽如此,“问题”是时间窗口的定义。如果您指定 1 分钟的时间窗口,但您的时间戳分辨率小于毫秒,则您的窗口将小于 1 分钟。作为一种解决方法,您可以将窗口放大,例如,1000 分钟或 1,000,000 分钟以实现微/纳秒分辨率。

    另一个“问题”是,代理只了解毫秒级的分辨率,而保留时间是基于此的。因此,您需要将保留时间设置得更高,以“欺骗”代理并避免它过早删除数据。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2019-11-25
      • 1970-01-01
      • 1970-01-01
      • 2023-04-07
      • 2013-11-27
      • 1970-01-01
      • 1970-01-01
      • 2015-01-28
      相关资源
      最近更新 更多