【发布时间】: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