【问题标题】:How to add a cooldown/rate-limit to a stream in Kafka Streams?如何为 Kafka Streams 中的流添加冷却时间/速率限制?
【发布时间】:2017-06-18 21:25:14
【问题描述】:

我是流式数据处理的新手,我觉得必须是一个非常基本的用例。

假设我有一个(User, Alert) 元组流。我想要的是对每个用户的流进行速率限制。 IE。我想要一个只为用户输出一次警报的流。在接下来的 60 分钟内,用户的任何传入警报都应该被吞下。在这 60 分钟之后,应再次触发传入警报。

我尝试了什么:

使用aggregate 作为有状态转换,但聚合状态是时间相关的。然而,即使生成的KTable 的聚合值没有变化,KTable(作为更改日志)将继续向下发送元素,因此无法达到“限速”流的预期效果

val fooStream: KStream[String, String] = builder.stream("foobar2")
fooStream
  .groupBy((key, string) => string)
  .aggregate(() => "constant",
    (aggKey: String, value: String, aggregate: String) => aggregate,
    stringSerde,
    "name")
  .print

提供以下输出:

[KSTREAM-AGGREGATE-0000000004]: string , (constant<-null)
[KSTREAM-AGGREGATE-0000000004]: string , (constant<-null)

我通常不清楚aggregate 如何/何时决定向下游发布元素。我最初的理解是它是即时的,但似乎并非如此。据我所知,开窗在这里应该没有帮助。

有没有可能是 Kafka Streams DSL 目前没有考虑这种状态转换用例,类似于 Spark 的 updateStateByKey 或 Akka 的 statefulMapConcat?我是否必须使用较低级别的处理器/变压器 API?

编辑:

Possible duplicate 确实解决了记录缓存如何导致聚合何时决定向下游发布元素的一些混乱的问题。然而,主要问题是如何在 DSL 中实现“速率限制”。正如@miguno 指出的那样,必须恢复到较低级别的处理器 API。下面我粘贴了相当冗长的方法:

  val logConfig = new util.HashMap[String, String]();
  // override min.insync.replicas
  logConfig.put("min.insyc.replicas", "1")

  case class StateRecord(alert: Alert, time: Long)

  val countStore = Stores.create("Limiter")
    .withKeys(integerSerde)
    .withValues(new JsonSerde[StateRecord])
    .persistent()
    .enableLogging(logConfig)
    .build();
  builder.addStateStore(countStore)

  class RateLimiter extends Transformer[Integer, Alert, KeyValue[Integer, Alert]] {
    var context: ProcessorContext = null;
    var store: KeyValueStore[Integer, StateRecord] = null;

    override def init(context: ProcessorContext) = {
      this.context = context
      this.store = context.getStateStore("Limiter").asInstanceOf[KeyValueStore[Integer, StateRecord]]
    }

    override def transform(key: Integer, value: Alert) = {
      val current = System.currentTimeMillis()
      val newRecord = StateRecord(value._1, value._2, current)
      store.get(key) match {
        case StateRecord(_, time) if time + 15.seconds.toMillis < current => {
          store.put(key, newRecord)
          (key, value)
        }
        case StateRecord(_, _) => null
        case null => {
          store.put(key, newRecord)
          (key, value)
        }
      }
    }
  }

【问题讨论】:

标签: scala stream apache-kafka apache-kafka-streams


【解决方案1】:

假设我有一个 (User, Alert) 元组流。我想要的是对每个用户的流进行速率限制。 IE。我想要一个只为用户输出一次警报的流。在接下来的 60 分钟内,用户的任何传入警报都应该被吞下。在这 60 分钟之后,应再次触发传入警报。

目前在使用 Kafka Streams 的 DSL 时这是不可能的。相反,您可以(并且需要)使用较低级别的处理器 API 手动实现此类行为。

仅供参考:我们一直在 Kafka 社区讨论是否将此类功能(通常称为“触发器”)添加到 DSL。到目前为止,我们决定暂时不使用此类功能。

我通常不清楚aggregate 如何/何时决定向下游发布元素。我最初的理解是它是立竿见影的,但似乎并非如此。

是的,这是 Kafka 0.10.0.0 的初始行为。从那时起(不确定您使用的是什么版本)我们引入了记录缓存;如果您禁用记录缓存,您将恢复初始行为,但据我所知,记录缓存会给您某种(间接)速率限制旋钮。因此,您可能希望保持启用缓存。

不幸的是,Apache Kafka 文档尚未涵盖记录缓存,与此同时,您可能想改为阅读 http://docs.confluent.io/current/streams/developer-guide.html#memory-management

【讨论】:

  • 感谢您及时的评论。我怀疑这在 DSL 中目前是不可能的。您是否有指向社区围绕此主题进行的讨论的链接?虽然处理器 API 已经足够好,但它确实看起来像是一个可以抽象的常见用例(就像其他框架一样)。我用我的解决方法更新了问题
  • 讨论:我不记得了。很可能在 kafka-dev 邮件列表中的讨论中。
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 2021-04-21
  • 1970-01-01
  • 2022-11-16
  • 2022-12-10
  • 1970-01-01
  • 2021-09-07
  • 2021-10-14
相关资源
最近更新 更多