【问题标题】:Kafka Streams - Persist messages by timestamp/sequence?Kafka Streams - 按时间戳/序列保留消息?
【发布时间】:2018-06-18 21:09:55
【问题描述】:

我正在通过 Kafka 流接收消息。它们由用户 ID 键入。他们进来时会得到一个序列号和时间戳。消息在 15 分钟后“过期”。用户可以根据给定的时间(最多 15 分钟)或顺序请求新消息。

我最初拥有的是这样的:

` StreamsBuilder streamsBuilder = new StreamsBuilder();

  KStream<String, Message> inboundStream = streamsBuilder.stream("incoming.topic");
  messageSupplier = Stores.persistentKeyValueStore("user.messages");

  KTable<String, MessageCache> messageTable = inboundStream
      .filter(this::userExists)
      .peek(this::recordInboundMessage)
      .map(this::markMessage)       // add sequence/timestamp
      .groupByKey()
      .aggregate(this::createMessageCache,
              this::addMessageToMessageCache,
              Materialized.as(messageSupplier));

  // ---> Some other setup stuff, then start the streams

`

MessageCache 保存消息列表(当我们将消息添加到缓存时删除过期消息)。当我收到消息请求时,我会浏览列表并过滤掉相应的消息。

我在想我可以使用其中一种窗口策略,但找不到实际持久化消息列表的示例。

这是最好的方法吗?还是我错过了一些能让这更容易/更好的东西?

【问题讨论】:

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


    【解决方案1】:

    这是最好的方法吗?还是我错过了一些能让这更容易/更好的东西?

    我认为您有一个使用本机 Java 类的简单解决方案,可以有效地将流应用程序与您的代码连接起来......为了简单起见,有很多话要说!我能看到的唯一缺点是,如果您的事件率太高,您的用户缓存可能会增长到超出您的内存大小。此外,如果您需要容错,流式应用程序将在另一个应用程序实例上重建状态存储的内容,以防万一发生故障。但如果这不是问题,那就去吧!

    但是,就如何在流应用上下文中执行此操作而言,您可以进行一些调整:

    1. 定义您希望支持的用户查询的粒度。分钟?秒?为了争论,我们说几分钟。根据该粒度窗口化您的流。

    2. 定义一个累加器,类似于您所拥有的,它将接受用户记录并将其添加到列表中。类似于UserRecordGroup 的东西,它有一个ListUserRecord,还有一个方法add(UserEvent evt),它将一个UserRecord 附加到List

    然后,您可以像这样构建您的流应用程序:

    KStream<String, Message> inboundStream = streamsBuilder.stream("incoming.topic");
     Materialized<String, UserRecordGroup, WindowStore<Bytes, byte[]>> userStore =
     Materialized.<String, UserRecordGroup, WindowStore<Bytes,byte[]>>as("user.messages")
      .withValueSerde(/* your serializers here */);
    
    
    KTable<String, MessageCache> messageTable = inboundStream
      .filter(this::userExists)
      .peek(this::recordInboundMessage)
      .map(this::markMessage)       // add sequence/timestamp
      .groupByKey()
      .windowedBy(TimeWindows.of(ONE_MINUTE_IN_MS))
      .aggregate(UserRecordGroup::new,
                (key, value, agg) -> { agg.add(value); },
                 userStore);
    

    最后,当你想从商店中提供查询时,你可以

    ReadOnlyWindowStore<Integer, UserRecordGroup> store =
       streams.store("user.messages", QueryableStoreTypes.windowStore());
    WindowStoreIterator<UserRecordGroup> windowIterator = 
         store.fetch(pathHash, startTimestamp, endTimeStamp);
    

    您的迭代器将包含不同窗口的所有记录的列表;将这些列表合并在一起,您就有了 startTimestamp 和 endTimestamp 之间用户活动的描述。

    【讨论】:

    • 这很有趣。感谢您的输入!在这种情况下,旧窗户会自行清除,还是仍然存在? (我必须做些什么吗?)还是与主题的保留配置有关?有些用户会非常频繁地收到消息,有些则不会。您提出了一个有用的观点,可能是多个应用程序实例。消息的调用是一个长拉,所以如果不存在消息,它会等待消息。我想在收到消息时通知线程,但如果调用是在不同的实例上,有没有办法做到这一点?
    猜你喜欢
    • 2020-12-21
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2021-07-13
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2021-08-17
    相关资源
    最近更新 更多