【问题标题】:Kafka PersistentWindowStore rebalancing mechanicsKafka PersistentWindowStore 再平衡机制
【发布时间】:2018-10-04 15:30:14
【问题描述】:

我正在松散地基于 this confluent code 为 Kafka Streams 应用程序创建一个 30 分钟的重复数据删除存储(以解决与 Kafka 的一次性处理保证不同的问题),并希望最大限度地减少拓扑启动时间。

此代码使用持久窗口存储,这要求我指定要使用的日志段的数量。假设我想使用 2 个段,并且使用 1GB 的默认段大小,这是否意味着在重新平衡期间,客户端必须在应用程序启动之前读取 2GB 的数据?

【问题讨论】:

    标签: apache-kafka apache-kafka-streams confluent-platform


    【解决方案1】:

    segment 参数在 Kafka Streams 中配置了一些不同的东西——它与代理中的段无关(只是同名)。

    使用窗口存储,存储的保留时间除以段数。如果所有数据段都比保留时间旧,则删除整个段并创建一个新的空段。这些段,只存在于客户端。

    需要恢复的记录数,仅取决于保留时间(和您的输入数据速率)。它与段大小无关。段大小仅定义了细粒度的旧记录过期的程度。

    【讨论】:

    • WindowStores 默认创建为 3 个段,AFAIK。使用更多段创建它们以更快地清除保留是否有意义(更多段,所有记录现在过期的可能性更大)或在读/写操作上使用多个段的过载导致保守的段数(例如 3 代替比方说 10)?
    • segment的本地开销比较大。在当前的 impl 中,我们为每个段创建一个 RocksDB 实例。因此,我们尝试通过每个默认仅使用 3 个段来保持这种开销较小。
    猜你喜欢
    • 2018-11-03
    • 1970-01-01
    • 2019-02-27
    • 1970-01-01
    • 2020-06-30
    • 2021-09-05
    • 2020-12-31
    • 2015-04-18
    • 1970-01-01
    相关资源
    最近更新 更多