Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -53,6 +53,7 @@
public abstract class DiskUsageStatisticUtil implements Closeable {

protected static final Logger logger = LoggerFactory.getLogger(DiskUsageStatisticUtil.class);
// Collected files retain a reader reference and its resource read lock until released.
protected Queue<TsFileResource> resourcesWithReadLock;
protected final long timePartition;
protected final Iterator<TsFileResource> iterator;
Expand Down Expand Up @@ -100,9 +101,9 @@ protected void acquireReadLocks(List<TsFileResource> resources) {
if (!resource.isClosed()) {
continue;
}
resource.readLock();
FileReaderManager.getInstance().increaseFileReaderReference(resource, true);
if (resource.isDeleted() || !resource.isClosed()) {
resource.readUnlock();
FileReaderManager.getInstance().decreaseFileReaderReference(resource, true);
continue;
}
resourcesWithReadLock.add(resource);
Expand All @@ -118,7 +119,7 @@ protected void releaseReadLocks() {
return;
}
for (TsFileResource resource : resourcesWithReadLock) {
resource.readUnlock();
FileReaderManager.getInstance().decreaseFileReaderReference(resource, true);
}
resourcesWithReadLock = null;
}
Expand All @@ -127,10 +128,9 @@ public void calculateNextFile() {
TsFileResource tsFileResource = iterator.next();
if (tsFileResource.isDeleted() || calculateWithoutOpenFile(tsFileResource)) {
iterator.remove();
tsFileResource.readUnlock();
FileReaderManager.getInstance().decreaseFileReaderReference(tsFileResource, true);
return;
}
FileReaderManager.getInstance().increaseFileReaderReference(tsFileResource, true);
try {
TsFileSequenceReader reader =
FileReaderManager.getInstance()
Expand All @@ -146,7 +146,7 @@ public void calculateNextFile() {
tsFileResource.getTsFile().getAbsolutePath(),
e);
} finally {
// this operation including readUnlock
// Release the reader reference registered during collection, including its read lock.
FileReaderManager.getInstance().decreaseFileReaderReference(tsFileResource, true);
iterator.remove();
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -240,6 +240,14 @@ public void testCalculateTableSizeFromFile() throws Exception {

DataRegionTableSizeQueryContext context = new DataRegionTableSizeQueryContext(true);
queryTableSize(context);
boolean writeLockAcquired = resource1.tryWriteLock();
try {
Assert.assertTrue("Disk usage statistics should release all read locks", writeLockAcquired);
} finally {
if (writeLockAcquired) {
resource1.writeUnlock();
}
}
Assert.assertEquals(1, context.getTimePartitionTableSizeQueryContextMap().size());
TimePartitionTableSizeQueryContext timePartitionContext =
context.getTimePartitionTableSizeQueryContextMap().values().iterator().next();
Expand Down
Loading