【问题标题】:Is there a Spark like Accumulator for Kafka Streams?Kafka Streams 有类似 Spark 的累加器吗?
【发布时间】:2018-09-22 10:43:03
【问题描述】:

Spark 有一个有用的 API,用于以线程安全的方式累积数据 https://spark.apache.org/docs/2.3.0/api/scala/index.html#org.apache.spark.util.AccumulatorV2,并带有一些开箱即用的有用累加器,例如多头https://spark.apache.org/docs/2.3.0/api/scala/index.html#org.apache.spark.util.LongAccumulator

我通常使用累加器将调试、分析、监视和诊断连接到 Spark 作业中。我通常在运行 Spark 作业之前启动 Future 以定期打印统计信息(例如 TPS、直方图、计数、计时等)

到目前为止,我找不到任何与 Kafka Streams 类似的东西。有什么存在吗?我想这至少对于 Kafka 应用程序的每个实例都是可能的,但是要跨多个实例进行这项工作需要创建一个中间主题。

【问题讨论】:

  • 推荐/查找工具或库的请求不在此处。
  • Kafka Streams 有聚合器和缩减器

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


【解决方案1】:

Kafka Streams 在设计上避免了并发性——如果累积的不需要容错,您可以在内存中完成并通过挂钟时间标点将其刷新。

如果需要容错,可以使用状态存储,并在标点符号中扫描整个存储以将其刷新。

这将为您提供任务级别的积累。不确定 Spark 的累加器如何详细工作,但如果它为您提供“全局”视图,我假设它需要通过网络发送数据,并且一个实例只能访问数据(或者可能是广播 - 不是当然,如何保证广播案例的一致性)。类似地,您可以将数据发送到一个主题(具有 1 个分区)以将所有数据全局收集到一个地方。

【讨论】:

  • > "类似地,您可以将数据发送到一个主题(有 1 个分区)以将所有数据全局收集到一个地方。"这可以工作(尽管我们必须手动创建主题,因为我们的集群没有启用自动创建)。正确阅读时需要汇总该主题吗?我看不出这怎么能被归类为并发,所以不明白为什么不能将它包装在一个不错的 Kafka Streams API 中。
  • 我不是说,它不能用一个好的 API 来包装——我的回答是针对你需要用当前 API 做的事情。随时提交功能请求 Jira :)
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 2016-09-30
  • 2017-02-21
  • 2018-07-15
  • 2019-08-07
  • 2016-06-27
  • 2019-12-15
  • 2019-03-30
相关资源
最近更新 更多