【问题标题】:KTable state store persistenceKTable 状态存储持久化
【发布时间】:2018-12-28 07:09:25
【问题描述】:

如果我在实现 KTable 时使用持久存储,状态存储是否会在应用程序重新启动时保持不变?例如,如果我使用以下内容:

StreamsBuilder builder = new StreamsBuilder();
KeyValueBytesStoreSupplier storeSupplier =      Stores.persistentKeyValueStore("queryable-store-name");
 KTable<Long,String> table = builder.table(
   "foo",
   Materialized.as(storeSupplier)
               .withKeySerde(Serdes.Long())
               .withValueSerde(Serdes.String())

状态存储“queryable-store-name”是否可以在重新启动时使用先前运行的状态访问?可以说,我向主题 foo 发送了 50 条记录,它在状态存储中实现了。然后应用程序重新启动,我仍然在状态存储中保留这 50 条记录吗?如果没有,有没有办法做到这一点?

谢谢!

【问题讨论】:

    标签: apache-kafka-streams


    【解决方案1】:

    是的,状态存储默认保存在磁盘上。当应用程序重新启动并且application-id 未更改时,将从磁盘恢复状态,包含所有 50 条记录。当应用程序被杀死/停止/重新启动时,将从偏移量添加新记录。

    编辑: 好像您缺少 KTable 之上的聚合操作,这是必需的。请参阅我的代码示例:

    final KStream<CustomerKey, ViewPage> viewPagesStream=builder.stream(INPUT_TOPIC);
    
    final KTable<Windowed<ViewPageCountKey>,Long>uniqueViewPageCount=viewPagesStream
            .map((key,value)->{
                ViewPageCountKey newKey=new ViewPageCountKey(key.getProjectId(),value.getUrl());
                return new KeyValue<>(newKey,value);
            })
            .filter((key,value)->key!=null)
            .groupByKey()
            .count(TimeWindows.of(WINDOW_SIZE).advanceBy(WINDOW_ADVANCE),STORE_NAME);
    

    【讨论】:

    • 感谢您的回答。当我从以前的数据重新启动时,我是否能够查询状态存储?我没有看到你描述的那样工作,也许我做错了什么。
    • 旁注:即使您使用内存存储,或者在磁盘上松散状态,它也会在处理恢复之前从更新日志主题中恢复。只要启用了日志记录,您的商店就是完全容错的。请注意,未启用日志记录的持久存储不是完全容错的。持久存储的想法是允许大于主内存的状态和更快的启动时间,因为不需要从更改日志主题重建存储。但是,磁盘上的本地存储数据不是出于容错原因写入的——这就是更改日志主题的目的。
    • 感谢@MatthiasJ.Sax。我在这里有一个后续问题:stackoverflow.com/questions/51461416/…。如果您能对此表示赞赏,将不胜感激。
    猜你喜欢
    • 1970-01-01
    • 2023-04-08
    • 2020-04-17
    • 2011-01-09
    • 1970-01-01
    • 1970-01-01
    • 2021-05-20
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多