如何消除Spring Cloud Kafka Binder中IndicativeQueryService的存储未找到警告?
我们的项目基于spring-cloud-stream-binder-kafka-streams依赖实现Kafka跨服务通信,application.yaml中的监听器配置如下:
spring: cloud: kafka: stream: bindings: listenToBookTopic-in-0: consumer: applicationId: someAppId listenToUserTopic-in-0: consumer: applicationId: someAppId2 materializedAs: user-store listenToRolesTopic-in-0: consumer: applicationId: someAppId3 materializedAs: role-store
处理Book对象的Kafka监听器代码:
@Bean public Function<KStream<String, OldBook>, KStream<String, NewBook>> listenToTopic() { return stream -> stream .map((key, oldBook) -> { try { final var newBook = bookService.processBook(oldBook); final String newKey = createNewKey(newBook); return KeyValue.pair("", null); } catch (Exception e) { log.error(e.getMessage()); return KeyValue.pair("", null); } }); }
Book Service代码:
@Autowired private final StateStoreService stateStoreService; public NewBook processBook(OldBook oldBook) { final var newBook = createBook(oldBook); newBook.setUser(stateStoreService.getUserById(oldBook.getId())); newBook.setRoles(stateStoreService.getRolesById(oldBook.getId())); return newBook; }
StateStoreService代码:
@Autowired private final InteractiveQueryService interactiveQueryService; public Optional<User> getUserById(final String id) { final ReadOnlyKeyValueStore<String, User> userStore = interactiveQueryService.getQueryableStore( properties.getUserStore(), QueryableStoreTypes.keyValueStore()); return Optional.ofNullable(userStore.get(id)); } public Optional<Role> getRoleById(final String id) { final ReadOnlyKeyValueStore<String, Role> store = interactiveQueryService .getQueryableStore(properties.getRoleStore(), QueryableStoreTypes.keyValueStore()); return Optional.ofNullable(store.get(id)); }
问题现象
从正常请求上下文调用StateStoreService.getUserById()时,无任何警告;但通过listenToBookTopic bean触发调用时,每接收一条Book消息就会弹出警告:
"Store storeName could not be found in Streams context, falling back to all known Streams instances"
排查结论
listenToBookTopic对应的Kafka Streams实例运行在独立线程中,调用getThreadContextSpecificKafkaStreams()会获取该线程绑定的上下文,而此实例并不包含user-store和role-store(这两个物化视图属于其他applicationId的Streams实例),因此Spring会触发全局回退查询,进而抛出警告。
方法1:统一所有绑定的applicationId
将所有KStream绑定的applicationId设置为同一个值,让所有Streams任务共享同一上下文,StateStore会对所有线程可见:
spring: cloud: kafka: stream: bindings: listenToBookTopic-in-0: consumer: applicationId: shared-app-id listenToUserTopic-in-0: consumer: applicationId: shared-app-id materializedAs: user-store listenToRolesTopic-in-0: consumer: applicationId: shared-app-id materializedAs: role-store
优势:实现简单,无需修改业务代码,彻底消除上下文不匹配问题。
方法2:显式指定查询目标的applicationId
使用InteractiveQueryService的重载方法,直接指定StateStore所属的applicationId,跳过线程上下文自动查找逻辑:
public Optional<User> getUserById(final String id) { final ReadOnlyKeyValueStore<String, User> userStore = interactiveQueryService.getQueryableStore( properties.getUserStore(), QueryableStoreTypes.keyValueStore(), "someAppId2"); // 指定user-store对应的applicationId return Optional.ofNullable(userStore.get(id)); } public Optional<Role> getRoleById(final String id) { final ReadOnlyKeyValueStore<String, Role> store = interactiveQueryService .getQueryableStore( properties.getRoleStore(), QueryableStoreTypes.keyValueStore(), "someAppId3"); // 指定role-store对应的applicationId return Optional.ofNullable(store.get(id)); }
优势:保留多applicationId的架构设计,精准定位目标StateStore实例,避免全局回退查询。
方法3:将StateStore查询嵌入KStream拓扑
遵循Kafka Streams原生设计,将StateStore查询逻辑直接放到流处理拓扑中,确保与流处理线程上下文一致:
@Bean public Function<KStream<String, OldBook>, KStream<String, NewBook>> listenToTopic() { return stream -> stream .transform(() -> new Transformer<String, OldBook, KeyValue<String, NewBook>>() { private ReadOnlyKeyValueStore<String, User> userStore; private ReadOnlyKeyValueStore<String, Role> roleStore; @Override public void init(ProcessorContext context) { // 从当前流处理上下文直接获取StateStore userStore = context.getStateStore("user-store"); roleStore = context.getStateStore("role-store"); } @Override public KeyValue<String, NewBook> transform(String key, OldBook oldBook) { try { NewBook newBook = createBook(oldBook); newBook.setUser(Optional.ofNullable(userStore.get(oldBook.getId()))); newBook.setRoles(Optional.ofNullable(roleStore.get(oldBook.getId()))); String newKey = createNewKey(newBook); return KeyValue.pair(newKey, newBook); } catch (Exception e) { log.error(e.getMessage()); return KeyValue.pair("", null); } } @Override public void close() {} }, "user-store", "role-store"); // 声明依赖的StateStore }
优势:完全符合Kafka Streams设计规范,性能最优,从根源上避免跨上下文查询问题。
内容的提问来源于stack exchange,提问作者Dima

