in core/src/main/java/org/apache/rocketmq/streams/core/running/WorkerThread.java [228:240]
void doCommit(HashSet<MessageQueue> set) throws Throwable {
if (set != null && set.size() != 0) {
this.stateStore.persist(set);
this.unionConsumer.commit(set, true);
for (MessageQueue messageQueue : set) {
logger.debug("committed messageQueue: [{}]", messageQueue);
}
set.clear();
}
}