【问题标题】:Kafka state-store on different scaled instances不同规模实例上的 Kafka 状态存储
【发布时间】:2020-01-07 07:28:00
【问题描述】:

我有 5 台不同的机器,每个机器都有 5 个使用 kafka-streams 应用程序的缩放的 spring boot 实例。我正在使用具有不同 2-3 个主题的 50 个分区压缩主题,并且我的每个实例都有 10 个并发。我正在使用 docker swarm 和 docker volume。使用这些主题 KTable 或 KStream 使用我的 kafka 流应用程序执行一些 flatMap、映射和连接操作。

    props.put(StreamsConfig.STATE_DIR_CONFIG, /tmp/kafka-streams);
    props.put(StreamsConfig.REPLICATION_FACTOR_CONFIG, 3);
    props.put(StreamsConfig.NUM_STANDBY_REPLICAS_CONFIG, 2);
    props.put(StreamsConfig.COMMIT_INTERVAL_MS_CONFIG, 100);
    props.put(StreamsConfig.PROCESSING_GUARANTEE_CONFIG, EXACTLY_ONCE);
    props.put("num.stream.threads", 10);
    props.put("application.id", applicationId);

如果一切正常,在我的应用程序中使用 .join() 操作没有任何问题或没有数据丢失,但是当我的一个实例关闭时,我的连接操作实际上无法执行连接。

我的问题是:当应用程序重新启动或重新部署(并且考虑到它在非持久性容器中工作)时,它的状态是否消失了?比我的加入操作不起作用。当我重新部署我的实例并使用最新实体从 elasticsearch 填充我的压缩主题时,我的连接操作就可以了。所以我认为当我的应用程序在新机器上启动时,我的本地状态商店不见了?但是kafka文档说:

如果任务在失败的机器上运行并在另一台机器上重新启动,则 Kafka Streams 通过在恢复处理新启动的任务之前重播相应的变更日志主题,保证将其关联的状态存储恢复到失败之前的内容。因此,故障处理对最终用户是完全透明的。 请注意,任务(重新)初始化的成本通常主要取决于通过重放状态存储的关联更改日志主题来恢复状态的时间。为了最大限度地缩短恢复时间,用户可以将他们的应用程序配置为具有本地状态的备用副本(即状态的完全复制副本)。当发生任务迁移时,Kafka Streams 会尝试将任务分配给已经存在此类备用副本的应用程序实例,以最小化任务(重新)初始化成本。请参阅 Kafka Streams 配置部分的 num.standby.replicas。 (https://kafka.apache.org/0102/documentation/streams/architecture)

我宕机的实例在启动时会刷新 kafka 状态存储吗?如果这就是我丢失数据并且我不知道的原因:/或者因为我的所有实例都使用相同的 applicationId 而因为 commit_offset 而无法重新加载状态存储?

谢谢!

【问题讨论】:

    标签: spring-boot apache-kafka apache-kafka-streams


    【解决方案1】:

    更改日志主题总是从最早的偏移量读取,并且它们被压缩,所以它们不会丢失数据。

    如果您要加入非紧凑主题,那么肯定会丢失数据,但这不仅限于 Kafka Streams 或您的特定用例...您需要配置主题以至少保留数据尽可能长的时间正如您认为的那样,它将带您解决主题停机时间的任何问题。在保留数据的同时,您可以随时向您的消费者寻求它

    如果您想要持久存储,例如,通过 Kubernetes 将卷挂载到您的容器,或者插入存储在容器外部的状态状态存储,例如 Redis:https://github.com/andreas-schroeder/redisks

    【讨论】:

    • 对不起,误解,我一直使用压缩主题,并在压缩主题之间进行连接操作。我正在使用已安装的 docker 卷。可以说 5 个服务器 5 个实例.. 每个实例的分区分配为 0..10, - 10 - 20, 20-30, 30 - 40, 40-50。当第五个实例在这个分区间隔内死亡(40-50)时,4个实例将共享分区对吗?但是第五台机器有卷和它的状态存储。分区分布在 4 个实例之间,所以我丢失了数据? :/ 如果你和卡夫卡文件说我不应该丢失数据,但我确实这样做了。我是否缺少配置
    • 您可以配置备用副本,但这并不能真正解决您的容器在没有任何数据的情况下停止和移动机器的问题。似乎您可能正在某处的生产者中丢弃消息,或者使用没有事务的低版本流。整个组应该重新平衡,所以你不能保证 10 个分区会被均匀地重新分配。 docker 上的 Kafka Streams 基本上需要持久化卷,这已经在文档之外说明(因为 kafka 不必说任何关于 docker 的具体内容)youtu.be/pFZBs_8hmyo
    猜你喜欢
    • 1970-01-01
    • 2019-02-13
    • 1970-01-01
    • 2023-03-28
    • 1970-01-01
    • 2019-07-03
    • 1970-01-01
    • 2014-09-19
    • 2021-09-27
    相关资源
    最近更新 更多