diff --git a/broker/src/main/java/org/apache/rocketmq/broker/config/v1/RocksDBConsumerOffsetManager.java b/broker/src/main/java/org/apache/rocketmq/broker/config/v1/RocksDBConsumerOffsetManager.java index 2d015afca3..4f63516777 100644 --- a/broker/src/main/java/org/apache/rocketmq/broker/config/v1/RocksDBConsumerOffsetManager.java +++ b/broker/src/main/java/org/apache/rocketmq/broker/config/v1/RocksDBConsumerOffsetManager.java @@ -157,17 +157,21 @@ public class RocksDBConsumerOffsetManager extends ConsumerOffsetManager { @Override public synchronized void persist() { - try (WriteBatch writeBatch = new WriteBatch()) { - for (Entry> entry : this.offsetTable.entrySet()) { - putWriteBatch(writeBatch, entry.getKey(), entry.getValue()); - if (writeBatch.getDataSize() >= 4 * 1024) { - this.rocksDBConfigManager.batchPutWithWal(writeBatch); + if (!rocksDBConfigManager.isStop) { + try (WriteBatch writeBatch = new WriteBatch()) { + for (Entry> entry : this.offsetTable.entrySet()) { + putWriteBatch(writeBatch, entry.getKey(), entry.getValue()); + if (writeBatch.getDataSize() >= 4 * 1024) { + this.rocksDBConfigManager.batchPutWithWal(writeBatch); + } } + this.rocksDBConfigManager.batchPutWithWal(writeBatch); + this.rocksDBConfigManager.flushWAL(); + } catch (Exception e) { + log.error("consumer offset persist Failed", e); } - this.rocksDBConfigManager.batchPutWithWal(writeBatch); - this.rocksDBConfigManager.flushWAL(); - } catch (Exception e) { - log.error("consumer offset persist Failed", e); + } else { + log.warn("RocksDBConsumerOffsetManager has been stopped, persist fail"); } }