diff --git a/codestyle/pmd-ruleset.xml b/codestyle/pmd-ruleset.xml index 7d7285d0a516..ef1858bf2d8b 100644 --- a/codestyle/pmd-ruleset.xml +++ b/codestyle/pmd-ruleset.xml @@ -29,4 +29,6 @@ This ruleset defines the PMD rules for the Apache Druid project. + + diff --git a/extensions-contrib/opentelemetry-emitter/src/main/java/org/apache/druid/emitter/opentelemetry/OpenTelemetryEmitter.java b/extensions-contrib/opentelemetry-emitter/src/main/java/org/apache/druid/emitter/opentelemetry/OpenTelemetryEmitter.java index 2c2691e64f0d..4fbb9a02f46c 100644 --- a/extensions-contrib/opentelemetry-emitter/src/main/java/org/apache/druid/emitter/opentelemetry/OpenTelemetryEmitter.java +++ b/extensions-contrib/opentelemetry-emitter/src/main/java/org/apache/druid/emitter/opentelemetry/OpenTelemetryEmitter.java @@ -82,7 +82,7 @@ private void emitQueryTimeEvent(ServiceMetricEvent event) { Context opentelemetryContext = propagator.extract(Context.current(), event, DRUID_CONTEXT_TEXT_MAP_GETTER); - try (Scope scope = opentelemetryContext.makeCurrent()) { + try (Scope ignoredScope = opentelemetryContext.makeCurrent()) { DateTime endTime = event.getCreatedTime(); DateTime startTime = endTime.minusMillis(event.getValue().intValue()); diff --git a/extensions-contrib/rabbit-stream-indexing-service/src/main/java/org/apache/druid/indexing/rabbitstream/RabbitStreamRecordSupplier.java b/extensions-contrib/rabbit-stream-indexing-service/src/main/java/org/apache/druid/indexing/rabbitstream/RabbitStreamRecordSupplier.java index b2a2432fc9c0..a23dd60a1d0f 100644 --- a/extensions-contrib/rabbit-stream-indexing-service/src/main/java/org/apache/druid/indexing/rabbitstream/RabbitStreamRecordSupplier.java +++ b/extensions-contrib/rabbit-stream-indexing-service/src/main/java/org/apache/druid/indexing/rabbitstream/RabbitStreamRecordSupplier.java @@ -275,11 +275,14 @@ private void filterBufferAndResetBackgroundFetch(Set> pa { this.stopBackgroundFetch(); // filter records in buffer and only retain ones whose partition was not seeked - BlockingQueue> newQ = new LinkedBlockingQueue<>( + final Set partitionIds = partitions.stream() + .map(StreamPartition::getPartitionId) + .collect(Collectors.toSet()); + final BlockingQueue> newQ = new LinkedBlockingQueue<>( recordBufferSize); queue.stream() - .filter(x -> !streamBuilders.containsKey(x.getPartitionId())) + .filter(x -> !partitionIds.contains(x.getPartitionId())) .forEachOrdered(newQ::offer); queue = newQ; diff --git a/extensions-contrib/rabbit-stream-indexing-service/src/test/java/org/apache/druid/indexing/rabbitstream/RabbitStreamRecordSupplierTest.java b/extensions-contrib/rabbit-stream-indexing-service/src/test/java/org/apache/druid/indexing/rabbitstream/RabbitStreamRecordSupplierTest.java index 7c6beda72cbb..70501b941ef6 100644 --- a/extensions-contrib/rabbit-stream-indexing-service/src/test/java/org/apache/druid/indexing/rabbitstream/RabbitStreamRecordSupplierTest.java +++ b/extensions-contrib/rabbit-stream-indexing-service/src/test/java/org/apache/druid/indexing/rabbitstream/RabbitStreamRecordSupplierTest.java @@ -27,6 +27,7 @@ import com.rabbitmq.stream.ConsumerBuilder; import com.rabbitmq.stream.Environment; import com.rabbitmq.stream.EnvironmentBuilder; +import com.rabbitmq.stream.Message; import com.rabbitmq.stream.MessageHandler; import com.rabbitmq.stream.OffsetSpecification; import com.rabbitmq.stream.codec.WrapperMessageBuilder; @@ -386,6 +387,51 @@ public void testSeek() } + @Test + public void testSeekRetainsBufferedRecordsForOtherPartitions() + { + final StreamPartition partition0 = StreamPartition.of(STREAM, PARTITION_ID0); + final StreamPartition partition1 = StreamPartition.of(STREAM, PARTITION_ID1); + final Set> partitions = ImmutableSet.of(partition0, partition1); + final RabbitStreamRecordSupplier recordSupplier = makeRecordSupplierWithMockedEnvironment(uri, null); + + EasyMock.expect(environmentBuilder.uri("rabbitmq-stream://localhost:5552")).andReturn(environmentBuilder).once(); + EasyMock.expect(environmentBuilder.build()).andStubReturn(environment); + + final ConsumerBuilder consumerBuilder0 = createMock(ConsumerBuilder.class); + EasyMock.expect(environment.consumerBuilder()).andReturn(consumerBuilder0).once(); + EasyMock.expect(consumerBuilder0.noTrackingStrategy()).andReturn(consumerBuilder0).once(); + EasyMock.expect(consumerBuilder0.stream(PARTITION_ID0)).andReturn(consumerBuilder0).once(); + EasyMock.expect(consumerBuilder0.messageHandler(recordSupplier)).andReturn(consumerBuilder0).once(); + + final ConsumerBuilder consumerBuilder1 = createMock(ConsumerBuilder.class); + EasyMock.expect(environment.consumerBuilder()).andReturn(consumerBuilder1).once(); + EasyMock.expect(consumerBuilder1.noTrackingStrategy()).andReturn(consumerBuilder1).once(); + EasyMock.expect(consumerBuilder1.stream(PARTITION_ID1)).andReturn(consumerBuilder1).once(); + EasyMock.expect(consumerBuilder1.messageHandler(recordSupplier)).andReturn(consumerBuilder1).once(); + + replayAll(); + recordSupplier.assign(partitions); + + final WrapperMessageBuilder messageBuilder = new WrapperMessageBuilder(); + messageBuilder.addData("record".getBytes(StandardCharsets.UTF_8)); + final Message rabbitMessage = messageBuilder.build(); + for (int i = 0; i < 50; i++) { + recordSupplier.handle(new MessageHandlerContext(i, 0, 0, PARTITION_ID0), rabbitMessage); + recordSupplier.handle(new MessageHandlerContext(i, 0, 0, PARTITION_ID1), rabbitMessage); + } + + recordSupplier.seek(partition0, 10L); + + final List> messages = recordSupplier.poll(0); + Assert.assertEquals(50, messages.size()); + Assert.assertTrue(messages.stream().allMatch(message -> PARTITION_ID1.equals(message.getPartitionId()))); + + recordSupplier.close(); + verifyAll(); + } + + @Test public void testPollBothPartitions() diff --git a/extensions-core/google-extensions/src/main/java/org/apache/druid/storage/google/GoogleTaskLogs.java b/extensions-core/google-extensions/src/main/java/org/apache/druid/storage/google/GoogleTaskLogs.java index 2825f4033350..f162fcb62b06 100644 --- a/extensions-core/google-extensions/src/main/java/org/apache/druid/storage/google/GoogleTaskLogs.java +++ b/extensions-core/google-extensions/src/main/java/org/apache/druid/storage/google/GoogleTaskLogs.java @@ -119,24 +119,24 @@ private void pushTaskFile(final File logFile, final String taskKey) throws IOExc public Optional streamTaskLog(final String taskid, final long offset) throws IOException { final String taskKey = getTaskLogKey(taskid); - return streamTaskFile(taskid, offset, taskKey); + return streamTaskFile(offset, taskKey); } @Override public Optional streamTaskReports(String taskid) throws IOException { final String taskKey = getTaskReportKey(taskid); - return streamTaskFile(taskid, 0, taskKey); + return streamTaskFile(0, taskKey); } @Override public Optional streamTaskStatus(String taskid) throws IOException { final String taskKey = getTaskStatusKey(taskid); - return streamTaskFile(taskid, 0, taskKey); + return streamTaskFile(0, taskKey); } - private Optional streamTaskFile(final String taskid, final long offset, String taskKey) + private Optional streamTaskFile(final long offset, String taskKey) throws IOException { try { diff --git a/indexing-service/src/main/java/org/apache/druid/indexing/overlord/sampler/InputSourceSampler.java b/indexing-service/src/main/java/org/apache/druid/indexing/overlord/sampler/InputSourceSampler.java index 38a79684dcbd..71afe07da87b 100644 --- a/indexing-service/src/main/java/org/apache/druid/indexing/overlord/sampler/InputSourceSampler.java +++ b/indexing-service/src/main/java/org/apache/druid/indexing/overlord/sampler/InputSourceSampler.java @@ -127,7 +127,7 @@ public SamplerResponse sample( ); try (final CloseableIterator iterator = reader.sample(); final IncrementalIndex index = buildIncrementalIndex(nonNullSamplerConfig, nonNullDataSchema); - final Closer closer1 = closer) { + final Closer ignoredCloser = closer) { List responseRows = new ArrayList<>(nonNullSamplerConfig.getNumRows()); int numRowsIndexed = 0; diff --git a/indexing-service/src/main/java/org/apache/druid/indexing/seekablestream/supervisor/SeekableStreamSupervisor.java b/indexing-service/src/main/java/org/apache/druid/indexing/seekablestream/supervisor/SeekableStreamSupervisor.java index e2a5a4b1ccdf..58f92ddaab7f 100644 --- a/indexing-service/src/main/java/org/apache/druid/indexing/seekablestream/supervisor/SeekableStreamSupervisor.java +++ b/indexing-service/src/main/java/org/apache/druid/indexing/seekablestream/supervisor/SeekableStreamSupervisor.java @@ -1924,9 +1924,6 @@ private List getCurrentParseErrors() } } - SeekableStreamIndexTaskTuningConfig ss = spec.getSpec().getTuningConfig().convertToTaskTuningConfig(); - SeekableStreamSupervisorIOConfig oo = spec.getSpec().getIOConfig(); - // store a limited number of parse exceptions, keeping the most recent ones int parseErrorLimit = spec.getSpec().getTuningConfig().convertToTaskTuningConfig().getMaxSavedParseExceptions() * spec.getSpec().getIOConfig().getTaskCount(); diff --git a/indexing-service/src/main/java/org/apache/druid/indexing/worker/shuffle/LocalIntermediaryDataManager.java b/indexing-service/src/main/java/org/apache/druid/indexing/worker/shuffle/LocalIntermediaryDataManager.java index 1974b096be24..f45e442a4bb4 100644 --- a/indexing-service/src/main/java/org/apache/druid/indexing/worker/shuffle/LocalIntermediaryDataManager.java +++ b/indexing-service/src/main/java/org/apache/druid/indexing/worker/shuffle/LocalIntermediaryDataManager.java @@ -301,7 +301,7 @@ public DataSegment addSegment(String supervisorTaskId, String subTaskId, DataSeg final BucketNumberedShardSpec bucketNumberedShardSpec = (BucketNumberedShardSpec) segment.getShardSpec(); //noinspection unused - try (final Closer resourceCloser = closer) { + try (final Closer ignoredCloser = closer) { FileUtils.mkdirp(taskTempDir); // Temporary compressed file. Will be removed when taskTempDir is deleted. diff --git a/processing/src/main/java/org/apache/druid/common/utils/SocketUtil.java b/processing/src/main/java/org/apache/druid/common/utils/SocketUtil.java index 8eb38776e9e7..0eac24c83fff 100644 --- a/processing/src/main/java/org/apache/druid/common/utils/SocketUtil.java +++ b/processing/src/main/java/org/apache/druid/common/utils/SocketUtil.java @@ -41,7 +41,7 @@ public static int findOpenPortFrom(int startPort) int currPort = startPort; while (currPort < 0xffff) { - try (ServerSocket socket = new ServerSocket(currPort)) { + try (ServerSocket ignoredSocket = new ServerSocket(currPort)) { return currPort; } catch (IOException e) { diff --git a/processing/src/main/java/org/apache/druid/frame/processor/FrameProcessors.java b/processing/src/main/java/org/apache/druid/frame/processor/FrameProcessors.java index 4e659b67cff7..09331d9c207b 100644 --- a/processing/src/main/java/org/apache/druid/frame/processor/FrameProcessors.java +++ b/processing/src/main/java/org/apache/druid/frame/processor/FrameProcessors.java @@ -77,8 +77,8 @@ public void cleanup() throws IOException { if (cleanedUp.compareAndSet(false, true)) { //noinspection EmptyTryBlock - try (Closeable ignore1 = baggage; - Closeable ignore2 = processor::cleanup) { + try (Closeable ignoredBaggage = baggage; + Closeable ignoredCleanup = processor::cleanup) { // piggy-back try-with-resources semantics } } diff --git a/processing/src/main/java/org/apache/druid/java/util/common/FileUtils.java b/processing/src/main/java/org/apache/druid/java/util/common/FileUtils.java index f852786cc8b2..7a302a5f4d2e 100644 --- a/processing/src/main/java/org/apache/druid/java/util/common/FileUtils.java +++ b/processing/src/main/java/org/apache/druid/java/util/common/FileUtils.java @@ -264,7 +264,7 @@ public static T writeAtomically(final File file, final File tmpDir, OutputSt final File tmpFile = new File(tmpDir, StringUtils.format(".%s.%s", file.getName(), UUID.randomUUID())); //noinspection unused - try (final Closeable deleter = () -> Files.deleteIfExists(tmpFile.toPath())) { + try (final Closeable ignoredDeleter = () -> Files.deleteIfExists(tmpFile.toPath())) { final T retVal; try ( diff --git a/processing/src/main/java/org/apache/druid/java/util/common/guava/ConcatSequence.java b/processing/src/main/java/org/apache/druid/java/util/common/guava/ConcatSequence.java index 3a02d84471fd..70980931b822 100644 --- a/processing/src/main/java/org/apache/druid/java/util/common/guava/ConcatSequence.java +++ b/processing/src/main/java/org/apache/druid/java/util/common/guava/ConcatSequence.java @@ -142,7 +142,7 @@ public boolean isDone() @Override public void close() throws IOException { - try (Closeable toClose = yielderYielder) { + try (Closeable ignoredYielder = yielderYielder) { yielder.close(); } } diff --git a/server/src/main/java/org/apache/druid/segment/realtime/appenderator/StreamAppenderator.java b/server/src/main/java/org/apache/druid/segment/realtime/appenderator/StreamAppenderator.java index 4b1d2facc40c..d8290aaa9b01 100644 --- a/server/src/main/java/org/apache/druid/segment/realtime/appenderator/StreamAppenderator.java +++ b/server/src/main/java/org/apache/druid/segment/realtime/appenderator/StreamAppenderator.java @@ -521,7 +521,7 @@ private Sink getOrCreateSink(final SegmentIdWithShardSpec identifier) tuningConfig.getIndexSpec(), Collections.emptyList() ); - bytesCurrentlyInMemory.addAndGet(calculateSinkMemoryInUsed(retVal)); + bytesCurrentlyInMemory.addAndGet(calculateSinkMemoryInUsed()); // Add sink prior to announcing it, to ensure it is immediately queryable. addSink(identifier, retVal); @@ -1526,7 +1526,7 @@ private ListenableFuture abandonSegment( // i.e. those that haven't been persisted for *InMemory counters, or pushed to deep storage for the total counter. rowsCurrentlyInMemory.addAndGet(-sink.getNumRowsInMemory()); bytesCurrentlyInMemory.addAndGet(-sink.getBytesInMemory()); - bytesCurrentlyInMemory.addAndGet(-calculateSinkMemoryInUsed(sink)); + bytesCurrentlyInMemory.addAndGet(-calculateSinkMemoryInUsed()); for (FireHydrant hydrant : sink) { // Decrement memory used by all Memory Mapped Hydrant if (!hydrant.equals(sink.getCurrHydrant())) { @@ -1801,7 +1801,7 @@ private int calculateMMappedHydrantMemoryInUsed(FireHydrant hydrant) return total; } - private int calculateSinkMemoryInUsed(Sink sink) + private int calculateSinkMemoryInUsed() { if (skipBytesInMemoryOverheadCheck) { return 0; diff --git a/server/src/main/java/org/apache/druid/server/coordinator/loading/StrategicSegmentAssigner.java b/server/src/main/java/org/apache/druid/server/coordinator/loading/StrategicSegmentAssigner.java index f65ed4c46f47..3fde88fdea67 100644 --- a/server/src/main/java/org/apache/druid/server/coordinator/loading/StrategicSegmentAssigner.java +++ b/server/src/main/java/org/apache/druid/server/coordinator/loading/StrategicSegmentAssigner.java @@ -807,7 +807,7 @@ private int dropReplicas( // Drop as many replicas as possible from decommissioning servers int remainingNumToDrop = numToDrop; int numDropsQueued = - dropReplicasFromServers(remainingNumToDrop, segment, eligibleDyingServers.iterator(), tier); + dropReplicasFromServers(remainingNumToDrop, segment, eligibleDyingServers.iterator()); // Drop replicas from active servers if required if (numToDrop > numDropsQueued) { @@ -816,7 +816,7 @@ private int dropReplicas( (useRoundRobinAssignment || eligibleLiveServers.size() <= remainingNumToDrop) ? eligibleLiveServers.iterator() : strategy.findServersToDropSegment(segment, new ArrayList<>(eligibleLiveServers)); - numDropsQueued += dropReplicasFromServers(remainingNumToDrop, segment, serverIterator, tier); + numDropsQueued += dropReplicasFromServers(remainingNumToDrop, segment, serverIterator); } return numDropsQueued; @@ -829,8 +829,7 @@ private int dropReplicas( private int dropReplicasFromServers( int numToDrop, DataSegment segment, - Iterator serverIterator, - String tier + Iterator serverIterator ) { int numDropsQueued = 0; diff --git a/sql/src/main/java/org/apache/druid/sql/avatica/DruidAvaticaJsonHandler.java b/sql/src/main/java/org/apache/druid/sql/avatica/DruidAvaticaJsonHandler.java index 2e80bfff4e79..3f20195de397 100644 --- a/sql/src/main/java/org/apache/druid/sql/avatica/DruidAvaticaJsonHandler.java +++ b/sql/src/main/java/org/apache/druid/sql/avatica/DruidAvaticaJsonHandler.java @@ -69,7 +69,7 @@ public boolean handle(Request request, Response response, Callback callback) thr String remoteAddr = Request.getRemoteAddr(request); DruidMeta.setThreadLocalRemoteAddress(remoteAddr); - try (Timer.Context ctx = this.requestTimer.start()) { + try (Timer.Context ignoredContext = this.requestTimer.start()) { if (AVATICA_PATH_NO_TRAILING_SLASH.equals(StringUtils.maybeRemoveTrailingSlash(requestURI))) { response.getHeaders().put("Content-Type", "application/json;charset=utf-8"); diff --git a/sql/src/main/java/org/apache/druid/sql/avatica/DruidAvaticaProtobufHandler.java b/sql/src/main/java/org/apache/druid/sql/avatica/DruidAvaticaProtobufHandler.java index 1aee0db7729b..be646c10cf5f 100644 --- a/sql/src/main/java/org/apache/druid/sql/avatica/DruidAvaticaProtobufHandler.java +++ b/sql/src/main/java/org/apache/druid/sql/avatica/DruidAvaticaProtobufHandler.java @@ -73,7 +73,7 @@ public boolean handle(Request request, Response response, Callback callback) thr try { if (AVATICA_PATH_NO_TRAILING_SLASH.equals(StringUtils.maybeRemoveTrailingSlash(requestURI))) { - try (Timer.Context ctx = this.requestTimer.start()) { + try (Timer.Context ignoredContext = this.requestTimer.start()) { if (!"POST".equals(request.getMethod())) { response.setStatus(405); response.write( diff --git a/sql/src/main/java/org/apache/druid/sql/calcite/rule/logical/UnnestInputCleanupRule.java b/sql/src/main/java/org/apache/druid/sql/calcite/rule/logical/UnnestInputCleanupRule.java index 6b3e7aac5ac9..97c2b6981c84 100644 --- a/sql/src/main/java/org/apache/druid/sql/calcite/rule/logical/UnnestInputCleanupRule.java +++ b/sql/src/main/java/org/apache/druid/sql/calcite/rule/logical/UnnestInputCleanupRule.java @@ -83,7 +83,7 @@ public void onMatch(RelOptRuleCall call) newProjects.set(inputIndex, null); - RexNode newUnnestExpr = unnestInput.accept(new ExpressionPullerRexShuttle(newProjects, inputIndex)); + RexNode newUnnestExpr = unnestInput.accept(new ExpressionPullerRexShuttle(newProjects)); if (newUnnestExpr instanceof RexInputRef) { // this won't make it simpler @@ -142,7 +142,7 @@ private static class ExpressionPullerRexShuttle extends RexShuttle { private final List projects; - private ExpressionPullerRexShuttle(List projects, int replaceableIndex) + private ExpressionPullerRexShuttle(List projects) { this.projects = projects; } diff --git a/sql/src/main/java/org/apache/druid/sql/calcite/schema/InformationSchema.java b/sql/src/main/java/org/apache/druid/sql/calcite/schema/InformationSchema.java index a9e3d2e31d2c..94e11de1c720 100644 --- a/sql/src/main/java/org/apache/druid/sql/calcite/schema/InformationSchema.java +++ b/sql/src/main/java/org/apache/druid/sql/calcite/schema/InformationSchema.java @@ -377,8 +377,7 @@ public Iterable apply(final String tableName) return generateColumnMetadata( schemaName, tableName, - table.getRowType(typeFactory), - typeFactory + table.getRowType(typeFactory) ); } } @@ -395,8 +394,7 @@ public Iterable apply(final String functionName) return generateColumnMetadata( schemaName, functionName, - viewMacro.apply(Collections.emptyList()).getRowType(typeFactory), - typeFactory + viewMacro.apply(Collections.emptyList()).getRowType(typeFactory) ); } catch (Exception e) { @@ -442,8 +440,7 @@ public TableType getJdbcTableType() private Iterable generateColumnMetadata( final String schemaName, final String tableName, - final RelDataType tableSchema, - final RelDataTypeFactory typeFactory + final RelDataType tableSchema ) { return FluentIterable