【问题标题】:A windowed aggregation on event count事件计数的窗口聚合
【发布时间】:2020-07-11 01:53:03
【问题描述】:

我已经对我的 kafka 事件进行了分组:

    private static void createImportStream(final StreamsBuilder builder, final Collection<String> topics) {
        final KStream<byte[], GraphEvent> stream = builder.stream(topics, Consumed.with(Serdes.ByteArray(), new UserEventThriftSerde()));
        stream.filter((key, request) -> {
            return Objects.nonNull(request);
        }).groupBy(
                (key, value) -> Integer.valueOf(value.getSourceType()),
                Grouped.with(Serdes.Integer(), new UserEventThriftSerde()))
              .aggregate(ArrayList::new, (key, value, aggregatedValue) -> {
                          aggregatedValue.add(value);
                          return aggregatedValue;
                      },
                      Materialized.with(Serdes.Integer(), new ArrayListSerde<UserEvent>(new UserEventThriftSerde()))
              ).toStream();
    }

如何添加window 但不是基于时间,而是基于事件数量。 原因是事件将是批量转储,时间窗口聚合不适合,因为所有事件都可能在相同的几秒钟内出现。

【问题讨论】:

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


    【解决方案1】:

    Kafka Streams 不支持开箱即用的基于计数的窗口,因为这些窗口是不确定的,并且很难处理乱序数据。

    您可以使用处理器 API 为您的用例构建自定义运算符,而不是使用 DSL。

    【讨论】:

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