【问题标题】:KStream batch process windowsKStream 批处理窗口
【发布时间】:2016-12-30 11:16:43
【问题描述】:

我想用 KStream 接口批量发送消息。

我有一个带有键/值的流 我试图在一个滚动窗口中收集它们,然后我想立即处理整个窗口。

builder.stream(longSerde, updateEventSerde, CONSUME_TOPIC)
                .aggregateByKey(
                        HashMap::new,
                        (aggKey, value, aggregate) -> {
                            aggregate.put(value.getUuid, value);
                            return aggregate;
                        },
                        TimeWindows.of("intentWindow", 100),
                        longSerde, mapSerde)
                .foreach((wk, values) -> {

事情是每次更新 KTable 时都会调用 foreach。 我想在完成后处理整个窗口。就像从 100 毫秒收集数据然后立即处理一样。各有千秋。

16:** - windows from 2016-08-23T10:56:26 to 2016-08-23T10:56:27, key 2016-07-21T14:38:16.288, value count: 294
16:** - windows from 2016-08-23T10:56:26 to 2016-08-23T10:56:27, key 2016-07-21T14:38:16.288, value count: 295
16:** - windows from 2016-08-23T10:56:26 to 2016-08-23T10:56:27, key 2016-07-21T14:38:16.288, value count: 296
16:** - windows from 2016-08-23T10:56:26 to 2016-08-23T10:56:27, key 2016-07-21T14:38:16.288, value count: 297
16:** - windows from 2016-08-23T10:56:26 to 2016-08-23T10:56:27, key 2016-07-21T14:38:16.288, value count: 298
16:** - windows from 2016-08-23T10:56:26 to 2016-08-23T10:56:27, key 2016-07-21T14:38:16.288, value count: 299
16:** - windows from 2016-08-23T10:56:27 to 2016-08-23T10:56:28, key 2016-07-21T14:38:16.288, value count: 1
16:** - windows from 2016-08-23T10:56:27 to 2016-08-23T10:56:28, key 2016-07-21T14:38:16.288, value count: 2
16:** - windows from 2016-08-23T10:56:27 to 2016-08-23T10:56:28, key 2016-07-21T14:38:16.288, value count: 3
16:** - windows from 2016-08-23T10:56:27 to 2016-08-23T10:56:28, key 2016-07-21T14:38:16.288, value count: 4
16:** - windows from 2016-08-23T10:56:27 to 2016-08-23T10:56:28, key 2016-07-21T14:38:16.288, value count: 5
16:** - windows from 2016-08-23T10:56:27 to 2016-08-23T10:56:28, key 2016-07-21T14:38:16.288, value count: 6

在某些时候,新窗口从地图中的 1 个条目开始。 所以我什至不知道什么时候窗口是满的。

关于在 kafka 流中进行批处理的任何提示

【问题讨论】:

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


【解决方案1】:

现在(从 Kafka 0.10.0.0 / 0.10.0.1 开始):您所描述的窗口行为是“按预期工作”。也就是说,如果您收到 1,000 条传入消息,您将(当前)始终看到 1,000 条最新版本的 Kafka / Kafka Streams 在下游进行更新。

展望未来:Kafka 社区正在开发新功能,以使这种更新率行为更加灵活(例如,允许您在上面描述的行为作为您想要的行为)。详情请见KIP-63: Unify store and downstream caching in streams

【讨论】:

    【解决方案2】:

    ====== 更新 ======

    在进一步测试中,这不起作用。 正确的方法是使用@friedrich-nietzsche 概述的处理器。我对自己的答案投了反对票.... grrrr.

    ====================

    我仍在与这个 API 搏斗(但我喜欢它,所以它的时间是值得花的 :)),我不确定你想从你的代码示例结束的地方在下游完成什么,但它看起来类似于我得到了什么工作。高级别是:

    从源读取的对象。它代表一个键和 1:∞ 事件数,我想每 5 秒发布一次每个键的事件总数(或 TP5s,每 5 秒的事务)。代码的开头看起来一样,但我使用的是:

    1. KStreamBuilder.stream
    2. reduceByKey
    3. 到一个窗口(5000)
    4. 发送到new stream,它每 5 秒获取每个键的累积值。
    5. map 流到一个新的 KeyValue 每个键
    6. to水槽话题。

    在我的情况下,每个窗口期,我可以将所有事件减少到每个键一个事件,所以这是可行的。如果您想保留每个窗口的所有单个事件,我假设可以使用 reduce 将每个实例映射到实例集合(可能使用相同的键,或者您可能需要一个新键)并在每个窗口期结束时,下游流将一次性获得一堆您的事件集合(或者可能只是所有事件的一个集合)。它看起来像这样,经过消毒和 Java 7-ish:

        builder.stream(STRING_SERDE, EVENT_SERDE, SOURCE_TOPICS)
            .reduceByKey(eventReducer, TimeWindows.of("EventMeterAccumulator", 5000), STRING_SERDE, EVENT_SERDE)            
            .toStream()
            .map(new KeyValueMapper<Windowed<String>, Event, KeyValue<String,Event>>() {
                public KeyValue<String, Event> apply(final Windowed<String> key, final Event finalEvent) {
                    return new KeyValue<String, Event>(key.key(), new Event(key.window().end(), finalEvent.getCount());
                }
        }).to(STRING_SERDE, EVENT_SERDE, SINK_TOPIC);
    

    【讨论】:

    • 所以这将每5秒输出1个值/键?我在第 4 步上仍然不是 100%。从我看到的情况来看,你会在那里得到一系列变化。或者这就是你的意思?
    • 我每 5 秒得到一个键/值。该值表示该时间段内的事件总数(每个源事件可以表示多个“子事件”)。但是如果你创建了一个聚合对象,比如像 Map> 这样的单个地图,你可以向下游传递一个对象。
    • toStream() 不是又把一切搞砸了。 reduceByKey 为您提供了一个 KTable,但没有任何缓存或 Buffer 再次将其转换为 KStream 只会为您对键的每个井更新产生一个更新事件。我想我可能不得不在处理窗口时在外面检测我的批次,只要 windowKey 更改我现在最后一个翻滚窗口已经完成并且我的批次已满。我不知道它对垃圾收集有什么影响,因为我不确定聚合是否始终是同一个对象或存在深层副本。
    • 公平的问题。需要一个实际的代码示例。要点马上就来了。
    • 已将答案更新为非答案。这是行不通的。处理器方法是要走的路。
    【解决方案3】:

    我的实际任务是将更新从流推送到 redis,但我不想单独读取/更新/写入,即使 redis 速度很快。 我现在的解决方案是使用 KStream.process() 提供一个处理器,该处理器添加到进程队列并实际处理队列。

    public class BatchedProcessor extends AbstractProcessor{
    
    ...
    BatchedProcessor(Writer writer, long schedulePeriodic)
    
    @Override
    public void init(ProcessorContext context) {
        super.init(context);
        context.schedule(schedulePeriodic);
    }
    
    @Override
    public void punctuate(long timestamp) {
        super.punctuate(timestamp);
        writer.processQueue();
        context().commit();
    }
    
    @Override
    public void process(Long aLong, IntentUpdateEvent intentUpdateEvent) {
        writer.addToQueue(intentUpdateEvent);
    }
    

    我仍然需要测试,但它解决了我遇到的问题。人们可以很容易地以一种非常通用的方式编写这样的处理器。 API 非常整洁干净,但是一个 processBatched((List batchedMessaages)-> ..., timeInterval OR countInterval) 只使用标点符号来处理批处理并在此时提交并在 Store 中收集 KeyValues 可能是一个有用的补充。

    但也许它的目的是使用处理器来解决这个问题,并将 API 保持在一次只消息中的低延迟焦点中。

    【讨论】:

    • 您如何处理稍后到达(即乱序)的记录?
    • 一些用例可能不关心迟到的数据(例如物联网传感器数据)。因此,如果能以这种方式调整流,那就太好了。
    • 如果您使用进程添加大量消息,提交消费者偏移量,然后应用程序在调用 punctuate 之前死掉,会发生什么情况?这些消息不是永远丢失了吗?
    • 我认为@bm1729 就在这里。看起来你会让自己暴露在丢失的数据中。
    猜你喜欢
    • 2017-06-02
    • 2019-05-06
    • 2016-03-18
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2018-09-18
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多