【问题标题】:Kafka Stream count on time window not reporting zero valuesKafka Stream 时间窗口计数不报告零值
【发布时间】:2017-07-14 07:52:19
【问题描述】:

我正在使用 Kafka 流来使用跳跃时间窗口计算过去 3 分钟内发生了多少事件:

public class ViewCountAggregator {

    void buildStream(KStreamBuilder builder) {      

        final Serde<String> stringSerde = Serdes.String();
        final Serde<Long> longSerde = Serdes.Long();

        KStream<String, String> views = builder.stream(stringSerde, stringSerde, "streams-view-count-input");
        KStream<String, Long> viewCount = views
            .groupBy((key, value) -> value)
            .count(TimeWindows.of(TimeUnit.MINUTES.toMillis(3)).advanceBy(TimeUnit.MINUTES.toMillis(1)))
            .toStream()
            .map((key, value) -> new KeyValue<>(key.key(), value));

        viewCount.to(stringSerde, longSerde, "streams-view-count-output");        
    }

    public static void main(String[] args) throws Exception {                   
        // some not so important initialization code
        ...  
    }

}

当运行消费者并将一些消息推送到输入主题时,随着时间的推移,它会收到以下更新:

single  1
single  1
single  1
five    1
five    4
five    5
five    4
five    1

这几乎是正确的,但它从未收到以下更新:

single  0
five    0

如果没有它,我更新计数器的消费者永远不会在较长时间没有事件时将其设置回零。我希望消费的消息看起来像这样:

single  1
single  1
single  1
single  0
five    1
five    4
five    5
five    4
five    1
five    0

是否有一些我遗漏的配置选项/参数可以帮助我实现这种行为?

【问题讨论】:

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


    【解决方案1】:

    这几乎是正确的,但它从未收到以下更新:

    首先,计算出的输出是正确的。

    二、为什么正确:

    如果您应用窗口聚合,则只会创建那些确实具有实际内容的窗口(我熟悉的所有其他系统都会产生相同的输出)。因此,如果对于某个键,在超过窗口大小的时间段内没有数据,则没有实例化窗口,因此也根本没有计数。

    如果没有内容则不实例化窗口的原因很简单:处理器无法知道所有键。在您的示例中,您有两个键,但稍后可能会出现第三个键。你希望从一开始就得到&lt;thirdKey,0&gt; 吗?此外,由于数据流本质上是无限的,因此密钥可能会消失并且永远不会再次出现。如果你记得所有见过的键,并且如果没有消失的键的数据就发出&lt;key,0&gt;,你会永远发出&lt;key,0&gt;吗?

    我不想说您的预期结果/语义没有意义。这只是您的一个非常具体的用例,一般不适用。因此,流处理器没有实现它。

    第三:你能做什么?

    有多种选择:

    1. 您的消费者可以跟踪它看到的键,并使用嵌入的记录时间戳确定键是否“丢失”,然后将此键的计数器设置为零(为此,它也可能有助于删除map 步骤并保留 Windowed&lt;K&gt; 类型作为键,以便消费者获取记录所属窗口的信息)
    2. 在您的 Stream 应用程序中添加一个有状态的 #transform() 步骤,该步骤与 (1) 中描述的相同。为此,注册一个标点符号回调可能会有所帮助。

    方法 (2) 应该更容易跟踪密钥,因为您可以将状态存储附加到转换步骤,因此不需要处理下游使用者中的状态(和故障/恢复)。

    但是,这两种方法的棘手部分仍然是确定 何时 缺少密钥,即您要等多长时间才能生成 &lt;key,0&gt;。请注意,数据可能迟到(也就是乱序),即使您确实发出了 &lt;key,0&gt;,迟到的记录也可能会在您的代码发出 @ 之后 产生 &lt;key,1&gt; 消息987654332@记录。但也许这对你的情况来说并不是一个真正的问题,因为你似乎只使用最新的窗口。

    最后但并非最不重要的另一条评论:您似乎只使用最新计数,并且较新的窗口会覆盖下游消费者中的旧窗口。因此,可能值得探索“交互式查询”以直接利用 count 运算符的状态,而不是使用主题并更新其他状态。这可能允许您重新设计和显着简化下游应用程序。查看docs 和一篇关于Interactive Queries 的非常好的博文了解更多详情。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2019-06-10
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2017-10-01
      相关资源
      最近更新 更多