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 @@ -38,6 +38,7 @@
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;
Expand Down Expand Up @@ -95,15 +96,15 @@ public void testGZUncompress() throws IOException, SegmentLoadingException

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);
}

final OSSObject object0 = new OSSObject();
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);
Expand All @@ -115,30 +116,33 @@ public void testGZUncompress() throws IOException, SegmentLoadingException

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
Expand All @@ -151,7 +155,8 @@ public void testGZUncompressRetries() throws IOException, SegmentLoadingExceptio

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);
}

Expand All @@ -160,44 +165,46 @@ public void testGZUncompressRetries() throws IOException, SegmentLoadingExceptio
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));

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());
}
}

}
Original file line number Diff line number Diff line change
Expand Up @@ -300,10 +300,11 @@ public void test_taskLog_fetch() throws IOException

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);
}
Expand All @@ -324,10 +325,11 @@ public void test_taskLog_fetch_withRange() throws IOException

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);
}
Expand All @@ -348,10 +350,11 @@ public void test_taskLog_fetch_withNegativeRange() throws IOException

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);
}
Expand All @@ -373,10 +376,11 @@ public void test_taskReport_fetch() throws IOException

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);
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -117,12 +117,17 @@
@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;
Expand Down Expand Up @@ -171,9 +176,10 @@
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)) {

Check warning

Code scanning / CodeQL

Unsafe use of getResource Warning test

The idiom getClass().getResource() is unsafe for classes that may be extended.
Comment thread
FrankChen021 marked this conversation as resolved.
Dismissed
config = reader.readValue(is, TestConfig.class);
}
}

/**
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -48,6 +48,7 @@

import java.io.File;
import java.io.IOException;
import java.io.InputStream;
import java.sql.Timestamp;
import java.util.ArrayList;
import java.util.List;
Expand Down Expand Up @@ -149,23 +150,28 @@ public void testSimpleDataIngestionAndGroupByTest() throws Exception
.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);

Expand Down
Loading
Loading