This is an automated email from the ASF dual-hosted git repository.
FrankChen021 pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/druid.git
The following commit(s) were added to refs/heads/master by this push:
new a6907dea28c fix: address CodeQL resource lifetime warnings (#19816)
a6907dea28c is described below
commit a6907dea28c777e4f870a7688f2290cf01c08da0
Author: Frank Chen <[email protected]>
AuthorDate: Tue Aug 4 16:27:09 2026 +0800
fix: address CodeQL resource lifetime warnings (#19816)
* Fix CodeQL resource lifetime warnings
* style: make Avatica statement final
* test: ensure cleanup after Avatica failures
* test: tighten resource fixture types
---
.../storage/aliyun/OssDataSegmentPullerTest.java | 123 ++++++++-------
.../druid/storage/aliyun/OssTaskLogsTest.java | 36 +++--
.../movingaverage/MovingAverageQueryTest.java | 18 ++-
.../TimestampGroupByAggregationTest.java | 38 +++--
.../druid/query/filter/BloomKFilterTest.java | 77 ++++-----
.../druid/storage/s3/S3DataSegmentPullerTest.java | 9 +-
.../apache/druid/storage/s3/S3TaskLogsTest.java | 45 +++---
.../task/batch/parallel/HttpShuffleClientTest.java | 22 ++-
.../SQLMetadataStorageActionHandlerTest.java | 9 +-
.../java/util/common/io/smoosh/FileSmoosher.java | 4 +-
.../druid/segment/file/SegmentFileBuilderV10.java | 2 +
.../data/input/impl/RetryingInputStreamTest.java | 30 ++--
.../java/util/common/CompressionUtilsTest.java | 155 +++++++++++++-----
.../query/aggregation/AggregationTestHelper.java | 49 +++---
.../apache/druid/metadata/input/SqlEntityTest.java | 8 +-
.../org/apache/druid/rpc/RequestBuilderTest.java | 20 ++-
.../org/apache/druid/server/QueryResourceTest.java | 30 ++--
.../druid/server/initialization/JettyTest.java | 34 ++--
.../druid/sql/avatica/DruidAvaticaHandlerTest.java | 175 +++++++++++----------
19 files changed, 530 insertions(+), 354 deletions(-)
diff --git
a/extensions-contrib/aliyun-oss-extensions/src/test/java/org/apache/druid/storage/aliyun/OssDataSegmentPullerTest.java
b/extensions-contrib/aliyun-oss-extensions/src/test/java/org/apache/druid/storage/aliyun/OssDataSegmentPullerTest.java
index e3d667a411d..0fb2ec9c7c3 100644
---
a/extensions-contrib/aliyun-oss-extensions/src/test/java/org/apache/druid/storage/aliyun/OssDataSegmentPullerTest.java
+++
b/extensions-contrib/aliyun-oss-extensions/src/test/java/org/apache/druid/storage/aliyun/OssDataSegmentPullerTest.java
@@ -38,6 +38,7 @@ import java.io.File;
import java.io.FileInputStream;
import java.io.FileOutputStream;
import java.io.IOException;
+import java.io.InputStream;
import java.io.OutputStream;
import java.net.URI;
import java.nio.charset.StandardCharsets;
@@ -95,7 +96,8 @@ public class OssDataSegmentPullerTest
final File tmpFile = temporaryFolder.newFile("gzTest.gz");
- try (OutputStream outputStream = new GZIPOutputStream(new
FileOutputStream(tmpFile))) {
+ try (final FileOutputStream fileOutputStream = new
FileOutputStream(tmpFile);
+ final OutputStream outputStream = new
GZIPOutputStream(fileOutputStream)) {
outputStream.write(value);
}
@@ -103,7 +105,6 @@ public class OssDataSegmentPullerTest
object0.setBucketName(bucket);
object0.setKey(keyPrefix + "/renames-0.gz");
object0.getObjectMetadata().setLastModified(new Date(0));
- object0.setObjectContent(new FileInputStream(tmpFile));
final OSSObjectSummary objectSummary = new OSSObjectSummary();
objectSummary.setBucketName(bucket);
@@ -115,30 +116,33 @@ public class OssDataSegmentPullerTest
final File tmpDir = temporaryFolder.newFolder("gzTestDir");
-
EasyMock.expect(ossClient.doesObjectExist(EasyMock.eq(object0.getBucketName()),
EasyMock.eq(object0.getKey())))
- .andReturn(true)
- .once();
- EasyMock.expect(ossClient.getObjectMetadata(object0.getBucketName(),
object0.getKey()))
- .andReturn(objectMetadata)
- .once();
- EasyMock.expect(ossClient.getObject(EasyMock.eq(object0.getBucketName()),
EasyMock.eq(object0.getKey())))
- .andReturn(object0)
- .once();
- OssDataSegmentPuller puller = new OssDataSegmentPuller(ossClient);
-
- EasyMock.replay(ossClient);
- FileUtils.FileCopyResult result = puller.getSegmentFiles(
- new CloudObjectLocation(
- bucket,
- object0.getKey()
- ), tmpDir
- );
- EasyMock.verify(ossClient);
-
- Assert.assertEquals(value.length, result.size());
- File expected = new File(tmpDir, "renames-0");
- Assert.assertTrue(expected.exists());
- Assert.assertEquals(value.length, expected.length());
+ try (final InputStream objectContent = new FileInputStream(tmpFile)) {
+ object0.setObjectContent(objectContent);
+
EasyMock.expect(ossClient.doesObjectExist(EasyMock.eq(object0.getBucketName()),
EasyMock.eq(object0.getKey())))
+ .andReturn(true)
+ .once();
+ EasyMock.expect(ossClient.getObjectMetadata(object0.getBucketName(),
object0.getKey()))
+ .andReturn(objectMetadata)
+ .once();
+
EasyMock.expect(ossClient.getObject(EasyMock.eq(object0.getBucketName()),
EasyMock.eq(object0.getKey())))
+ .andReturn(object0)
+ .once();
+ final OssDataSegmentPuller puller = new OssDataSegmentPuller(ossClient);
+
+ EasyMock.replay(ossClient);
+ final FileUtils.FileCopyResult result = puller.getSegmentFiles(
+ new CloudObjectLocation(
+ bucket,
+ object0.getKey()
+ ), tmpDir
+ );
+ EasyMock.verify(ossClient);
+
+ Assert.assertEquals(value.length, result.size());
+ final File expected = new File(tmpDir, "renames-0");
+ Assert.assertTrue(expected.exists());
+ Assert.assertEquals(value.length, expected.length());
+ }
}
@Test
@@ -151,7 +155,8 @@ public class OssDataSegmentPullerTest
final File tmpFile = temporaryFolder.newFile("gzTest.gz");
- try (OutputStream outputStream = new GZIPOutputStream(new
FileOutputStream(tmpFile))) {
+ try (final FileOutputStream fileOutputStream = new
FileOutputStream(tmpFile);
+ final OutputStream outputStream = new
GZIPOutputStream(fileOutputStream)) {
outputStream.write(value);
}
@@ -160,7 +165,6 @@ public class OssDataSegmentPullerTest
object0.setBucketName(bucket);
object0.setKey(keyPrefix + "/renames-0.gz");
object0.getObjectMetadata().setLastModified(new Date(0));
- object0.setObjectContent(new FileInputStream(tmpFile));
final ObjectMetadata objectMetadata = new ObjectMetadata();
objectMetadata.setLastModified(new Date(0));
@@ -168,36 +172,39 @@ public class OssDataSegmentPullerTest
File tmpDir = temporaryFolder.newFolder("gzTestDir");
OSSException exception = new OSSException("OssDataSegmentPullerTest",
"NoSuchKey", null, null, null, null, null);
-
EasyMock.expect(ossClient.doesObjectExist(EasyMock.eq(object0.getBucketName()),
EasyMock.eq(object0.getKey())))
- .andReturn(true)
- .once();
- EasyMock.expect(ossClient.getObjectMetadata(bucket, object0.getKey()))
- .andReturn(objectMetadata)
- .once();
- EasyMock.expect(ossClient.getObject(EasyMock.eq(bucket),
EasyMock.eq(object0.getKey())))
- .andThrow(exception)
- .once();
- EasyMock.expect(ossClient.getObjectMetadata(bucket, object0.getKey()))
- .andReturn(objectMetadata)
- .once();
- EasyMock.expect(ossClient.getObject(EasyMock.eq(bucket),
EasyMock.eq(object0.getKey())))
- .andReturn(object0)
- .once();
- OssDataSegmentPuller puller = new OssDataSegmentPuller(ossClient);
-
- EasyMock.replay(ossClient);
- FileUtils.FileCopyResult result = puller.getSegmentFiles(
- new CloudObjectLocation(
- bucket,
- object0.getKey()
- ), tmpDir
- );
- EasyMock.verify(ossClient);
-
- Assert.assertEquals(value.length, result.size());
- File expected = new File(tmpDir, "renames-0");
- Assert.assertTrue(expected.exists());
- Assert.assertEquals(value.length, expected.length());
+ try (final InputStream objectContent = new FileInputStream(tmpFile)) {
+ object0.setObjectContent(objectContent);
+
EasyMock.expect(ossClient.doesObjectExist(EasyMock.eq(object0.getBucketName()),
EasyMock.eq(object0.getKey())))
+ .andReturn(true)
+ .once();
+ EasyMock.expect(ossClient.getObjectMetadata(bucket, object0.getKey()))
+ .andReturn(objectMetadata)
+ .once();
+ EasyMock.expect(ossClient.getObject(EasyMock.eq(bucket),
EasyMock.eq(object0.getKey())))
+ .andThrow(exception)
+ .once();
+ EasyMock.expect(ossClient.getObjectMetadata(bucket, object0.getKey()))
+ .andReturn(objectMetadata)
+ .once();
+ EasyMock.expect(ossClient.getObject(EasyMock.eq(bucket),
EasyMock.eq(object0.getKey())))
+ .andReturn(object0)
+ .once();
+ final OssDataSegmentPuller puller = new OssDataSegmentPuller(ossClient);
+
+ EasyMock.replay(ossClient);
+ final FileUtils.FileCopyResult result = puller.getSegmentFiles(
+ new CloudObjectLocation(
+ bucket,
+ object0.getKey()
+ ), tmpDir
+ );
+ EasyMock.verify(ossClient);
+
+ Assert.assertEquals(value.length, result.size());
+ final File expected = new File(tmpDir, "renames-0");
+ Assert.assertTrue(expected.exists());
+ Assert.assertEquals(value.length, expected.length());
+ }
}
}
diff --git
a/extensions-contrib/aliyun-oss-extensions/src/test/java/org/apache/druid/storage/aliyun/OssTaskLogsTest.java
b/extensions-contrib/aliyun-oss-extensions/src/test/java/org/apache/druid/storage/aliyun/OssTaskLogsTest.java
index 16b09866ec4..d88aa7af468 100644
---
a/extensions-contrib/aliyun-oss-extensions/src/test/java/org/apache/druid/storage/aliyun/OssTaskLogsTest.java
+++
b/extensions-contrib/aliyun-oss-extensions/src/test/java/org/apache/druid/storage/aliyun/OssTaskLogsTest.java
@@ -300,10 +300,11 @@ public class OssTaskLogsTest extends EasyMockSupport
OssTaskLogs ossTaskLogs = getOssTaskLogs();
Optional<InputStream> inputStreamOptional =
ossTaskLogs.streamTaskLog(KEY_1, 0);
- String taskLogs = new BufferedReader(
- new InputStreamReader(inputStreamOptional.get(),
StandardCharsets.UTF_8))
- .lines()
- .collect(Collectors.joining("\n"));
+ final String taskLogs;
+ try (final BufferedReader reader = new BufferedReader(
+ new InputStreamReader(inputStreamOptional.get(),
StandardCharsets.UTF_8))) {
+ taskLogs = reader.lines().collect(Collectors.joining("\n"));
+ }
Assert.assertEquals(LOG_CONTENTS, taskLogs);
}
@@ -324,10 +325,11 @@ public class OssTaskLogsTest extends EasyMockSupport
OssTaskLogs ossTaskLogs = getOssTaskLogs();
Optional<InputStream> inputStreamOptional =
ossTaskLogs.streamTaskLog(KEY_1, 1);
- String taskLogs = new BufferedReader(
- new InputStreamReader(inputStreamOptional.get(),
StandardCharsets.UTF_8))
- .lines()
- .collect(Collectors.joining("\n"));
+ final String taskLogs;
+ try (final BufferedReader reader = new BufferedReader(
+ new InputStreamReader(inputStreamOptional.get(),
StandardCharsets.UTF_8))) {
+ taskLogs = reader.lines().collect(Collectors.joining("\n"));
+ }
Assert.assertEquals(LOG_CONTENTS.substring(1), taskLogs);
}
@@ -348,10 +350,11 @@ public class OssTaskLogsTest extends EasyMockSupport
OssTaskLogs ossTaskLogs = getOssTaskLogs();
Optional<InputStream> inputStreamOptional =
ossTaskLogs.streamTaskLog(KEY_1, -1 * (LOG_CONTENTS.length() - 1));
- String taskLogs = new BufferedReader(
- new InputStreamReader(inputStreamOptional.get(),
StandardCharsets.UTF_8))
- .lines()
- .collect(Collectors.joining("\n"));
+ final String taskLogs;
+ try (final BufferedReader reader = new BufferedReader(
+ new InputStreamReader(inputStreamOptional.get(),
StandardCharsets.UTF_8))) {
+ taskLogs = reader.lines().collect(Collectors.joining("\n"));
+ }
Assert.assertEquals(LOG_CONTENTS.substring(1), taskLogs);
}
@@ -373,10 +376,11 @@ public class OssTaskLogsTest extends EasyMockSupport
OssTaskLogs ossTaskLogs = getOssTaskLogs();
Optional<InputStream> inputStreamOptional =
ossTaskLogs.streamTaskReports(KEY_1);
- String report = new BufferedReader(
- new InputStreamReader(inputStreamOptional.get(),
StandardCharsets.UTF_8))
- .lines()
- .collect(Collectors.joining("\n"));
+ final String report;
+ try (final BufferedReader reader = new BufferedReader(
+ new InputStreamReader(inputStreamOptional.get(),
StandardCharsets.UTF_8))) {
+ report = reader.lines().collect(Collectors.joining("\n"));
+ }
Assert.assertEquals(REPORT_CONTENTS, report);
}
diff --git
a/extensions-contrib/moving-average-query/src/test/java/org/apache/druid/query/movingaverage/MovingAverageQueryTest.java
b/extensions-contrib/moving-average-query/src/test/java/org/apache/druid/query/movingaverage/MovingAverageQueryTest.java
index 7072b974fac..a7ddabd6a4d 100644
---
a/extensions-contrib/moving-average-query/src/test/java/org/apache/druid/query/movingaverage/MovingAverageQueryTest.java
+++
b/extensions-contrib/moving-average-query/src/test/java/org/apache/druid/query/movingaverage/MovingAverageQueryTest.java
@@ -118,12 +118,17 @@ public class MovingAverageQueryTest extends
InitializedNullHandlingTest
@Parameters(name = "{0}")
public static Iterable<String[]> data() throws IOException
{
- BufferedReader testReader = new BufferedReader(
- new
InputStreamReader(MovingAverageQueryTest.class.getResourceAsStream("/queryTests"),
StandardCharsets.UTF_8));
List<String[]> tests = new ArrayList<>();
- for (String line = testReader.readLine(); line != null; line =
testReader.readLine()) {
- tests.add(new String[]{line});
+ try (final BufferedReader testReader = new BufferedReader(
+ new InputStreamReader(
+ MovingAverageQueryTest.class.getResourceAsStream("/queryTests"),
+ StandardCharsets.UTF_8
+ )
+ )) {
+ for (String line = testReader.readLine(); line != null; line =
testReader.readLine()) {
+ tests.add(new String[]{line});
+ }
}
return tests;
@@ -172,9 +177,10 @@ public class MovingAverageQueryTest extends
InitializedNullHandlingTest
retryConfig = injector.getInstance(RetryQueryRunnerConfig.class);
serverConfig = injector.getInstance(ServerConfig.class);
- InputStream is = getClass().getResourceAsStream("/queryTests/" + yamlFile);
ObjectMapper reader = new ObjectMapper(new YAMLFactory());
- config = reader.readValue(is, TestConfig.class);
+ try (final InputStream is = getClass().getResourceAsStream("/queryTests/"
+ yamlFile)) {
+ config = reader.readValue(is, TestConfig.class);
+ }
}
/**
diff --git
a/extensions-contrib/time-min-max/src/test/java/org/apache/druid/query/aggregation/TimestampGroupByAggregationTest.java
b/extensions-contrib/time-min-max/src/test/java/org/apache/druid/query/aggregation/TimestampGroupByAggregationTest.java
index 8aa47963fd7..45c8a791719 100644
---
a/extensions-contrib/time-min-max/src/test/java/org/apache/druid/query/aggregation/TimestampGroupByAggregationTest.java
+++
b/extensions-contrib/time-min-max/src/test/java/org/apache/druid/query/aggregation/TimestampGroupByAggregationTest.java
@@ -48,6 +48,7 @@ import org.junit.runners.Parameterized;
import java.io.File;
import java.io.IOException;
+import java.io.InputStream;
import java.sql.Timestamp;
import java.util.ArrayList;
import java.util.List;
@@ -149,23 +150,28 @@ public class TimestampGroupByAggregationTest
.setInterval("2011-01-01T00:00:00.000Z/2011-05-01T00:00:00.000Z")
.build();
- ZipFile zip = new ZipFile(new
File(this.getClass().getClassLoader().getResource("druid.sample.tsv.zip").toURI()));
- Sequence<ResultRow> seq = helper.createIndexAndRunQueryOnSegment(
- zip.getInputStream(zip.getEntry("druid.sample.tsv")),
- new InputRowSchema(
- new TimestampSpec("timestamp", "auto", null),
- new
DimensionsSpec(DimensionsSpec.getDefaultSchemas(List.of("product"))),
- ColumnsFilter.all()
- ),
- DelimitedInputFormat.forColumns(
- List.of("timestamp", "cat", "product", "prefer", "prefer2",
"pty_country")
- ),
- aggregators,
- 0,
- Granularities.MONTH,
- 100,
- groupByQuery
+ final Sequence<ResultRow> seq;
+ try (final ZipFile zip = new ZipFile(
+ new
File(this.getClass().getClassLoader().getResource("druid.sample.tsv.zip").toURI())
);
+ final InputStream inputStream =
zip.getInputStream(zip.getEntry("druid.sample.tsv"))) {
+ seq = helper.createIndexAndRunQueryOnSegment(
+ inputStream,
+ new InputRowSchema(
+ new TimestampSpec("timestamp", "auto", null),
+ new
DimensionsSpec(DimensionsSpec.getDefaultSchemas(List.of("product"))),
+ ColumnsFilter.all()
+ ),
+ DelimitedInputFormat.forColumns(
+ List.of("timestamp", "cat", "product", "prefer", "prefer2",
"pty_country")
+ ),
+ aggregators,
+ 0,
+ Granularities.MONTH,
+ 100,
+ groupByQuery
+ );
+ }
int groupByFieldNumber =
groupByQuery.getResultRowSignature().indexOf(groupByField);
diff --git
a/extensions-core/druid-bloom-filter/src/test/java/org/apache/druid/query/filter/BloomKFilterTest.java
b/extensions-core/druid-bloom-filter/src/test/java/org/apache/druid/query/filter/BloomKFilterTest.java
index 63bd658cfc7..bc87c2c2f5c 100644
---
a/extensions-core/druid-bloom-filter/src/test/java/org/apache/druid/query/filter/BloomKFilterTest.java
+++
b/extensions-core/druid-bloom-filter/src/test/java/org/apache/druid/query/filter/BloomKFilterTest.java
@@ -52,28 +52,28 @@ public class BloomKFilterTest
bf.add(val);
BloomKFilter.add(buffer, val);
- BloomKFilter rehydrated = BloomKFilter.deserialize(new
ByteBufferInputStream(buffer));
+ BloomKFilter rehydrated = deserializeBloomFilter(buffer);
buffer.position(0);
Assert.assertTrue(rehydrated.test(val));
Assert.assertFalse(rehydrated.test(val1));
Assert.assertFalse(rehydrated.test(val2));
Assert.assertFalse(rehydrated.test(val3));
BloomKFilter.add(buffer, val1);
- rehydrated = BloomKFilter.deserialize(new ByteBufferInputStream(buffer));
+ rehydrated = deserializeBloomFilter(buffer);
buffer.position(0);
Assert.assertTrue(rehydrated.test(val));
Assert.assertTrue(rehydrated.test(val1));
Assert.assertFalse(rehydrated.test(val2));
Assert.assertFalse(rehydrated.test(val3));
BloomKFilter.add(buffer, val2);
- rehydrated = BloomKFilter.deserialize(new ByteBufferInputStream(buffer));
+ rehydrated = deserializeBloomFilter(buffer);
buffer.position(0);
Assert.assertTrue(rehydrated.test(val));
Assert.assertTrue(rehydrated.test(val1));
Assert.assertTrue(rehydrated.test(val2));
Assert.assertFalse(rehydrated.test(val3));
BloomKFilter.add(buffer, val3);
- rehydrated = BloomKFilter.deserialize(new ByteBufferInputStream(buffer));
+ rehydrated = deserializeBloomFilter(buffer);
buffer.position(0);
Assert.assertTrue(rehydrated.test(val));
Assert.assertTrue(rehydrated.test(val1));
@@ -86,7 +86,7 @@ public class BloomKFilterTest
BloomKFilter.add(buffer, randVal);
}
// last value should be present
- rehydrated = BloomKFilter.deserialize(new ByteBufferInputStream(buffer));
+ rehydrated = deserializeBloomFilter(buffer);
Assert.assertTrue(rehydrated.test(randVal));
// most likely this value should not exist
randVal[0] = 0;
@@ -114,28 +114,28 @@ public class BloomKFilterTest
byte val3 = Byte.MAX_VALUE;
BloomKFilter.addLong(buffer, val);
- BloomKFilter rehydrated = BloomKFilter.deserialize(new
ByteBufferInputStream(buffer));
+ BloomKFilter rehydrated = deserializeBloomFilter(buffer);
buffer.position(0);
Assert.assertTrue(rehydrated.testLong(val));
Assert.assertFalse(rehydrated.testLong(val1));
Assert.assertFalse(rehydrated.testLong(val2));
Assert.assertFalse(rehydrated.testLong(val3));
BloomKFilter.addLong(buffer, val1);
- rehydrated = BloomKFilter.deserialize(new ByteBufferInputStream(buffer));
+ rehydrated = deserializeBloomFilter(buffer);
buffer.position(0);
Assert.assertTrue(rehydrated.testLong(val));
Assert.assertTrue(rehydrated.testLong(val1));
Assert.assertFalse(rehydrated.testLong(val2));
Assert.assertFalse(rehydrated.testLong(val3));
BloomKFilter.addLong(buffer, val2);
- rehydrated = BloomKFilter.deserialize(new ByteBufferInputStream(buffer));
+ rehydrated = deserializeBloomFilter(buffer);
buffer.position(0);
Assert.assertTrue(rehydrated.testLong(val));
Assert.assertTrue(rehydrated.testLong(val1));
Assert.assertTrue(rehydrated.testLong(val2));
Assert.assertFalse(rehydrated.testLong(val3));
BloomKFilter.addLong(buffer, val3);
- rehydrated = BloomKFilter.deserialize(new ByteBufferInputStream(buffer));
+ rehydrated = deserializeBloomFilter(buffer);
buffer.position(0);
Assert.assertTrue(rehydrated.testLong(val));
Assert.assertTrue(rehydrated.testLong(val1));
@@ -148,7 +148,7 @@ public class BloomKFilterTest
BloomKFilter.addLong(buffer, randVal);
}
- rehydrated = BloomKFilter.deserialize(new ByteBufferInputStream(buffer));
+ rehydrated = deserializeBloomFilter(buffer);
// last value should be present
Assert.assertTrue(rehydrated.testLong(randVal));
@@ -173,28 +173,28 @@ public class BloomKFilterTest
int val3 = Integer.MAX_VALUE;
BloomKFilter.addLong(buffer, val);
- BloomKFilter rehydrated = BloomKFilter.deserialize(new
ByteBufferInputStream(buffer));
+ BloomKFilter rehydrated = deserializeBloomFilter(buffer);
buffer.position(0);
Assert.assertTrue(rehydrated.testLong(val));
Assert.assertFalse(rehydrated.testLong(val1));
Assert.assertFalse(rehydrated.testLong(val2));
Assert.assertFalse(rehydrated.testLong(val3));
BloomKFilter.addLong(buffer, val1);
- rehydrated = BloomKFilter.deserialize(new ByteBufferInputStream(buffer));
+ rehydrated = deserializeBloomFilter(buffer);
buffer.position(0);
Assert.assertTrue(rehydrated.testLong(val));
Assert.assertTrue(rehydrated.testLong(val1));
Assert.assertFalse(rehydrated.testLong(val2));
Assert.assertFalse(rehydrated.testLong(val3));
BloomKFilter.addLong(buffer, val2);
- rehydrated = BloomKFilter.deserialize(new ByteBufferInputStream(buffer));
+ rehydrated = deserializeBloomFilter(buffer);
buffer.position(0);
Assert.assertTrue(rehydrated.testLong(val));
Assert.assertTrue(rehydrated.testLong(val1));
Assert.assertTrue(rehydrated.testLong(val2));
Assert.assertFalse(rehydrated.testLong(val3));
BloomKFilter.addLong(buffer, val3);
- rehydrated = BloomKFilter.deserialize(new ByteBufferInputStream(buffer));
+ rehydrated = deserializeBloomFilter(buffer);
buffer.position(0);
Assert.assertTrue(rehydrated.testLong(val));
Assert.assertTrue(rehydrated.testLong(val1));
@@ -206,7 +206,7 @@ public class BloomKFilterTest
randVal = rand.nextInt();
BloomKFilter.addLong(buffer, randVal);
}
- rehydrated = BloomKFilter.deserialize(new ByteBufferInputStream(buffer));
+ rehydrated = deserializeBloomFilter(buffer);
// last value should be present
Assert.assertTrue(rehydrated.testLong(randVal));
// most likely this value should not exist
@@ -230,28 +230,28 @@ public class BloomKFilterTest
long val3 = Long.MAX_VALUE;
BloomKFilter.addLong(buffer, val);
- BloomKFilter rehydrated = BloomKFilter.deserialize(new
ByteBufferInputStream(buffer));
+ BloomKFilter rehydrated = deserializeBloomFilter(buffer);
buffer.position(0);
Assert.assertTrue(rehydrated.testLong(val));
Assert.assertFalse(rehydrated.testLong(val1));
Assert.assertFalse(rehydrated.testLong(val2));
Assert.assertFalse(rehydrated.testLong(val3));
BloomKFilter.addLong(buffer, val1);
- rehydrated = BloomKFilter.deserialize(new ByteBufferInputStream(buffer));
+ rehydrated = deserializeBloomFilter(buffer);
buffer.position(0);
Assert.assertTrue(rehydrated.testLong(val));
Assert.assertTrue(rehydrated.testLong(val1));
Assert.assertFalse(rehydrated.testLong(val2));
Assert.assertFalse(rehydrated.testLong(val3));
BloomKFilter.addLong(buffer, val2);
- rehydrated = BloomKFilter.deserialize(new ByteBufferInputStream(buffer));
+ rehydrated = deserializeBloomFilter(buffer);
buffer.position(0);
Assert.assertTrue(rehydrated.testLong(val));
Assert.assertTrue(rehydrated.testLong(val1));
Assert.assertTrue(rehydrated.testLong(val2));
Assert.assertFalse(rehydrated.testLong(val3));
BloomKFilter.addLong(buffer, val3);
- rehydrated = BloomKFilter.deserialize(new ByteBufferInputStream(buffer));
+ rehydrated = deserializeBloomFilter(buffer);
buffer.position(0);
Assert.assertTrue(rehydrated.testLong(val));
Assert.assertTrue(rehydrated.testLong(val1));
@@ -263,7 +263,7 @@ public class BloomKFilterTest
randVal = rand.nextInt();
BloomKFilter.addLong(buffer, randVal);
}
- rehydrated = BloomKFilter.deserialize(new ByteBufferInputStream(buffer));
+ rehydrated = deserializeBloomFilter(buffer);
// last value should be present
Assert.assertTrue(rehydrated.testLong(randVal));
// most likely this value should not exist
@@ -287,28 +287,28 @@ public class BloomKFilterTest
float val3 = Float.POSITIVE_INFINITY;
BloomKFilter.addFloat(buffer, val);
- BloomKFilter rehydrated = BloomKFilter.deserialize(new
ByteBufferInputStream(buffer));
+ BloomKFilter rehydrated = deserializeBloomFilter(buffer);
buffer.position(0);
Assert.assertTrue(rehydrated.testFloat(val));
Assert.assertFalse(rehydrated.testFloat(val1));
Assert.assertFalse(rehydrated.testFloat(val2));
Assert.assertFalse(rehydrated.testFloat(val3));
BloomKFilter.addFloat(buffer, val1);
- rehydrated = BloomKFilter.deserialize(new ByteBufferInputStream(buffer));
+ rehydrated = deserializeBloomFilter(buffer);
buffer.position(0);
Assert.assertTrue(rehydrated.testFloat(val));
Assert.assertTrue(rehydrated.testFloat(val1));
Assert.assertFalse(rehydrated.testFloat(val2));
Assert.assertFalse(rehydrated.testFloat(val3));
BloomKFilter.addFloat(buffer, val2);
- rehydrated = BloomKFilter.deserialize(new ByteBufferInputStream(buffer));
+ rehydrated = deserializeBloomFilter(buffer);
buffer.position(0);
Assert.assertTrue(rehydrated.testFloat(val));
Assert.assertTrue(rehydrated.testFloat(val1));
Assert.assertTrue(rehydrated.testFloat(val2));
Assert.assertFalse(rehydrated.testFloat(val3));
BloomKFilter.addFloat(buffer, val3);
- rehydrated = BloomKFilter.deserialize(new ByteBufferInputStream(buffer));
+ rehydrated = deserializeBloomFilter(buffer);
buffer.position(0);
Assert.assertTrue(rehydrated.testFloat(val));
Assert.assertTrue(rehydrated.testFloat(val1));
@@ -320,7 +320,7 @@ public class BloomKFilterTest
randVal = rand.nextFloat();
BloomKFilter.addFloat(buffer, randVal);
}
- rehydrated = BloomKFilter.deserialize(new ByteBufferInputStream(buffer));
+ rehydrated = deserializeBloomFilter(buffer);
// last value should be present
Assert.assertTrue(rehydrated.testFloat(randVal));
@@ -345,28 +345,28 @@ public class BloomKFilterTest
double val3 = Double.POSITIVE_INFINITY;
BloomKFilter.addDouble(buffer, val);
- BloomKFilter rehydrated = BloomKFilter.deserialize(new
ByteBufferInputStream(buffer));
+ BloomKFilter rehydrated = deserializeBloomFilter(buffer);
buffer.position(0);
Assert.assertTrue(rehydrated.testDouble(val));
Assert.assertFalse(rehydrated.testDouble(val1));
Assert.assertFalse(rehydrated.testDouble(val2));
Assert.assertFalse(rehydrated.testDouble(val3));
BloomKFilter.addDouble(buffer, val1);
- rehydrated = BloomKFilter.deserialize(new ByteBufferInputStream(buffer));
+ rehydrated = deserializeBloomFilter(buffer);
buffer.position(0);
Assert.assertTrue(rehydrated.testDouble(val));
Assert.assertTrue(rehydrated.testDouble(val1));
Assert.assertFalse(rehydrated.testDouble(val2));
Assert.assertFalse(rehydrated.testDouble(val3));
BloomKFilter.addDouble(buffer, val2);
- rehydrated = BloomKFilter.deserialize(new ByteBufferInputStream(buffer));
+ rehydrated = deserializeBloomFilter(buffer);
buffer.position(0);
Assert.assertTrue(rehydrated.testDouble(val));
Assert.assertTrue(rehydrated.testDouble(val1));
Assert.assertTrue(rehydrated.testDouble(val2));
Assert.assertFalse(rehydrated.testDouble(val3));
BloomKFilter.addDouble(buffer, val3);
- rehydrated = BloomKFilter.deserialize(new ByteBufferInputStream(buffer));
+ rehydrated = deserializeBloomFilter(buffer);
buffer.position(0);
Assert.assertTrue(rehydrated.testDouble(val));
Assert.assertTrue(rehydrated.testDouble(val1));
@@ -378,7 +378,7 @@ public class BloomKFilterTest
randVal = rand.nextDouble();
BloomKFilter.addDouble(buffer, randVal);
}
- rehydrated = BloomKFilter.deserialize(new ByteBufferInputStream(buffer));
+ rehydrated = deserializeBloomFilter(buffer);
// last value should be present
Assert.assertTrue(rehydrated.testDouble(randVal));
@@ -403,28 +403,28 @@ public class BloomKFilterTest
String val3 = "cuckoo filter";
BloomKFilter.addString(buffer, val);
- BloomKFilter rehydrated = BloomKFilter.deserialize(new
ByteBufferInputStream(buffer));
+ BloomKFilter rehydrated = deserializeBloomFilter(buffer);
buffer.position(0);
Assert.assertTrue(rehydrated.testString(val));
Assert.assertFalse(rehydrated.testString(val1));
Assert.assertFalse(rehydrated.testString(val2));
Assert.assertFalse(rehydrated.testString(val3));
BloomKFilter.addString(buffer, val1);
- rehydrated = BloomKFilter.deserialize(new ByteBufferInputStream(buffer));
+ rehydrated = deserializeBloomFilter(buffer);
buffer.position(0);
Assert.assertTrue(rehydrated.testString(val));
Assert.assertTrue(rehydrated.testString(val1));
Assert.assertFalse(rehydrated.testString(val2));
Assert.assertFalse(rehydrated.testString(val3));
BloomKFilter.addString(buffer, val2);
- rehydrated = BloomKFilter.deserialize(new ByteBufferInputStream(buffer));
+ rehydrated = deserializeBloomFilter(buffer);
buffer.position(0);
Assert.assertTrue(rehydrated.testString(val));
Assert.assertTrue(rehydrated.testString(val1));
Assert.assertTrue(rehydrated.testString(val2));
Assert.assertFalse(rehydrated.testString(val3));
BloomKFilter.addString(buffer, val3);
- rehydrated = BloomKFilter.deserialize(new ByteBufferInputStream(buffer));
+ rehydrated = deserializeBloomFilter(buffer);
buffer.position(0);
Assert.assertTrue(rehydrated.testString(val));
Assert.assertTrue(rehydrated.testString(val1));
@@ -436,7 +436,7 @@ public class BloomKFilterTest
randVal = rand.nextLong();
BloomKFilter.addString(buffer, Long.toString(randVal));
}
- rehydrated = BloomKFilter.deserialize(new ByteBufferInputStream(buffer));
+ rehydrated = deserializeBloomFilter(buffer);
// last value should be present
Assert.assertTrue(rehydrated.testString(Long.toString(randVal)));
// most likely this value should not exist
@@ -536,4 +536,11 @@ public class BloomKFilterTest
BloomKFilter.getNumSetBits(bufWithValues, 0) >
BloomKFilter.getNumSetBits(bufWithNull, 0)
);
}
+
+ private static BloomKFilter deserializeBloomFilter(final ByteBuffer buffer)
throws IOException
+ {
+ try (final ByteBufferInputStream inputStream = new
ByteBufferInputStream(buffer)) {
+ return BloomKFilter.deserialize(inputStream);
+ }
+ }
}
diff --git
a/extensions-core/s3-extensions/src/test/java/org/apache/druid/storage/s3/S3DataSegmentPullerTest.java
b/extensions-core/s3-extensions/src/test/java/org/apache/druid/storage/s3/S3DataSegmentPullerTest.java
index 289c60219af..75c241f3c11 100644
---
a/extensions-core/s3-extensions/src/test/java/org/apache/druid/storage/s3/S3DataSegmentPullerTest.java
+++
b/extensions-core/s3-extensions/src/test/java/org/apache/druid/storage/s3/S3DataSegmentPullerTest.java
@@ -95,7 +95,8 @@ public class S3DataSegmentPullerTest
final File tmpFile = temporaryFolder.newFile("gzTest.gz");
- try (OutputStream outputStream = new GZIPOutputStream(new
FileOutputStream(tmpFile))) {
+ try (final FileOutputStream fileOutputStream = new
FileOutputStream(tmpFile);
+ final OutputStream outputStream = new
GZIPOutputStream(fileOutputStream)) {
outputStream.write(value);
}
@@ -135,7 +136,8 @@ public class S3DataSegmentPullerTest
final File tmpFile = temporaryFolder.newFile("gzTest.gz");
- try (OutputStream outputStream = new GZIPOutputStream(new
FileOutputStream(tmpFile))) {
+ try (final FileOutputStream fileOutputStream = new
FileOutputStream(tmpFile);
+ final OutputStream outputStream = new
GZIPOutputStream(fileOutputStream)) {
outputStream.write(value);
}
@@ -181,7 +183,8 @@ public class S3DataSegmentPullerTest
final File tmpFile = temporaryFolder.newFile("gzTest.gz");
- try (OutputStream outputStream = new GZIPOutputStream(new
FileOutputStream(tmpFile))) {
+ try (final FileOutputStream fileOutputStream = new
FileOutputStream(tmpFile);
+ final OutputStream outputStream = new
GZIPOutputStream(fileOutputStream)) {
outputStream.write(value);
}
diff --git
a/extensions-core/s3-extensions/src/test/java/org/apache/druid/storage/s3/S3TaskLogsTest.java
b/extensions-core/s3-extensions/src/test/java/org/apache/druid/storage/s3/S3TaskLogsTest.java
index 2516b62ba4c..3da6bf53ce2 100644
---
a/extensions-core/s3-extensions/src/test/java/org/apache/druid/storage/s3/S3TaskLogsTest.java
+++
b/extensions-core/s3-extensions/src/test/java/org/apache/druid/storage/s3/S3TaskLogsTest.java
@@ -454,10 +454,11 @@ public class S3TaskLogsTest extends EasyMockSupport
S3TaskLogs s3TaskLogs = getS3TaskLogs();
Optional<InputStream> inputStreamOptional =
s3TaskLogs.streamTaskLog(KEY_1, 0);
- String taskLogs = new BufferedReader(
- new InputStreamReader(inputStreamOptional.get(),
StandardCharsets.UTF_8))
- .lines()
- .collect(Collectors.joining("\n"));
+ final String taskLogs;
+ try (final BufferedReader reader = new BufferedReader(
+ new InputStreamReader(inputStreamOptional.get(),
StandardCharsets.UTF_8))) {
+ taskLogs = reader.lines().collect(Collectors.joining("\n"));
+ }
Assert.assertEquals(LOG_CONTENTS, taskLogs);
}
@@ -483,10 +484,11 @@ public class S3TaskLogsTest extends EasyMockSupport
S3TaskLogs s3TaskLogs = getS3TaskLogs();
Optional<InputStream> inputStreamOptional =
s3TaskLogs.streamTaskLog(KEY_1, 1);
- String taskLogs = new BufferedReader(
- new InputStreamReader(inputStreamOptional.get(),
StandardCharsets.UTF_8))
- .lines()
- .collect(Collectors.joining("\n"));
+ final String taskLogs;
+ try (final BufferedReader reader = new BufferedReader(
+ new InputStreamReader(inputStreamOptional.get(),
StandardCharsets.UTF_8))) {
+ taskLogs = reader.lines().collect(Collectors.joining("\n"));
+ }
Assert.assertEquals(LOG_CONTENTS.substring(1), taskLogs);
}
@@ -512,10 +514,11 @@ public class S3TaskLogsTest extends EasyMockSupport
S3TaskLogs s3TaskLogs = getS3TaskLogs();
Optional<InputStream> inputStreamOptional =
s3TaskLogs.streamTaskLog(KEY_1, -1 * (LOG_CONTENTS.length() - 1));
- String taskLogs = new BufferedReader(
- new InputStreamReader(inputStreamOptional.get(),
StandardCharsets.UTF_8))
- .lines()
- .collect(Collectors.joining("\n"));
+ final String taskLogs;
+ try (final BufferedReader reader = new BufferedReader(
+ new InputStreamReader(inputStreamOptional.get(),
StandardCharsets.UTF_8))) {
+ taskLogs = reader.lines().collect(Collectors.joining("\n"));
+ }
Assert.assertEquals(LOG_CONTENTS.substring(1), taskLogs);
}
@@ -541,10 +544,11 @@ public class S3TaskLogsTest extends EasyMockSupport
S3TaskLogs s3TaskLogs = getS3TaskLogs();
Optional<InputStream> inputStreamOptional =
s3TaskLogs.streamTaskReports(KEY_1);
- String report = new BufferedReader(
- new InputStreamReader(inputStreamOptional.get(),
StandardCharsets.UTF_8))
- .lines()
- .collect(Collectors.joining("\n"));
+ final String report;
+ try (final BufferedReader reader = new BufferedReader(
+ new InputStreamReader(inputStreamOptional.get(),
StandardCharsets.UTF_8))) {
+ report = reader.lines().collect(Collectors.joining("\n"));
+ }
Assert.assertEquals(REPORT_CONTENTS, report);
}
@@ -569,10 +573,11 @@ public class S3TaskLogsTest extends EasyMockSupport
S3TaskLogs s3TaskLogs = getS3TaskLogs();
Optional<InputStream> inputStreamOptional =
s3TaskLogs.streamTaskStatus(KEY_1);
- String report = new BufferedReader(
- new InputStreamReader(inputStreamOptional.get(),
StandardCharsets.UTF_8))
- .lines()
- .collect(Collectors.joining("\n"));
+ final String report;
+ try (final BufferedReader reader = new BufferedReader(
+ new InputStreamReader(inputStreamOptional.get(),
StandardCharsets.UTF_8))) {
+ report = reader.lines().collect(Collectors.joining("\n"));
+ }
Assert.assertEquals(STATUS_CONTENTS, report);
}
diff --git
a/indexing-service/src/test/java/org/apache/druid/indexing/common/task/batch/parallel/HttpShuffleClientTest.java
b/indexing-service/src/test/java/org/apache/druid/indexing/common/task/batch/parallel/HttpShuffleClientTest.java
index 34e2b33b30c..b0e1ca12aaa 100644
---
a/indexing-service/src/test/java/org/apache/druid/indexing/common/task/batch/parallel/HttpShuffleClientTest.java
+++
b/indexing-service/src/test/java/org/apache/druid/indexing/common/task/batch/parallel/HttpShuffleClientTest.java
@@ -20,10 +20,12 @@
package org.apache.druid.indexing.common.task.batch.parallel;
import com.google.common.util.concurrent.Futures;
+import com.google.common.util.concurrent.ListenableFuture;
import org.apache.druid.java.util.common.Intervals;
import org.apache.druid.java.util.common.StringUtils;
import org.apache.druid.java.util.common.concurrent.Execs;
import org.apache.druid.java.util.http.client.HttpClient;
+import
org.apache.druid.java.util.http.client.response.InputStreamResponseHandler;
import org.apache.druid.utils.CompressionUtils;
import org.easymock.EasyMock;
import org.joda.time.Interval;
@@ -38,6 +40,7 @@ import java.io.File;
import java.io.FileInputStream;
import java.io.FileNotFoundException;
import java.io.IOException;
+import java.io.InputStream;
import java.io.Writer;
import java.nio.charset.StandardCharsets;
import java.nio.file.Files;
@@ -183,21 +186,28 @@ public class HttpShuffleClientTest
{
HttpClient httpClient = EasyMock.strictMock(HttpClient.class);
if (numFailures == 0) {
- EasyMock.expect(httpClient.go(EasyMock.anyObject(),
EasyMock.anyObject()))
+ EasyMock.expect(httpClient.go(EasyMock.anyObject(),
EasyMock.<InputStreamResponseHandler>anyObject()))
// should return different instances of input stream
- .andReturn(Futures.immediateFuture(new
FileInputStream(segmentFile)))
- .andReturn(Futures.immediateFuture(new
FileInputStream(segmentFile)));
+ .andReturn(openSegmentFile())
+ .andReturn(openSegmentFile());
} else {
- EasyMock.expect(httpClient.go(EasyMock.anyObject(),
EasyMock.anyObject()))
+ EasyMock.expect(httpClient.go(EasyMock.anyObject(),
EasyMock.<InputStreamResponseHandler>anyObject()))
.andReturn(Futures.immediateFailedFuture(new
RuntimeException())).times(numFailures)
// should return different instances of input stream
- .andReturn(Futures.immediateFuture(new
FileInputStream(segmentFile)))
- .andReturn(Futures.immediateFuture(new
FileInputStream(segmentFile)));
+ .andReturn(openSegmentFile())
+ .andReturn(openSegmentFile());
}
EasyMock.replay(httpClient);
return new HttpShuffleClient(httpClient);
}
+ private ListenableFuture<InputStream> openSegmentFile() throws
FileNotFoundException
+ {
+ // Ownership passes through HttpShuffleClient to FileUtils.copyLarge,
which closes the response stream.
+ // codeql[java/input-resource-leak]
+ return Futures.immediateFuture(new FileInputStream(segmentFile));
+ }
+
private static class TestPartitionLocation extends GenericPartitionLocation
{
private TestPartitionLocation()
diff --git
a/indexing-service/src/test/java/org/apache/druid/metadata/SQLMetadataStorageActionHandlerTest.java
b/indexing-service/src/test/java/org/apache/druid/metadata/SQLMetadataStorageActionHandlerTest.java
index 604f682fb66..84f958b56e6 100644
---
a/indexing-service/src/test/java/org/apache/druid/metadata/SQLMetadataStorageActionHandlerTest.java
+++
b/indexing-service/src/test/java/org/apache/druid/metadata/SQLMetadataStorageActionHandlerTest.java
@@ -51,6 +51,7 @@ import org.junit.Rule;
import org.junit.Test;
import java.sql.ResultSet;
+import java.sql.Statement;
import java.util.HashMap;
import java.util.HashSet;
import java.util.List;
@@ -454,9 +455,11 @@ public class SQLMetadataStorageActionHandlerTest
"SELECT COUNT(*) FROM %s WHERE type is NULL or group_id is NULL",
entryTable
);
- ResultSet resultSet =
handle.getConnection().createStatement().executeQuery(sql);
- resultSet.next();
- return resultSet.getInt(1);
+ try (final Statement statement =
handle.getConnection().createStatement();
+ final ResultSet resultSet = statement.executeQuery(sql)) {
+ resultSet.next();
+ return resultSet.getInt(1);
+ }
}
);
}
diff --git
a/processing/src/main/java/org/apache/druid/java/util/common/io/smoosh/FileSmoosher.java
b/processing/src/main/java/org/apache/druid/java/util/common/io/smoosh/FileSmoosher.java
index f6166e65648..4b84e4cd7c5 100644
---
a/processing/src/main/java/org/apache/druid/java/util/common/io/smoosh/FileSmoosher.java
+++
b/processing/src/main/java/org/apache/druid/java/util/common/io/smoosh/FileSmoosher.java
@@ -447,7 +447,9 @@ public class FileSmoosher implements SegmentFileBuilder
this.outFile = outFile;
this.maxLength = maxLength;
- FileOutputStream outStream = closer.register(new
FileOutputStream(outFile)); // lgtm [java/output-resource-leak]
+ // The closer owns this stream for the lifetime of the writer and
releases it in close().
+ // codeql[java/output-resource-leak]
+ final FileOutputStream outStream = closer.register(new
FileOutputStream(outFile));
this.channel = closer.register(outStream.getChannel());
}
diff --git
a/processing/src/main/java/org/apache/druid/segment/file/SegmentFileBuilderV10.java
b/processing/src/main/java/org/apache/druid/segment/file/SegmentFileBuilderV10.java
index d5281ab1186..230b705b4bb 100644
---
a/processing/src/main/java/org/apache/druid/segment/file/SegmentFileBuilderV10.java
+++
b/processing/src/main/java/org/apache/druid/segment/file/SegmentFileBuilderV10.java
@@ -685,6 +685,8 @@ public class SegmentFileBuilderV10 implements
SegmentFileBuilder
this.file = file;
this.bundle = bundle;
this.maxSize = maxSize;
+ // The closer owns this stream for the lifetime of the writer and
releases it in close().
+ // codeql[java/output-resource-leak]
final FileOutputStream outStream = closer.register(new
FileOutputStream(file));
this.channel = closer.register(outStream.getChannel());
}
diff --git
a/processing/src/test/java/org/apache/druid/data/input/impl/RetryingInputStreamTest.java
b/processing/src/test/java/org/apache/druid/data/input/impl/RetryingInputStreamTest.java
index ac5cff961bb..a4855b91d58 100644
---
a/processing/src/test/java/org/apache/druid/data/input/impl/RetryingInputStreamTest.java
+++
b/processing/src/test/java/org/apache/druid/data/input/impl/RetryingInputStreamTest.java
@@ -128,15 +128,15 @@ public class RetryingInputStreamTest
public void testRetryOnCustomException() throws IOException
{
throwCustomExceptions = 1;
- final RetryingInputStream<File> retryingInputStream = new
RetryingInputStream<>(
+ try (final RetryingInputStream<File> retryingInputStream = new
RetryingInputStream<>(
testFile,
objectOpenFunction,
t -> t instanceof CustomException,
MAX_RETRY,
false
- );
-
- retryHelper(retryingInputStream);
+ )) {
+ retryHelper(retryingInputStream);
+ }
Assertions.assertEquals(0, throwCustomExceptions);
}
@@ -168,15 +168,15 @@ public class RetryingInputStreamTest
readBytesBeforeExceptions = 1000;
throwCustomExceptions = 100;
- final RetryingInputStream<File> retryingInputStream = new
RetryingInputStream<>(
+ try (final RetryingInputStream<File> retryingInputStream = new
RetryingInputStream<>(
testFile,
objectOpenFunction,
t -> true, // always retry
MAX_RETRY,
false
- );
-
- retryHelper(retryingInputStream);
+ )) {
+ retryHelper(retryingInputStream);
+ }
// Tried more than MAX_RETRY times because progress was being made.
(MAX_RETRIES applies to each call individually.)
Assertions.assertEquals(81, throwCustomExceptions);
@@ -207,15 +207,15 @@ public class RetryingInputStreamTest
{
throwCustomExceptions = 1;
throwIOExceptions = 1;
- final RetryingInputStream<File> retryingInputStream = new
RetryingInputStream<>(
+ try (final RetryingInputStream<File> retryingInputStream = new
RetryingInputStream<>(
testFile,
objectOpenFunction,
t -> t instanceof IOException || t instanceof CustomException,
MAX_RETRY,
false
- );
-
- retryHelper(retryingInputStream);
+ )) {
+ retryHelper(retryingInputStream);
+ }
Assertions.assertEquals(0, throwCustomExceptions);
Assertions.assertEquals(0, throwIOExceptions);
@@ -242,13 +242,15 @@ public class RetryingInputStreamTest
}
}).when(objectOpenFunction).open(any(), anyLong());
- new RetryingInputStream<>(
+ try (final RetryingInputStream<File> ignored = new RetryingInputStream<>(
testFile,
objectOpenFunction,
t -> t instanceof CustomException,
MAX_RETRY,
false
- );
+ )) {
+ // Construction itself exercises the retry behavior.
+ }
verify(objectOpenFunction, times(3)).open(any(), anyLong());
Assertions.assertEquals(0, throwCustomExceptions);
}
diff --git
a/processing/src/test/java/org/apache/druid/java/util/common/CompressionUtilsTest.java
b/processing/src/test/java/org/apache/druid/java/util/common/CompressionUtilsTest.java
index 937c51458d9..50a65c56832 100644
---
a/processing/src/test/java/org/apache/druid/java/util/common/CompressionUtilsTest.java
+++
b/processing/src/test/java/org/apache/druid/java/util/common/CompressionUtilsTest.java
@@ -57,6 +57,7 @@ import java.util.HashMap;
import java.util.Iterator;
import java.util.List;
import java.util.Map;
+import java.util.Objects;
import java.util.Scanner;
import java.util.concurrent.atomic.AtomicInteger;
import java.util.concurrent.atomic.AtomicLong;
@@ -75,10 +76,15 @@ public class CompressionUtilsTest
static {
final StringBuilder builder = new StringBuilder();
- try (InputStream stream =
CompressionUtilsTest.class.getClassLoader().getResourceAsStream("white-rabbit.txt"))
{
- final Iterator<String> it = new Scanner(
- new InputStreamReader(stream, StandardCharsets.UTF_8)
- ).useDelimiter(Pattern.quote(System.lineSeparator()));
+ try (
+ final InputStream stream = Objects.requireNonNull(
+
CompressionUtilsTest.class.getClassLoader().getResourceAsStream("white-rabbit.txt"),
+ "Missing test resource: white-rabbit.txt"
+ );
+ final InputStreamReader reader = new InputStreamReader(stream,
StandardCharsets.UTF_8);
+ final Scanner scanner = new
Scanner(reader).useDelimiter(Pattern.quote(System.lineSeparator()))
+ ) {
+ final Iterator<String> it = scanner;
while (it.hasNext()) {
builder.append(it.next());
}
@@ -193,7 +199,10 @@ public class CompressionUtilsTest
Assert.assertFalse(gzFile.exists());
CompressionUtils.gzip(testFile, gzFile);
Assert.assertTrue(gzFile.exists());
- try (final InputStream inputStream = new GZIPInputStream(new
FileInputStream(gzFile))) {
+ try (
+ final InputStream fileInputStream = new FileInputStream(gzFile);
+ final InputStream inputStream = new GZIPInputStream(fileInputStream)
+ ) {
assertGoodDataStream(inputStream);
}
testFile.delete();
@@ -210,10 +219,15 @@ public class CompressionUtilsTest
{
final File tmpDir = temporaryFolder.newFolder("testGoodZipStream");
final File zipFile = new File(tmpDir, "compressionUtilTest.zip");
- CompressionUtils.zip(testDir, new FileOutputStream(zipFile));
+ try (final OutputStream outputStream = new FileOutputStream(zipFile)) {
+ CompressionUtils.zip(testDir, outputStream);
+ }
final File newDir = new File(tmpDir, "newDir");
newDir.mkdir();
- final FileUtils.FileCopyResult result = CompressionUtils.unzip(new
FileInputStream(zipFile), newDir);
+ final FileUtils.FileCopyResult result;
+ try (final InputStream inputStream = new FileInputStream(zipFile)) {
+ result = CompressionUtils.unzip(inputStream, newDir);
+ }
verifyUnzip(newDir, result, ImmutableMap.of(testFile.getName(),
StringUtils.toUtf8(CONTENT)));
}
@@ -232,7 +246,9 @@ public class CompressionUtilsTest
}
}
- CompressionUtils.zip(srcDir, new FileOutputStream(zipFile));
+ try (final OutputStream outputStream = new FileOutputStream(zipFile)) {
+ CompressionUtils.zip(srcDir, outputStream);
+ }
return expectedFiles;
}
@@ -308,7 +324,10 @@ public class CompressionUtilsTest
Assert.assertFalse(gzFile.exists());
CompressionUtils.gzip(Files.asByteSource(testFile),
Files.asByteSink(gzFile), Predicates.alwaysTrue());
Assert.assertTrue(gzFile.exists());
- try (final InputStream inputStream = CompressionUtils.decompress(new
FileInputStream(gzFile), gzFile.getName())) {
+ try (
+ final InputStream fileInputStream = new FileInputStream(gzFile);
+ final InputStream inputStream =
CompressionUtils.decompress(fileInputStream, gzFile.getName())
+ ) {
assertGoodDataStream(inputStream);
}
if (!testFile.delete()) {
@@ -328,10 +347,17 @@ public class CompressionUtilsTest
final File tmpDir = temporaryFolder.newFolder("testDecompressBzip2");
final File bzFile = new File(tmpDir, testFile.getName() + ".bz2");
Assert.assertFalse(bzFile.exists());
- try (final OutputStream out = new BZip2CompressorOutputStream(new
FileOutputStream(bzFile))) {
- ByteStreams.copy(new FileInputStream(testFile), out);
+ try (
+ final OutputStream fileOutputStream = new FileOutputStream(bzFile);
+ final OutputStream out = new
BZip2CompressorOutputStream(fileOutputStream);
+ final InputStream in = new FileInputStream(testFile)
+ ) {
+ ByteStreams.copy(in, out);
}
- try (final InputStream inputStream = CompressionUtils.decompress(new
FileInputStream(bzFile), bzFile.getName())) {
+ try (
+ final InputStream fileInputStream = new FileInputStream(bzFile);
+ final InputStream inputStream =
CompressionUtils.decompress(fileInputStream, bzFile.getName())
+ ) {
assertGoodDataStream(inputStream);
}
}
@@ -342,10 +368,17 @@ public class CompressionUtilsTest
final File tmpDir = temporaryFolder.newFolder("testDecompressXz");
final File xzFile = new File(tmpDir, testFile.getName() + ".xz");
Assert.assertFalse(xzFile.exists());
- try (final OutputStream out = new XZCompressorOutputStream(new
FileOutputStream(xzFile))) {
- ByteStreams.copy(new FileInputStream(testFile), out);
+ try (
+ final OutputStream fileOutputStream = new FileOutputStream(xzFile);
+ final OutputStream out = new
XZCompressorOutputStream(fileOutputStream);
+ final InputStream in = new FileInputStream(testFile)
+ ) {
+ ByteStreams.copy(in, out);
}
- try (final InputStream inputStream = CompressionUtils.decompress(new
FileInputStream(xzFile), xzFile.getName())) {
+ try (
+ final InputStream fileInputStream = new FileInputStream(xzFile);
+ final InputStream inputStream =
CompressionUtils.decompress(fileInputStream, xzFile.getName())
+ ) {
assertGoodDataStream(inputStream);
}
}
@@ -356,10 +389,17 @@ public class CompressionUtilsTest
final File tmpDir = temporaryFolder.newFolder("testDecompressSnappy");
final File snappyFile = new File(tmpDir, testFile.getName() + ".sz");
Assert.assertFalse(snappyFile.exists());
- try (final OutputStream out = new FramedSnappyCompressorOutputStream(new
FileOutputStream(snappyFile))) {
- ByteStreams.copy(new FileInputStream(testFile), out);
+ try (
+ final OutputStream fileOutputStream = new FileOutputStream(snappyFile);
+ final OutputStream out = new
FramedSnappyCompressorOutputStream(fileOutputStream);
+ final InputStream in = new FileInputStream(testFile)
+ ) {
+ ByteStreams.copy(in, out);
}
- try (final InputStream inputStream = CompressionUtils.decompress(new
FileInputStream(snappyFile), snappyFile.getName())) {
+ try (
+ final InputStream fileInputStream = new FileInputStream(snappyFile);
+ final InputStream inputStream =
CompressionUtils.decompress(fileInputStream, snappyFile.getName())
+ ) {
assertGoodDataStream(inputStream);
}
}
@@ -370,10 +410,17 @@ public class CompressionUtilsTest
final File tmpDir = temporaryFolder.newFolder("testDecompressZstd");
final File zstdFile = new File(tmpDir, testFile.getName() + ".zst");
Assert.assertFalse(zstdFile.exists());
- try (final OutputStream out = new ZstdCompressorOutputStream(new
FileOutputStream(zstdFile))) {
- ByteStreams.copy(new FileInputStream(testFile), out);
+ try (
+ final OutputStream fileOutputStream = new FileOutputStream(zstdFile);
+ final OutputStream out = new
ZstdCompressorOutputStream(fileOutputStream);
+ final InputStream in = new FileInputStream(testFile)
+ ) {
+ ByteStreams.copy(in, out);
}
- try (final InputStream inputStream = CompressionUtils.decompress(new
FileInputStream(zstdFile), zstdFile.getName())) {
+ try (
+ final InputStream fileInputStream = new FileInputStream(zstdFile);
+ final InputStream inputStream =
CompressionUtils.decompress(fileInputStream, zstdFile.getName())
+ ) {
assertGoodDataStream(inputStream);
}
}
@@ -384,12 +431,19 @@ public class CompressionUtilsTest
final File tmpDir = temporaryFolder.newFolder("testDecompressZip");
final File zipFile = new File(tmpDir, testFile.getName() + ".zip");
Assert.assertFalse(zipFile.exists());
- try (final ZipOutputStream out = new ZipOutputStream(new
FileOutputStream(zipFile))) {
+ try (
+ final OutputStream fileOutputStream = new FileOutputStream(zipFile);
+ final ZipOutputStream out = new ZipOutputStream(fileOutputStream);
+ final InputStream in = new FileInputStream(testFile)
+ ) {
out.putNextEntry(new ZipEntry("cool.file"));
- ByteStreams.copy(new FileInputStream(testFile), out);
+ ByteStreams.copy(in, out);
out.closeEntry();
}
- try (final InputStream inputStream = CompressionUtils.decompress(new
FileInputStream(zipFile), zipFile.getName())) {
+ try (
+ final InputStream fileInputStream = new FileInputStream(zipFile);
+ final InputStream inputStream =
CompressionUtils.decompress(fileInputStream, zipFile.getName())
+ ) {
assertGoodDataStream(inputStream);
}
}
@@ -401,7 +455,10 @@ public class CompressionUtilsTest
final File zipFile = new File(tmpDir, testFile.getName() + ".zip");
writeZipWithManyFiles(zipFile);
- try (final InputStream inputStream = CompressionUtils.decompress(new
FileInputStream(zipFile), zipFile.getName())) {
+ try (
+ final InputStream fileInputStream = new FileInputStream(zipFile);
+ final InputStream inputStream =
CompressionUtils.decompress(fileInputStream, zipFile.getName())
+ ) {
// Should read the first file, which contains a single null byte.
Assert.assertArrayEquals(new byte[]{0},
ByteStreams.toByteArray(inputStream));
}
@@ -453,16 +510,26 @@ public class CompressionUtilsTest
final File tmpDir = temporaryFolder.newFolder("testGoodGZStream");
final File gzFile = new File(tmpDir, testFile.getName() + ".gz");
Assert.assertFalse(gzFile.exists());
- CompressionUtils.gzip(new FileInputStream(testFile), new
FileOutputStream(gzFile));
+ try (
+ final InputStream inputStream = new FileInputStream(testFile);
+ final OutputStream outputStream = new FileOutputStream(gzFile)
+ ) {
+ CompressionUtils.gzip(inputStream, outputStream);
+ }
Assert.assertTrue(gzFile.exists());
- try (final InputStream inputStream = new GZIPInputStream(new
FileInputStream(gzFile))) {
+ try (
+ final InputStream fileInputStream = new FileInputStream(gzFile);
+ final InputStream inputStream = new GZIPInputStream(fileInputStream)
+ ) {
assertGoodDataStream(inputStream);
}
if (!testFile.delete()) {
throw new IOE("Unable to delete file [%s]", testFile.getAbsolutePath());
}
Assert.assertFalse(testFile.exists());
- CompressionUtils.gunzip(new FileInputStream(gzFile), testFile);
+ try (final InputStream inputStream = new FileInputStream(gzFile)) {
+ CompressionUtils.gunzip(inputStream, testFile);
+ }
Assert.assertTrue(testFile.exists());
try (final InputStream inputStream = new FileInputStream(testFile)) {
assertGoodDataStream(inputStream);
@@ -504,8 +571,8 @@ public class CompressionUtilsTest
java.nio.file.Files.deleteIfExists(evilZip.toPath());
CompressionUtilsTest.makeEvilZip(evilZip);
- try {
- CompressionUtils.unzip(new FileInputStream(evilZip), tmpDir);
+ try (final InputStream inputStream = new FileInputStream(evilZip)) {
+ CompressionUtils.unzip(inputStream, tmpDir);
}
catch (ISE ise) {
Assert.assertTrue(ise.getMessage().contains("does not start with
outDir"));
@@ -741,7 +808,10 @@ public class CompressionUtilsTest
}, Predicates.alwaysTrue()
);
Assert.assertTrue(gzFile.exists());
- try (final InputStream inputStream = CompressionUtils.decompress(new
FileInputStream(gzFile), "file.gz")) {
+ try (
+ final InputStream fileInputStream = new FileInputStream(gzFile);
+ final InputStream inputStream =
CompressionUtils.decompress(fileInputStream, "file.gz")
+ ) {
assertGoodDataStream(inputStream);
}
if (!testFile.delete()) {
@@ -763,8 +833,9 @@ public class CompressionUtilsTest
final File gzFile = new File(tmpDir, testFile.getName() + ".gz");
Assert.assertFalse(gzFile.exists());
final AtomicLong flushes = new AtomicLong(0L);
- CompressionUtils.gzip(
- new FileInputStream(testFile), new FileOutputStream(gzFile)
+ try (
+ final InputStream inputStream = new FileInputStream(testFile);
+ final OutputStream outputStream = new FileOutputStream(gzFile)
{
@Override
public void flush() throws IOException
@@ -776,7 +847,9 @@ public class CompressionUtilsTest
}
}
}
- );
+ ) {
+ CompressionUtils.gzip(inputStream, outputStream);
+ }
}
@Test(expected = IOException.class)
@@ -787,7 +860,10 @@ public class CompressionUtilsTest
Assert.assertFalse(gzFile.exists());
CompressionUtils.gzip(Files.asByteSource(testFile),
Files.asByteSink(gzFile), Predicates.alwaysTrue());
Assert.assertTrue(gzFile.exists());
- try (final InputStream inputStream = CompressionUtils.decompress(new
FileInputStream(gzFile), "file.gz")) {
+ try (
+ final InputStream fileInputStream = new FileInputStream(gzFile);
+ final InputStream inputStream =
CompressionUtils.decompress(fileInputStream, "file.gz")
+ ) {
assertGoodDataStream(inputStream);
}
if (testFile.exists() && !testFile.delete()) {
@@ -795,8 +871,9 @@ public class CompressionUtilsTest
}
Assert.assertFalse(testFile.exists());
final AtomicLong flushes = new AtomicLong(0L);
- CompressionUtils.gunzip(
- new FileInputStream(gzFile), new FilterOutputStream(
+ try (
+ final InputStream inputStream = new FileInputStream(gzFile);
+ final OutputStream outputStream = new FilterOutputStream(
new FileOutputStream(testFile)
{
@Override
@@ -810,7 +887,9 @@ public class CompressionUtilsTest
}
}
)
- );
+ ) {
+ CompressionUtils.gunzip(inputStream, outputStream);
+ }
}
private void verifyUnzip(
diff --git
a/processing/src/test/java/org/apache/druid/query/aggregation/AggregationTestHelper.java
b/processing/src/test/java/org/apache/druid/query/aggregation/AggregationTestHelper.java
index b8a31fe584b..9269a04b1a2 100644
---
a/processing/src/test/java/org/apache/druid/query/aggregation/AggregationTestHelper.java
+++
b/processing/src/test/java/org/apache/druid/query/aggregation/AggregationTestHelper.java
@@ -366,17 +366,19 @@ public class AggregationTestHelper implements Closeable
int maxRowCount
) throws Exception
{
- createIndex(
- new FileInputStream(inputDataFile),
- inputSchema,
- inputFormat,
- aggregators,
- outDir,
- minTimestamp,
- gran,
- maxRowCount,
- true
- );
+ try (final InputStream inputStream = new FileInputStream(inputDataFile)) {
+ createIndex(
+ inputStream,
+ inputSchema,
+ inputFormat,
+ aggregators,
+ outDir,
+ minTimestamp,
+ gran,
+ maxRowCount,
+ true
+ );
+ }
}
public void createIndex(
@@ -391,17 +393,19 @@ public class AggregationTestHelper implements Closeable
boolean rollup
) throws Exception
{
- createIndex(
- new FileInputStream(inputDataFile),
- inputSchema,
- inputFormat,
- aggregators,
- outDir,
- minTimestamp,
- gran,
- maxRowCount,
- rollup
- );
+ try (final InputStream inputStream = new FileInputStream(inputDataFile)) {
+ createIndex(
+ inputStream,
+ inputSchema,
+ inputFormat,
+ aggregators,
+ outDir,
+ minTimestamp,
+ gran,
+ maxRowCount,
+ rollup
+ );
+ }
}
public void createIndex(
@@ -704,4 +708,3 @@ public class AggregationTestHelper implements Closeable
resourceCloser.close();
}
}
-
diff --git
a/server/src/test/java/org/apache/druid/metadata/input/SqlEntityTest.java
b/server/src/test/java/org/apache/druid/metadata/input/SqlEntityTest.java
index 4053ec9ecf2..f86abf12326 100644
--- a/server/src/test/java/org/apache/druid/metadata/input/SqlEntityTest.java
+++ b/server/src/test/java/org/apache/druid/metadata/input/SqlEntityTest.java
@@ -66,15 +66,17 @@ public class SqlEntityTest
SqlTestUtils testUtils = new SqlTestUtils(derbyConnector);
final InputRow expectedRow = testUtils.createTableWithRows(TABLE_NAME_1,
1).get(0);
File tmpFile = File.createTempFile("testQueryResults", "");
- InputEntity.CleanableFile queryResult = SqlEntity.openCleanableFile(
+ final String actualJson;
+ try (final InputEntity.CleanableFile queryResult =
SqlEntity.openCleanableFile(
VALID_SQL,
testUtils.getDerbyInputSourceConnector(),
mapper,
true,
tmpFile
);
- InputStream queryInputStream = new FileInputStream(queryResult.file());
- String actualJson = IOUtils.toString(queryInputStream,
StandardCharsets.UTF_8);
+ final InputStream queryInputStream = new
FileInputStream(queryResult.file())) {
+ actualJson = IOUtils.toString(queryInputStream, StandardCharsets.UTF_8);
+ }
String expectedJson = mapper.writeValueAsString(
Collections.singletonList(((MapBasedInputRow) expectedRow).getEvent())
);
diff --git a/server/src/test/java/org/apache/druid/rpc/RequestBuilderTest.java
b/server/src/test/java/org/apache/druid/rpc/RequestBuilderTest.java
index ca95284ac56..03ca63d6198 100644
--- a/server/src/test/java/org/apache/druid/rpc/RequestBuilderTest.java
+++ b/server/src/test/java/org/apache/druid/rpc/RequestBuilderTest.java
@@ -129,10 +129,12 @@ public class RequestBuilderTest
Assert.assertTrue(request.hasContent());
// Read and verify content.
- Assert.assertEquals(
- json,
- StringUtils.fromUtf8(ByteStreams.toByteArray(new
ChannelBufferInputStream(request.getContent())))
- );
+ try (final ChannelBufferInputStream inputStream = new
ChannelBufferInputStream(request.getContent())) {
+ Assert.assertEquals(
+ json,
+ StringUtils.fromUtf8(ByteStreams.toByteArray(inputStream))
+ );
+ }
}
@Test
@@ -151,10 +153,12 @@ public class RequestBuilderTest
Assert.assertTrue(request.hasContent());
// Read and verify content.
- Assert.assertEquals(
- "{\"foo\":3}",
- StringUtils.fromUtf8(ByteStreams.toByteArray(new
ChannelBufferInputStream(request.getContent())))
- );
+ try (final ChannelBufferInputStream inputStream = new
ChannelBufferInputStream(request.getContent())) {
+ Assert.assertEquals(
+ "{\"foo\":3}",
+ StringUtils.fromUtf8(ByteStreams.toByteArray(inputStream))
+ );
+ }
}
@Test
diff --git
a/server/src/test/java/org/apache/druid/server/QueryResourceTest.java
b/server/src/test/java/org/apache/druid/server/QueryResourceTest.java
index e01e6fe073a..6c39227d9d0 100644
--- a/server/src/test/java/org/apache/druid/server/QueryResourceTest.java
+++ b/server/src/test/java/org/apache/druid/server/QueryResourceTest.java
@@ -1038,11 +1038,16 @@ public class QueryResourceTest
@Test
public void testResourceLimitExceeded() throws IOException
{
- Response response = queryResource.doPost(
- new ExceptionalInputStream(() -> new
ResourceLimitExceededException("You require too much of something")),
- null /*pretty*/,
- testServletRequest
- );
+ final Response response;
+ try (final ExceptionalInputStream inputStream = new ExceptionalInputStream(
+ () -> new ResourceLimitExceededException("You require too much of
something")
+ )) {
+ response = queryResource.doPost(
+ inputStream,
+ null /*pretty*/,
+ testServletRequest
+ );
+ }
Assert.assertNotNull(response);
Assert.assertEquals(Status.BAD_REQUEST.getStatusCode(),
response.getStatus());
QueryException e = jsonMapper.readValue((byte[]) response.getEntity(),
QueryException.class);
@@ -1054,11 +1059,16 @@ public class QueryResourceTest
public void testUnsupportedQueryThrowsException() throws IOException
{
String errorMessage = "This will be support in Druid 9999";
- Response response = queryResource.doPost(
- new ExceptionalInputStream(() -> new
QueryUnsupportedException(errorMessage)),
- null /*pretty*/,
- testServletRequest
- );
+ final Response response;
+ try (final ExceptionalInputStream inputStream = new ExceptionalInputStream(
+ () -> new QueryUnsupportedException(errorMessage)
+ )) {
+ response = queryResource.doPost(
+ inputStream,
+ null /*pretty*/,
+ testServletRequest
+ );
+ }
Assert.assertNotNull(response);
Assert.assertEquals(QueryUnsupportedException.STATUS_CODE,
response.getStatus());
QueryException ex = jsonMapper.readValue((byte[]) response.getEntity(),
QueryException.class);
diff --git
a/server/src/test/java/org/apache/druid/server/initialization/JettyTest.java
b/server/src/test/java/org/apache/druid/server/initialization/JettyTest.java
index d2f3c5078bc..67b56318093 100644
--- a/server/src/test/java/org/apache/druid/server/initialization/JettyTest.java
+++ b/server/src/test/java/org/apache/druid/server/initialization/JettyTest.java
@@ -290,31 +290,39 @@ public class JettyTest extends BaseJettyTest
final HttpURLConnection get = (HttpURLConnection) url.openConnection();
get.setRequestProperty("Accept-Encoding", "gzip");
Assert.assertEquals("gzip", get.getContentEncoding());
- Assert.assertEquals(
- DEFAULT_RESPONSE_CONTENT,
- IOUtils.toString(new GZIPInputStream(get.getInputStream()),
StandardCharsets.UTF_8)
- );
+ try (final InputStream inputStream = new
GZIPInputStream(get.getInputStream())) {
+ Assert.assertEquals(
+ DEFAULT_RESPONSE_CONTENT,
+ IOUtils.toString(inputStream, StandardCharsets.UTF_8)
+ );
+ }
final HttpURLConnection post = (HttpURLConnection) url.openConnection();
post.setRequestProperty("Accept-Encoding", "gzip");
post.setRequestMethod("POST");
Assert.assertEquals("gzip", post.getContentEncoding());
- Assert.assertEquals(
- DEFAULT_RESPONSE_CONTENT,
- IOUtils.toString(new GZIPInputStream(post.getInputStream()),
StandardCharsets.UTF_8)
- );
+ try (final InputStream inputStream = new
GZIPInputStream(post.getInputStream())) {
+ Assert.assertEquals(
+ DEFAULT_RESPONSE_CONTENT,
+ IOUtils.toString(inputStream, StandardCharsets.UTF_8)
+ );
+ }
final HttpURLConnection getNoGzip = (HttpURLConnection)
url.openConnection();
Assert.assertNotEquals("gzip", getNoGzip.getContentEncoding());
- Assert.assertEquals(DEFAULT_RESPONSE_CONTENT,
IOUtils.toString(getNoGzip.getInputStream(), StandardCharsets.UTF_8));
+ try (final InputStream inputStream = getNoGzip.getInputStream()) {
+ Assert.assertEquals(DEFAULT_RESPONSE_CONTENT,
IOUtils.toString(inputStream, StandardCharsets.UTF_8));
+ }
final HttpURLConnection postNoGzip = (HttpURLConnection)
url.openConnection();
postNoGzip.setRequestMethod("POST");
Assert.assertNotEquals("gzip", postNoGzip.getContentEncoding());
- Assert.assertEquals(
- DEFAULT_RESPONSE_CONTENT,
- IOUtils.toString(postNoGzip.getInputStream(), StandardCharsets.UTF_8)
- );
+ try (final InputStream inputStream = postNoGzip.getInputStream()) {
+ Assert.assertEquals(
+ DEFAULT_RESPONSE_CONTENT,
+ IOUtils.toString(inputStream, StandardCharsets.UTF_8)
+ );
+ }
}
// Tests that threads are not stuck when partial chunk is not finalized
diff --git
a/sql/src/test/java/org/apache/druid/sql/avatica/DruidAvaticaHandlerTest.java
b/sql/src/test/java/org/apache/druid/sql/avatica/DruidAvaticaHandlerTest.java
index 20c46f3a108..416f04624c1 100644
---
a/sql/src/test/java/org/apache/druid/sql/avatica/DruidAvaticaHandlerTest.java
+++
b/sql/src/test/java/org/apache/druid/sql/avatica/DruidAvaticaHandlerTest.java
@@ -206,7 +206,7 @@ public class DruidAvaticaHandlerTest extends CalciteTestBase
);
}
- private class ServerWrapper
+ private class ServerWrapper implements AutoCloseable
{
final DruidMeta druidMeta;
final Server server;
@@ -248,6 +248,7 @@ public class DruidAvaticaHandlerTest extends CalciteTestBase
// return DriverManager.getConnection(url);
//}
+ @Override
public void close() throws Exception
{
druidMeta.closeAllConnections();
@@ -963,7 +964,9 @@ public class DruidAvaticaHandlerTest extends CalciteTestBase
@Test
public void testTooManyStatements() throws SQLException
{
+ // Leave these statements open until tearDown closes client so the test
reaches the configured limit.
for (int i = 0; i < STATEMENT_LIMIT; i++) {
+ // codeql[java/database-resource-leak]
client.createStatement();
}
@@ -978,7 +981,9 @@ public class DruidAvaticaHandlerTest extends CalciteTestBase
public void testNotTooManyStatementsWhenYouCloseThem() throws SQLException
{
for (int i = 0; i < STATEMENT_LIMIT * 2; i++) {
- client.createStatement().close();
+ try (final Statement ignored = client.createStatement()) {
+ // Closing each statement is the behavior under test.
+ }
}
}
@@ -1037,7 +1042,7 @@ public class DruidAvaticaHandlerTest extends
CalciteTestBase
public void testNotTooManyStatementsWhenClosed()
{
for (int i = 0; i < 50; i++) {
- try (Statement statement = client.createStatement()) {
+ try (final Statement statement = client.createStatement()) {
statement.executeQuery("SELECT SUM(nonexistent) FROM druid.foo");
Assert.fail();
}
@@ -1051,11 +1056,13 @@ public class DruidAvaticaHandlerTest extends
CalciteTestBase
public void testAutoReconnectOnNoSuchConnection() throws SQLException
{
for (int i = 0; i < 50; i++) {
- final ResultSet resultSet =
client.createStatement().executeQuery("SELECT COUNT(*) AS cnt FROM druid.foo");
- Assert.assertEquals(
- ImmutableList.of(ImmutableMap.of("cnt", 6L)),
- getRows(resultSet)
- );
+ try (final Statement statement = client.createStatement()) {
+ final ResultSet resultSet = statement.executeQuery("SELECT COUNT(*) AS
cnt FROM druid.foo");
+ Assert.assertEquals(
+ ImmutableList.of(ImmutableMap.of("cnt", 6L)),
+ getRows(resultSet)
+ );
+ }
server.druidMeta.closeAllConnections();
}
}
@@ -1063,9 +1070,14 @@ public class DruidAvaticaHandlerTest extends
CalciteTestBase
@Test
public void testTooManyConnections() throws SQLException
{
+ // Keep one statement open on each connection until tearDown so all
connection slots remain occupied.
+ // codeql[java/database-resource-leak]
client.createStatement();
+ // codeql[java/database-resource-leak]
clientLosAngeles.createStatement();
+ // codeql[java/database-resource-leak]
superuserClient.createStatement();
+ // codeql[java/database-resource-leak]
clientNoTrailingSlash.createStatement();
AvaticaClientRuntimeException ex = Assert.assertThrows(
@@ -1090,11 +1102,15 @@ public class DruidAvaticaHandlerTest extends
CalciteTestBase
for (int i = 0; i < CONNECTION_LIMIT * 2; i++) {
try (Connection connection = server.getUserConnection()) {
// Note: NOT in a try-catch block. Let the connection close the
statement
+ // codeql[java/database-resource-leak]
final Statement statement = connection.createStatement();
// Again, NOT in a try-catch block: let the statement close the
// result set.
- final ResultSet resultSet = statement.executeQuery("SELECT COUNT(*) AS
cnt FROM druid.foo");
+ // codeql[java/database-resource-leak]
+ final ResultSet resultSet = statement.executeQuery(
+ "SELECT COUNT(*) AS cnt FROM druid.foo"
+ );
Assert.assertTrue(resultSet.next());
}
}
@@ -1157,30 +1173,27 @@ public class DruidAvaticaHandlerTest extends
CalciteTestBase
}
};
- ServerWrapper server = new ServerWrapper(smallFrameDruidMeta);
- Connection smallFrameClient = server.getUserConnection();
-
- final ResultSet resultSet =
smallFrameClient.createStatement().executeQuery(
- "SELECT dim1 FROM druid.foo"
- );
- List<Map<String, Object>> rows = getRows(resultSet);
- Assert.assertEquals(2, frames.size());
- Assert.assertEquals(
- ImmutableList.of(
- ImmutableMap.of("dim1", ""),
- ImmutableMap.of("dim1", "10.1"),
- ImmutableMap.of("dim1", "2"),
- ImmutableMap.of("dim1", "1"),
- ImmutableMap.of("dim1", "def"),
- ImmutableMap.of("dim1", "abc")
- ),
- rows
- );
-
- resultSet.close();
- smallFrameClient.close();
- exec.shutdown();
- server.close();
+ try (final ServerWrapper server = new ServerWrapper(smallFrameDruidMeta);
+ final Connection smallFrameClient = server.getUserConnection();
+ final Statement statement = smallFrameClient.createStatement();
+ final ResultSet resultSet = statement.executeQuery("SELECT dim1 FROM
druid.foo")) {
+ final List<Map<String, Object>> rows = getRows(resultSet);
+ Assert.assertEquals(2, frames.size());
+ Assert.assertEquals(
+ ImmutableList.of(
+ ImmutableMap.of("dim1", ""),
+ ImmutableMap.of("dim1", "10.1"),
+ ImmutableMap.of("dim1", "2"),
+ ImmutableMap.of("dim1", "1"),
+ ImmutableMap.of("dim1", "def"),
+ ImmutableMap.of("dim1", "abc")
+ ),
+ rows
+ );
+ }
+ finally {
+ exec.shutdown();
+ }
}
@Test
@@ -1219,33 +1232,32 @@ public class DruidAvaticaHandlerTest extends
CalciteTestBase
}
};
- ServerWrapper server = new ServerWrapper(smallFrameDruidMeta);
- Connection smallFrameClient = server.getUserConnection();
-
- // use a prepared statement because Avatica currently ignores fetchSize on
the initial fetch of a Statement
- PreparedStatement statement = smallFrameClient.prepareStatement("SELECT
dim1 FROM druid.foo");
- // set a fetch size below the minimum configured threshold
- statement.setFetchSize(2);
- final ResultSet resultSet = statement.executeQuery();
- List<Map<String, Object>> rows = getRows(resultSet);
- // expect minimum threshold to be used, which should be enough to do this
all in first fetch
- Assert.assertEquals(0, frames.size());
- Assert.assertEquals(
- ImmutableList.of(
- ImmutableMap.of("dim1", ""),
- ImmutableMap.of("dim1", "10.1"),
- ImmutableMap.of("dim1", "2"),
- ImmutableMap.of("dim1", "1"),
- ImmutableMap.of("dim1", "def"),
- ImmutableMap.of("dim1", "abc")
- ),
- rows
- );
-
- resultSet.close();
- smallFrameClient.close();
- exec.shutdown();
- server.close();
+ try (final ServerWrapper server = new ServerWrapper(smallFrameDruidMeta);
+ final Connection smallFrameClient = server.getUserConnection();
+ final PreparedStatement statement =
smallFrameClient.prepareStatement("SELECT dim1 FROM druid.foo")) {
+ // use a prepared statement because Avatica currently ignores fetchSize
on the initial fetch of a Statement
+ // set a fetch size below the minimum configured threshold
+ statement.setFetchSize(2);
+ try (final ResultSet resultSet = statement.executeQuery()) {
+ final List<Map<String, Object>> rows = getRows(resultSet);
+ // expect minimum threshold to be used, which should be enough to do
this all in first fetch
+ Assert.assertEquals(0, frames.size());
+ Assert.assertEquals(
+ ImmutableList.of(
+ ImmutableMap.of("dim1", ""),
+ ImmutableMap.of("dim1", "10.1"),
+ ImmutableMap.of("dim1", "2"),
+ ImmutableMap.of("dim1", "1"),
+ ImmutableMap.of("dim1", "def"),
+ ImmutableMap.of("dim1", "abc")
+ ),
+ rows
+ );
+ }
+ }
+ finally {
+ exec.shutdown();
+ }
}
@Test
@@ -1830,29 +1842,30 @@ public class DruidAvaticaHandlerTest extends
CalciteTestBase
}
};
- ServerWrapper server = new ServerWrapper(druidMeta);
- try (Connection conn = server.getUserConnection()) {
-
- // Test with plain JDBC
- try (ResultSet resultSet = conn.createStatement().executeQuery(
- "SELECT dim1 FROM druid.foo")) {
- List<Map<String, Object>> rows = getRows(resultSet);
- Assert.assertEquals(6, rows.size());
- Assert.assertEquals(6, frames.size()); // 3 empty frames and then 3
frames of 2 rows each
-
- Assert.assertFalse(frames.get(0).rows.iterator().hasNext());
- Assert.assertFalse(frames.get(1).rows.iterator().hasNext());
- Assert.assertFalse(frames.get(2).rows.iterator().hasNext());
- Assert.assertTrue(frames.get(3).rows.iterator().hasNext());
- Assert.assertTrue(frames.get(4).rows.iterator().hasNext());
- Assert.assertTrue(frames.get(5).rows.iterator().hasNext());
+ try (final ServerWrapper server = new ServerWrapper(druidMeta)) {
+ try (final Connection conn = server.getUserConnection()) {
+
+ // Test with plain JDBC
+ try (final Statement statement = conn.createStatement();
+ final ResultSet resultSet = statement.executeQuery("SELECT dim1
FROM druid.foo")) {
+ final List<Map<String, Object>> rows = getRows(resultSet);
+ Assert.assertEquals(6, rows.size());
+ Assert.assertEquals(6, frames.size()); // 3 empty frames and then 3
frames of 2 rows each
+
+ Assert.assertFalse(frames.get(0).rows.iterator().hasNext());
+ Assert.assertFalse(frames.get(1).rows.iterator().hasNext());
+ Assert.assertFalse(frames.get(2).rows.iterator().hasNext());
+ Assert.assertTrue(frames.get(3).rows.iterator().hasNext());
+ Assert.assertTrue(frames.get(4).rows.iterator().hasNext());
+ Assert.assertTrue(frames.get(5).rows.iterator().hasNext());
+ }
}
- }
-
- testWithJDBI(server.url);
- exec.shutdown();
- server.close();
+ testWithJDBI(server.url);
+ }
+ finally {
+ exec.shutdown();
+ }
}
@Test
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]