【问题标题】:How do I implement timeseries rollups in Kafka?如何在 Kafka 中实现时间序列汇总?
【发布时间】:2019-04-15 11:53:46
【问题描述】:

我想使用 Kafka 在公司内分配高频金融市场价格。来自不同提供商的数据以每秒 2000-3000 个数字的速度传入。消费者对最新的点感兴趣,因为那是最近的价格,然而,他们往往也对获取价格的历史感兴趣。

现在,像美元/欧元汇率 (EURUSD) 这样的高流动性系列可能会导致每秒最多 100 条消息。当消费者想要历史数据时,他们想要一个抽样系列,而不是整个消息日志,因为那将是巨大的。例如,他们可能只想要过去每 5 分钟的价格历史记录,比如 10 天,即在过去的 8600 万次报价(10 天 * 24 小时 * 3600秒 * 100/秒 = 日志中的 8640 万条消息)。

每 30000 个解析整个 10 天的日志,这肯定是一项非常昂贵的操作。显然,我可以有一个消费者这样做,然后每 5 分钟重新发布到另一个主题,但是我现在有两个不同的主题用于同一个股票代码 (EURUSD),它再次引入了一种“批量与实时”架构。而且,我不想这么快用完空间。每秒存储 100 个滴答声太多了。同时,我还希望在不运行两个主题的情况下获得最新价格。

如何解决?理想情况下,我希望随时发布实时价格,但也希望在返回日志时,每 5 分钟左右才获得一次历史消息。这是可行的/可行的,没有昂贵的扫描? Kafka 是否可以推出未存储在日志中的消息(即丢失不是什么大问题的消息),但每 5 分钟存储一次?这将如何实现?

【问题讨论】:

    标签: apache-kafka


    【解决方案1】:

    您可以使用offsetsForTime 获取所需分区的偏移量图并从那里进行搜索。据我所知,这可以通过引入基于时间的索引来实现(参见https://cwiki.apache.org/confluence/display/KAFKA/KIP-33+-+Add+a+time+based+log+index#KIP-33-Addatimebasedlogindex-Enforcetimebasedlogretention)——所以我认为它在可能的范围内是有效的。

    你不能告诉 Kafka 根据时间戳有选择地存储。如果您的主题仅包含这些选定的消息,您应该复制到一个新主题

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 2018-05-15
      • 2016-04-25
      • 2016-09-12
      • 1970-01-01
      • 2015-06-03
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多