From aff397e16e0e9618d387438c99d652fee9ca851f Mon Sep 17 00:00:00 2001 From: Zoltan Haindrich Date: Thu, 8 Oct 2026 14:21:09 +0200 Subject: [PATCH] fix: release the Coordinator segment response in sys.segments fetchAllUsedSegmentsWithOvershadowedStatus hands out a ResourceHolder that closes the streamed response; MetadataSegmentView and SystemSchema close it, also on early exit, and QueryHandler now closes the bound Interpreter so scans under Filter/Project nodes are released too. No segment left behind. --- .../schema/SysSegmentsTableBenchmark.java | 9 +- .../server/CoordinatorClientTest.java | 9 +- .../client/coordinator/CoordinatorClient.java | 7 +- .../coordinator/CoordinatorClientImpl.java | 22 ++- .../coordinator/NoopCoordinatorClient.java | 5 +- .../CoordinatorClientImplTest.java | 44 ++++-- .../sql/calcite/planner/QueryHandler.java | 90 ++++++------ .../calcite/schema/MetadataSegmentView.java | 62 +++++--- .../sql/calcite/schema/SystemSchema.java | 10 +- .../sql/calcite/planner/QueryHandlerTest.java | 57 ++++++++ .../schema/MetadataSegmentViewTest.java | 135 +++++++++++++++++- .../sql/calcite/schema/SystemSchemaTest.java | 46 +++++- 12 files changed, 397 insertions(+), 99 deletions(-) create mode 100644 sql/src/test/java/org/apache/druid/sql/calcite/planner/QueryHandlerTest.java diff --git a/benchmarks/src/test/java/org/apache/druid/sql/calcite/schema/SysSegmentsTableBenchmark.java b/benchmarks/src/test/java/org/apache/druid/sql/calcite/schema/SysSegmentsTableBenchmark.java index d75f46c89eef..3b85e24f9d1a 100644 --- a/benchmarks/src/test/java/org/apache/druid/sql/calcite/schema/SysSegmentsTableBenchmark.java +++ b/benchmarks/src/test/java/org/apache/druid/sql/calcite/schema/SysSegmentsTableBenchmark.java @@ -35,12 +35,12 @@ import org.apache.druid.client.InternalQueryConfig; import org.apache.druid.client.TimelineServerView; import org.apache.druid.client.coordinator.NoopCoordinatorClient; +import org.apache.druid.collections.ResourceHolder; +import org.apache.druid.collections.StupidResourceHolder; import org.apache.druid.error.NotYetImplemented; import org.apache.druid.jackson.DefaultObjectMapper; -import org.apache.druid.java.util.common.CloseableIterators; import org.apache.druid.java.util.common.Intervals; import org.apache.druid.java.util.common.StringUtils; -import org.apache.druid.java.util.common.parsers.CloseableIterator; import org.apache.druid.segment.join.JoinableFactory; import org.apache.druid.segment.metadata.CentralizedDatasourceSchemaConfig; import org.apache.druid.server.QueryLifecycleFactory; @@ -73,6 +73,7 @@ import java.util.ArrayList; import java.util.Comparator; import java.util.EnumMap; +import java.util.Iterator; import java.util.List; import java.util.Map; import java.util.Set; @@ -180,12 +181,12 @@ public void setup() final NoopCoordinatorClient coordinatorClient = new NoopCoordinatorClient() { @Override - public ListenableFuture> fetchAllUsedSegmentsWithOvershadowedStatus( + public ListenableFuture>> fetchAllUsedSegmentsWithOvershadowedStatus( Set watchedDataSources, boolean includeRealtimeSegments ) { - return Futures.immediateFuture(CloseableIterators.withEmptyBaggage(publishedSegments.iterator())); + return Futures.immediateFuture(StupidResourceHolder.create(publishedSegments.iterator())); } }; diff --git a/embedded-tests/src/test/java/org/apache/druid/testing/embedded/server/CoordinatorClientTest.java b/embedded-tests/src/test/java/org/apache/druid/testing/embedded/server/CoordinatorClientTest.java index 065963f25a37..980b2b256abf 100644 --- a/embedded-tests/src/test/java/org/apache/druid/testing/embedded/server/CoordinatorClientTest.java +++ b/embedded-tests/src/test/java/org/apache/druid/testing/embedded/server/CoordinatorClientTest.java @@ -21,11 +21,11 @@ import org.apache.druid.client.ImmutableSegmentLoadInfo; import org.apache.druid.client.coordinator.CoordinatorClient; +import org.apache.druid.collections.ResourceHolder; import org.apache.druid.common.utils.IdUtils; import org.apache.druid.indexing.common.task.Task; import org.apache.druid.indexing.common.task.TaskBuilder; import org.apache.druid.java.util.common.Intervals; -import org.apache.druid.java.util.common.parsers.CloseableIterator; import org.apache.druid.query.SegmentDescriptor; import org.apache.druid.server.coordinator.rules.ForeverBroadcastDistributionRule; import org.apache.druid.server.coordinator.rules.Rule; @@ -44,9 +44,9 @@ import org.junit.jupiter.api.Test; import org.junit.jupiter.api.Timeout; -import java.io.IOException; import java.net.URI; import java.util.ArrayList; +import java.util.Iterator; import java.util.List; import java.util.Map; import java.util.Set; @@ -154,13 +154,14 @@ public void test_fetchUsedSegments() @Test @Timeout(20) - public void test_fetchAllUsedSegmentsWithOvershadowedStatus() throws IOException + public void test_fetchAllUsedSegmentsWithOvershadowedStatus() { runIndexTask(); - try (CloseableIterator iterator = cluster.callApi().onLeaderCoordinator( + try (ResourceHolder> segments = cluster.callApi().onLeaderCoordinator( c -> c.fetchAllUsedSegmentsWithOvershadowedStatus(Set.of(dataSource), true)) ) { + final Iterator iterator = segments.get(); Assertions.assertTrue(iterator.hasNext()); SegmentStatusInCluster segmentStatus = iterator.next(); Assertions.assertEquals(dataSource, segmentStatus.getDataSegment().getDataSource()); diff --git a/server/src/main/java/org/apache/druid/client/coordinator/CoordinatorClient.java b/server/src/main/java/org/apache/druid/client/coordinator/CoordinatorClient.java index fb23fadf4894..bc2e0a71511c 100644 --- a/server/src/main/java/org/apache/druid/client/coordinator/CoordinatorClient.java +++ b/server/src/main/java/org/apache/druid/client/coordinator/CoordinatorClient.java @@ -22,7 +22,7 @@ import com.google.common.util.concurrent.ListenableFuture; import org.apache.druid.client.BootstrapSegmentsResponse; import org.apache.druid.client.ImmutableSegmentLoadInfo; -import org.apache.druid.java.util.common.parsers.CloseableIterator; +import org.apache.druid.collections.ResourceHolder; import org.apache.druid.query.SegmentDescriptor; import org.apache.druid.query.lookup.LookupExtractorFactoryContainer; import org.apache.druid.rpc.ServiceRetryPolicy; @@ -37,6 +37,7 @@ import javax.annotation.Nullable; import java.net.URI; +import java.util.Iterator; import java.util.List; import java.util.Map; import java.util.Set; @@ -137,12 +138,14 @@ public interface CoordinatorClient /** * Returns an iterator over the metadata segments of multiple datasources in the cluster, fetching them in one go. *

+ * The caller is responsible for closing the holder. + *

* API: {@code GET /druid/coordinator/v1/metadata/segments?includeOvershadowedStatus} * * @param watchedDataSources Optional datasources to filter the segments by. If null or empty, all segments are returned. * @param includeRealtimeSegments If true, includes realtime segments in the result. */ - ListenableFuture> fetchAllUsedSegmentsWithOvershadowedStatus( + ListenableFuture>> fetchAllUsedSegmentsWithOvershadowedStatus( @Nullable Set watchedDataSources, boolean includeRealtimeSegments ); diff --git a/server/src/main/java/org/apache/druid/client/coordinator/CoordinatorClientImpl.java b/server/src/main/java/org/apache/druid/client/coordinator/CoordinatorClientImpl.java index 1fed06decc92..9dfa0bab5b73 100644 --- a/server/src/main/java/org/apache/druid/client/coordinator/CoordinatorClientImpl.java +++ b/server/src/main/java/org/apache/druid/client/coordinator/CoordinatorClientImpl.java @@ -28,10 +28,10 @@ import org.apache.druid.client.BootstrapSegmentsResponse; import org.apache.druid.client.ImmutableSegmentLoadInfo; import org.apache.druid.client.JsonParserIterator; +import org.apache.druid.collections.ResourceHolder; import org.apache.druid.common.guava.FutureUtils; import org.apache.druid.java.util.common.StringUtils; import org.apache.druid.java.util.common.jackson.JacksonUtils; -import org.apache.druid.java.util.common.parsers.CloseableIterator; import org.apache.druid.java.util.http.client.response.BytesFullResponseHandler; import org.apache.druid.java.util.http.client.response.BytesFullResponseHolder; import org.apache.druid.java.util.http.client.response.InputStreamResponseHandler; @@ -52,6 +52,7 @@ import org.apache.druid.server.coordinator.rules.Rule; import org.apache.druid.timeline.DataSegment; import org.apache.druid.timeline.SegmentStatusInCluster; +import org.apache.druid.utils.CloseableUtils; import org.joda.time.Interval; import javax.annotation.Nullable; @@ -59,6 +60,7 @@ import java.net.URISyntaxException; import java.nio.charset.StandardCharsets; import java.util.ArrayList; +import java.util.Iterator; import java.util.List; import java.util.Map; import java.util.Set; @@ -321,7 +323,7 @@ public Map fetchLookupsForTierSync(Stri } @Override - public ListenableFuture> fetchAllUsedSegmentsWithOvershadowedStatus( + public ListenableFuture>> fetchAllUsedSegmentsWithOvershadowedStatus( @Nullable Set watchedDataSources, boolean includeRealtimeSegments ) @@ -346,11 +348,25 @@ public ListenableFuture> fetchAllUsedS new InputStreamResponseHandler() ), inputStream -> { - return new JsonParserIterator<>( + final JsonParserIterator segments = new JsonParserIterator<>( jsonMapper.getTypeFactory().constructType(SegmentStatusInCluster.class), Futures.immediateFuture(inputStream), jsonMapper ); + return new ResourceHolder>() + { + @Override + public Iterator get() + { + return segments; + } + + @Override + public void close() + { + CloseableUtils.closeAndWrapExceptions(() -> CloseableUtils.closeAll(segments, inputStream)); + } + }; } ); } diff --git a/server/src/main/java/org/apache/druid/client/coordinator/NoopCoordinatorClient.java b/server/src/main/java/org/apache/druid/client/coordinator/NoopCoordinatorClient.java index b8aa55ec1e4a..24881e6994ba 100644 --- a/server/src/main/java/org/apache/druid/client/coordinator/NoopCoordinatorClient.java +++ b/server/src/main/java/org/apache/druid/client/coordinator/NoopCoordinatorClient.java @@ -22,7 +22,7 @@ import com.google.common.util.concurrent.ListenableFuture; import org.apache.druid.client.BootstrapSegmentsResponse; import org.apache.druid.client.ImmutableSegmentLoadInfo; -import org.apache.druid.java.util.common.parsers.CloseableIterator; +import org.apache.druid.collections.ResourceHolder; import org.apache.druid.query.SegmentDescriptor; import org.apache.druid.query.lookup.LookupExtractorFactoryContainer; import org.apache.druid.rpc.ServiceRetryPolicy; @@ -37,6 +37,7 @@ import javax.annotation.Nullable; import java.net.URI; +import java.util.Iterator; import java.util.List; import java.util.Map; import java.util.Set; @@ -131,7 +132,7 @@ public Map fetchLookupsForTierSync( } @Override - public ListenableFuture> fetchAllUsedSegmentsWithOvershadowedStatus( + public ListenableFuture>> fetchAllUsedSegmentsWithOvershadowedStatus( @Nullable Set watchedDataSources, boolean includeOvershadowed ) diff --git a/server/src/test/java/org/apache/druid/client/coordinator/CoordinatorClientImplTest.java b/server/src/test/java/org/apache/druid/client/coordinator/CoordinatorClientImplTest.java index 206e5df84462..d7c3c5a8e0bf 100644 --- a/server/src/test/java/org/apache/druid/client/coordinator/CoordinatorClientImplTest.java +++ b/server/src/test/java/org/apache/druid/client/coordinator/CoordinatorClientImplTest.java @@ -35,14 +35,15 @@ import org.apache.druid.client.BootstrapSegmentsResponse; import org.apache.druid.client.DruidServer; import org.apache.druid.client.ImmutableSegmentLoadInfo; +import org.apache.druid.collections.ResourceHolder; import org.apache.druid.common.guava.FutureUtils; import org.apache.druid.guice.StartupInjectorBuilder; import org.apache.druid.initialization.CoreInjectorBuilder; import org.apache.druid.jackson.DefaultObjectMapper; import org.apache.druid.java.util.common.Intervals; import org.apache.druid.java.util.common.StringUtils; -import org.apache.druid.java.util.common.parsers.CloseableIterator; import org.apache.druid.java.util.http.client.response.StringFullResponseHolder; +import org.apache.druid.query.QueryInterruptedException; import org.apache.druid.query.SegmentDescriptor; import org.apache.druid.query.lookup.LookupExtractorFactory; import org.apache.druid.query.lookup.LookupExtractorFactoryContainer; @@ -78,6 +79,7 @@ import java.nio.charset.StandardCharsets; import java.util.ArrayList; import java.util.Collections; +import java.util.Iterator; import java.util.LinkedHashSet; import java.util.List; import java.util.Map; @@ -575,10 +577,10 @@ public void test_fetchAllUsedSegmentsWithOvershadowedStatus_includeRealtime() th jsonMapper.writeValueAsBytes(segments) ); - CloseableIterator iterator = FutureUtils.getUnchecked( + Iterator iterator = FutureUtils.getUnchecked( coordinatorClient.fetchAllUsedSegmentsWithOvershadowedStatus(null, true), true - ); + ).get(); List actualSegments = new ArrayList<>(); while (iterator.hasNext()) { actualSegments.add(iterator.next()); @@ -606,10 +608,10 @@ public void test_fetchAllUsedSegmentsWithOvershadowedStatus_noParams() throws Js jsonMapper.writeValueAsBytes(segments) ); - CloseableIterator iterator = FutureUtils.getUnchecked( + Iterator iterator = FutureUtils.getUnchecked( coordinatorClient.fetchAllUsedSegmentsWithOvershadowedStatus(null, false), true - ); + ).get(); List actualSegments = new ArrayList<>(); while (iterator.hasNext()) { actualSegments.add(iterator.next()); @@ -622,6 +624,26 @@ public void test_fetchAllUsedSegmentsWithOvershadowedStatus_noParams() throws Js ); } + @Test + public void test_fetchAllUsedSegmentsWithOvershadowedStatus_closeBeforeIteratingReleasesTheResponse() + throws JsonProcessingException + { + serviceClient.expectAndRespond( + new RequestBuilder(HttpMethod.GET, "/druid/coordinator/v1/metadata/segments?includeOvershadowedStatus"), + HttpResponseStatus.OK, + ImmutableMap.of(HttpHeaders.CONTENT_TYPE, MediaType.APPLICATION_JSON), + jsonMapper.writeValueAsBytes(ImmutableList.of(SEGMENT1)) + ); + + final ResourceHolder> segments = FutureUtils.getUnchecked( + coordinatorClient.fetchAllUsedSegmentsWithOvershadowedStatus(null, false), + true + ); + segments.close(); + + Assertions.assertThrows(QueryInterruptedException.class, () -> segments.get().hasNext()); + } + @Test public void test_fetchAllUsedSegmentsWithOvershadowedStatus_filterByDataSource() throws Exception { @@ -635,10 +657,10 @@ public void test_fetchAllUsedSegmentsWithOvershadowedStatus_filterByDataSource() jsonMapper.writeValueAsBytes(ImmutableList.of(SEGMENT3)) ); - CloseableIterator iterator = FutureUtils.getUnchecked( + Iterator iterator = FutureUtils.getUnchecked( coordinatorClient.fetchAllUsedSegmentsWithOvershadowedStatus(Set.of("abc"), true), true - ); + ).get(); List actualSegments = new ArrayList<>(); while (iterator.hasNext()) { @@ -666,10 +688,10 @@ public void test_fetchAllUsedSegmentsWithOvershadowedStatus_filterByDataSources( ); Set dataSources = new LinkedHashSet<>(List.of("xyz", "abc")); - CloseableIterator iterator = FutureUtils.getUnchecked( + Iterator iterator = FutureUtils.getUnchecked( coordinatorClient.fetchAllUsedSegmentsWithOvershadowedStatus(dataSources, true), true - ); + ).get(); List actualSegments = new ArrayList<>(); while (iterator.hasNext()) { @@ -696,10 +718,10 @@ public void test_fetchAllUsedSegmentsWithOvershadowedStatus_filterByDataSourceOn jsonMapper.writeValueAsBytes(List.of(SEGMENT3)) ); - CloseableIterator iterator = FutureUtils.getUnchecked( + Iterator iterator = FutureUtils.getUnchecked( coordinatorClient.fetchAllUsedSegmentsWithOvershadowedStatus(Set.of("abc"), false), true - ); + ).get(); List actualSegments = new ArrayList<>(); while (iterator.hasNext()) { diff --git a/sql/src/main/java/org/apache/druid/sql/calcite/planner/QueryHandler.java b/sql/src/main/java/org/apache/druid/sql/calcite/planner/QueryHandler.java index 77a5d01dcaea..0a7ca689a605 100644 --- a/sql/src/main/java/org/apache/druid/sql/calcite/planner/QueryHandler.java +++ b/sql/src/main/java/org/apache/druid/sql/calcite/planner/QueryHandler.java @@ -23,6 +23,7 @@ import com.fasterxml.jackson.databind.ObjectMapper; import com.fasterxml.jackson.databind.node.ArrayNode; import com.fasterxml.jackson.databind.node.ObjectNode; +import com.google.common.annotations.VisibleForTesting; import com.google.common.base.Joiner; import com.google.common.base.Preconditions; import com.google.common.base.Supplier; @@ -59,6 +60,8 @@ import org.apache.druid.error.InvalidSqlInput; import org.apache.druid.jackson.DefaultObjectMapper; import org.apache.druid.java.util.common.guava.BaseSequence; +import org.apache.druid.java.util.common.guava.Sequence; +import org.apache.druid.java.util.common.guava.SequenceWrapper; import org.apache.druid.java.util.common.guava.Sequences; import org.apache.druid.java.util.emitter.EmittingLogger; import org.apache.druid.query.Query; @@ -333,45 +336,48 @@ private PlannerResult planWithBindableConvention() planner.getTypeFactory(), plannerContext.getParameters() ); - final Supplier> resultsSupplier = () -> { - final Enumerable enumerable = theRel.bind(dataContext); - final Enumerator enumerator = enumerable.enumerator(); - return QueryResponse.withEmptyContext( - Sequences.withBaggage(new BaseSequence<>( - new BaseSequence.IteratorMaker>() - { - @Override - public QueryHandler.EnumeratorIterator make() - { - return new QueryHandler.EnumeratorIterator<>(new Iterator<>() - { - @Override - public boolean hasNext() - { - return enumerator.moveNext(); - } - - @Override - public Object[] next() - { - return (Object[]) enumerator.current(); - } - }); - } - - @Override - public void cleanup(QueryHandler.EnumeratorIterator iterFromMake) - { - - } - } - ), enumerator::close) - ); - }; + final Supplier> resultsSupplier = + () -> QueryResponse.withEmptyContext(enumerate(theRel.bind(dataContext))); return new PlannerResult(resultsSupplier, rootQueryRel.validatedRowType); } } + /** + * Rows of a bound {@link BindableRel}; closing the sequence releases everything the binding holds. + */ + @VisibleForTesting + static Sequence enumerate(final Enumerable enumerable) + { + final Sequence rows = new BaseSequence<>( + new BaseSequence.IteratorMaker() + { + @Override + public EnumeratorIterator make() + { + return new EnumeratorIterator(enumerable.enumerator()); + } + + @Override + public void cleanup(EnumeratorIterator iterFromMake) + { + iterFromMake.enumerator.close(); + } + } + ); + + if (!(enumerable instanceof AutoCloseable closeable)) { + return rows; + } + return Sequences.wrap(rows, new SequenceWrapper() + { + @Override + public void after(boolean isDone, Throwable thrown) throws Exception + { + closeable.close(); + } + }); + } + /** * Construct a {@link PlannerResult} for an 'explain' query from a {@link RelNode} and root {@link RelRoot} */ @@ -750,25 +756,25 @@ protected QueryMaker buildQueryMaker(final RelRoot rootQueryRel) throws Validati } } - private static class EnumeratorIterator implements Iterator + private static class EnumeratorIterator implements Iterator { - private final Iterator it; + private final Enumerator enumerator; - EnumeratorIterator(Iterator it) + EnumeratorIterator(Enumerator enumerator) { - this.it = it; + this.enumerator = enumerator; } @Override public boolean hasNext() { - return it.hasNext(); + return enumerator.moveNext(); } @Override - public T next() + public Object[] next() { - return it.next(); + return (Object[]) enumerator.current(); } } } diff --git a/sql/src/main/java/org/apache/druid/sql/calcite/schema/MetadataSegmentView.java b/sql/src/main/java/org/apache/druid/sql/calcite/schema/MetadataSegmentView.java index b292e5e8e04a..e623efedbef9 100644 --- a/sql/src/main/java/org/apache/druid/sql/calcite/schema/MetadataSegmentView.java +++ b/sql/src/main/java/org/apache/druid/sql/calcite/schema/MetadataSegmentView.java @@ -29,9 +29,11 @@ import org.apache.druid.client.BrokerSegmentWatcherConfig; import org.apache.druid.client.DataSegmentInterner; import org.apache.druid.client.coordinator.CoordinatorClient; +import org.apache.druid.collections.ResourceHolder; import org.apache.druid.common.guava.FutureUtils; import org.apache.druid.concurrent.LifecycleLock; import org.apache.druid.guice.ManageLifecycle; +import org.apache.druid.java.util.common.CloseableIterators; import org.apache.druid.java.util.common.ISE; import org.apache.druid.java.util.common.Stopwatch; import org.apache.druid.java.util.common.concurrent.Execs; @@ -47,6 +49,7 @@ import org.apache.druid.timeline.DataSegment; import org.apache.druid.timeline.SegmentId; import org.apache.druid.timeline.SegmentStatusInCluster; +import org.apache.druid.utils.CloseableUtils; import org.checkerframework.checker.nullness.qual.MonotonicNonNull; import org.joda.time.Duration; @@ -154,26 +157,35 @@ private void poll() { log.info("Polling segments from coordinator"); final Stopwatch syncTime = Stopwatch.createStarted(); - final CloseableIterator metadataSegments = fetchSegmentMetadataFromCoordinator(); final ImmutableSortedSet.Builder builder = ImmutableSortedSet.naturalOrder(); - while (metadataSegments.hasNext()) { - final SegmentStatusInCluster segment = metadataSegments.next(); - final DataSegment interned = DataSegmentInterner.intern(segment.getDataSegment()); - Integer replicationFactor = segment.getReplicationFactor(); - if (replicationFactor == null) { - replicationFactor = segmentIdToReplicationFactor.getIfPresent(segment.getDataSegment().getId()); - } else { - segmentIdToReplicationFactor.put(segment.getDataSegment().getId(), segment.getReplicationFactor()); + final ResourceHolder> metadataSegments = fetchSegmentMetadataFromCoordinator(); + try { + final Iterator segments = metadataSegments.get(); + while (segments.hasNext()) { + final SegmentStatusInCluster segment = segments.next(); + final DataSegment interned = DataSegmentInterner.intern(segment.getDataSegment()); + Integer replicationFactor = segment.getReplicationFactor(); + if (replicationFactor == null) { + replicationFactor = segmentIdToReplicationFactor.getIfPresent(segment.getDataSegment().getId()); + } else { + segmentIdToReplicationFactor.put(segment.getDataSegment().getId(), segment.getReplicationFactor()); + } + final SegmentStatusInCluster segmentStatusInCluster = new SegmentStatusInCluster( + interned, + segment.isOvershadowed(), + replicationFactor, + segment.getNumRows(), + segment.isRealtime() + ); + builder.add(segmentStatusInCluster); } - final SegmentStatusInCluster segmentStatusInCluster = new SegmentStatusInCluster( - interned, - segment.isOvershadowed(), - replicationFactor, - segment.getNumRows(), - segment.isRealtime() + } + finally { + CloseableUtils.closeAndSuppressExceptions( + metadataSegments, + e -> log.warn(e, "Failed to close the segment metadata response from the Coordinator") ); - builder.add(segmentStatusInCluster); } publishedSegments = builder.build(); cachePopulated.countDown(); @@ -187,7 +199,7 @@ private void poll() * {@link BrokerSegmentMetadataCacheConfig#isMetadataSegmentCacheEnable()}) * OR by querying the Coordinator on the fly. */ - Iterator getSegments() + CloseableIterator getSegments() { return getSegments(null); } @@ -196,23 +208,27 @@ Iterator getSegments() * Returns published (and, with centralized schema, realtime) segment metadata, optionally * restricted to {@code dataSources}. */ - Iterator getSegments(@Nullable Set dataSources) + CloseableIterator getSegments(@Nullable Set dataSources) { - final Iterator base; + final CloseableIterator base; if (isCacheEnabled) { Uninterruptibles.awaitUninterruptibly(cachePopulated); - base = publishedSegments.iterator(); + base = CloseableIterators.withEmptyBaggage(publishedSegments.iterator()); } else { // Cache disabled: the Coordinator returns all used segments; filter client-side to preserve semantics. - base = fetchSegmentMetadataFromCoordinator(); + final ResourceHolder> fetched = fetchSegmentMetadataFromCoordinator(); + base = CloseableIterators.wrap(fetched.get(), fetched); } return dataSources == null ? base - : Iterators.filter(base, s -> dataSources.contains(s.getDataSegment().getDataSource())); + : CloseableIterators.wrap( + Iterators.filter(base, s -> dataSources.contains(s.getDataSegment().getDataSource())), + base + ); } // Note that coordinator must be up to get segments - private CloseableIterator fetchSegmentMetadataFromCoordinator() + private ResourceHolder> fetchSegmentMetadataFromCoordinator() { // includeRealtimeSegments flag would additionally request realtime segments // note that realtime segments are returned only when druid.centralizedDatasourceSchema.enabled is set on the Coordinator diff --git a/sql/src/main/java/org/apache/druid/sql/calcite/schema/SystemSchema.java b/sql/src/main/java/org/apache/druid/sql/calcite/schema/SystemSchema.java index 885fc1ca3522..09e175d24151 100644 --- a/sql/src/main/java/org/apache/druid/sql/calcite/schema/SystemSchema.java +++ b/sql/src/main/java/org/apache/druid/sql/calcite/schema/SystemSchema.java @@ -55,6 +55,7 @@ import org.apache.druid.indexing.overlord.supervisor.SupervisorStatus; import org.apache.druid.java.util.common.ISE; import org.apache.druid.java.util.common.StringUtils; +import org.apache.druid.java.util.common.io.Closer; import org.apache.druid.java.util.common.jackson.JacksonUtils; import org.apache.druid.java.util.common.parsers.CloseableIterator; import org.apache.druid.java.util.http.client.HttpClient; @@ -441,10 +442,11 @@ public Enumerable scan( : new HashSet<>(); // Get segments from metadata segment cache (if enabled in SQL planner config), else directly from - // Coordinator. This may include both published and realtime segments. - final Iterator metadataStoreSegments = metadataView.getSegments(dataSourceFilter); + // Coordinator. This may include both published and realtime segments. Fetched on first read, so a scan that is + // never read holds no Coordinator response. + final Closer closer = Closer.create(); final FluentIterable publishedSegments = FluentIterable - .from(() -> getAuthorizedPublishedSegments(metadataStoreSegments)) + .from(() -> getAuthorizedPublishedSegments(closer.register(metadataView.getSegments(dataSourceFilter)))) .transform(val -> { final DataSegment segment = val.getDataSegment(); final AvailableSegmentMetadata availableSegmentMetadata = @@ -549,7 +551,7 @@ public Enumerable scan( Iterables.concat(publishedSegments, availableSegments) ); - return Linq4j.asEnumerable(allSegments) + return Linq4j.asEnumerable(() -> wrap(allSegments.iterator(), closer)) .where(Objects::nonNull) .select(row -> projectSegmentsRow(row, projects, jsonMapper)); } diff --git a/sql/src/test/java/org/apache/druid/sql/calcite/planner/QueryHandlerTest.java b/sql/src/test/java/org/apache/druid/sql/calcite/planner/QueryHandlerTest.java new file mode 100644 index 000000000000..93d5f35c68b3 --- /dev/null +++ b/sql/src/test/java/org/apache/druid/sql/calcite/planner/QueryHandlerTest.java @@ -0,0 +1,57 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ + +package org.apache.druid.sql.calcite.planner; + +import org.apache.calcite.interpreter.Interpreter; +import org.apache.calcite.linq4j.Enumerator; +import org.apache.calcite.linq4j.Linq4j; +import org.apache.druid.java.util.common.ISE; +import org.junit.jupiter.api.Assertions; +import org.junit.jupiter.api.Test; +import org.mockito.Mockito; + +import java.util.List; + +public class QueryHandlerTest +{ + private final Interpreter interpreter = Mockito.mock(Interpreter.class); + + @Test + public void testEnumerateClosesTheEnumeratorAndTheInterpreterOnceTheResultsAreConsumed() + { + final Enumerator enumerator = Mockito.spy(Linq4j.enumerator(List.of(new Object[]{1L}))); + Mockito.when(interpreter.enumerator()).thenReturn(enumerator); + + Assertions.assertEquals(1, QueryHandler.enumerate(interpreter).toList().size()); + + Mockito.verify(enumerator).close(); + Mockito.verify(interpreter).close(); + } + + @Test + public void testEnumerateClosesTheInterpreterIfEnumerationFails() + { + Mockito.when(interpreter.enumerator()).thenThrow(new ISE("Interpreter node failed")); + + Assertions.assertThrows(ISE.class, () -> QueryHandler.enumerate(interpreter).toList()); + + Mockito.verify(interpreter).close(); + } +} diff --git a/sql/src/test/java/org/apache/druid/sql/calcite/schema/MetadataSegmentViewTest.java b/sql/src/test/java/org/apache/druid/sql/calcite/schema/MetadataSegmentViewTest.java index 08b8d111b961..b84ef6818fea 100644 --- a/sql/src/test/java/org/apache/druid/sql/calcite/schema/MetadataSegmentViewTest.java +++ b/sql/src/test/java/org/apache/druid/sql/calcite/schema/MetadataSegmentViewTest.java @@ -23,7 +23,9 @@ import com.google.common.util.concurrent.Futures; import org.apache.druid.client.BrokerSegmentWatcherConfig; import org.apache.druid.client.coordinator.CoordinatorClient; -import org.apache.druid.java.util.common.CloseableIterators; +import org.apache.druid.collections.ResourceHolder; +import org.apache.druid.collections.StupidResourceHolder; +import org.apache.druid.java.util.common.ISE; import org.apache.druid.java.util.metrics.NoopTaskHolder; import org.apache.druid.metadata.segment.cache.Metric; import org.apache.druid.segment.TestDataSource; @@ -31,14 +33,23 @@ import org.apache.druid.server.metrics.LatchableEmitter; import org.apache.druid.server.metrics.LatchableEmitterConfig; import org.apache.druid.timeline.SegmentStatusInCluster; +import org.apache.druid.utils.CloseableUtils; import org.junit.jupiter.api.Assertions; import org.junit.jupiter.api.BeforeEach; import org.junit.jupiter.api.Test; import org.mockito.ArgumentMatchers; import org.mockito.Mockito; +import java.io.Closeable; +import java.io.IOException; import java.util.ArrayList; +import java.util.Iterator; import java.util.List; +import java.util.NoSuchElementException; +import java.util.Set; +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicBoolean; public class MetadataSegmentViewTest { @@ -84,7 +95,7 @@ public void test_start_triggersSegmentPollFromCoordinator_ifCacheIsEnabled() ArgumentMatchers.eq(true) ) ).thenReturn( - Futures.immediateFuture(CloseableIterators.withEmptyBaggage(expectedSegments.iterator())) + Futures.immediateFuture(StupidResourceHolder.create(expectedSegments.iterator())) ); // Start the test target and wait for it to sync with the Coordinator @@ -96,4 +107,124 @@ public void test_start_triggersSegmentPollFromCoordinator_ifCacheIsEnabled() Assertions.assertEquals(10, observedSegments.size()); Assertions.assertEquals(expectedSegments, observedSegments); } + + @Test + public void test_poll_releasesTheCoordinatorResponse() + { + final AtomicBoolean closed = new AtomicBoolean(false); + + Mockito.when( + coordinatorClient.fetchAllUsedSegmentsWithOvershadowedStatus( + ArgumentMatchers.eq(null), + ArgumentMatchers.eq(true) + ) + ).thenAnswer( + invocation -> Futures.immediateFuture( + holder(List.of().iterator(), () -> closed.set(true)) + ) + ); + + segmentView.start(); + emitter.waitForEvent(event -> event.hasMetricName(Metric.SYNC_DURATION_MILLIS)); + + Assertions.assertTrue(closed.get(), "poll must close the streamed response"); + } + + @Test + public void test_poll_releasesTheCoordinatorResponse_whenIterationFails() throws Exception + { + final CountDownLatch closed = new CountDownLatch(1); + + final Iterator failsMidStream = new Iterator<>() + { + @Override + public boolean hasNext() + { + throw new ISE("Coordinator response ended mid-stream"); + } + + @Override + public SegmentStatusInCluster next() + { + throw new NoSuchElementException(); + } + }; + + Mockito.when( + coordinatorClient.fetchAllUsedSegmentsWithOvershadowedStatus( + ArgumentMatchers.eq(null), + ArgumentMatchers.eq(true) + ) + ).thenAnswer( + invocation -> Futures.immediateFuture(holder(failsMidStream, closed::countDown)) + ); + + segmentView.start(); + + try { + Assertions.assertTrue( + closed.await(30, TimeUnit.SECONDS), + "a failed poll must still release the streamed response" + ); + } + finally { + segmentView.stop(); + } + } + + @Test + public void test_getSegments_releasesTheCoordinatorResponse_ifCacheIsDisabledAndFiltered() throws IOException + { + final AtomicBoolean closed = new AtomicBoolean(false); + + Mockito.when( + coordinatorClient.fetchAllUsedSegmentsWithOvershadowedStatus( + ArgumentMatchers.eq(null), + ArgumentMatchers.eq(true) + ) + ).thenReturn( + Futures.immediateFuture( + holder(List.of().iterator(), () -> closed.set(true)) + ) + ); + + final MetadataSegmentView uncachedView = new MetadataSegmentView( + coordinatorClient, + new BrokerSegmentWatcherConfig(), + new BrokerSegmentMetadataCacheConfig() + { + @Override + public boolean isMetadataSegmentCacheEnable() + { + return false; + } + }, + emitter + ); + + uncachedView.getSegments(Set.of(TestDataSource.WIKI)).close(); + + Assertions.assertTrue(closed.get()); + } + + private static ResourceHolder> holder( + final Iterator segments, + final Closeable response + ) + { + return new ResourceHolder<>() + { + @Override + public Iterator get() + { + return segments; + } + + @Override + public void close() + { + CloseableUtils.closeAndWrapExceptions(response); + } + }; + } } diff --git a/sql/src/test/java/org/apache/druid/sql/calcite/schema/SystemSchemaTest.java b/sql/src/test/java/org/apache/druid/sql/calcite/schema/SystemSchemaTest.java index ea3e981a8d90..28945c6c066b 100644 --- a/sql/src/test/java/org/apache/druid/sql/calcite/schema/SystemSchemaTest.java +++ b/sql/src/test/java/org/apache/druid/sql/calcite/schema/SystemSchemaTest.java @@ -35,6 +35,7 @@ import org.apache.calcite.DataContext; import org.apache.calcite.adapter.java.JavaTypeFactory; import org.apache.calcite.jdbc.JavaTypeFactoryImpl; +import org.apache.calcite.linq4j.Enumerator; import org.apache.calcite.linq4j.QueryProvider; import org.apache.calcite.rel.type.RelDataType; import org.apache.calcite.rel.type.RelDataTypeField; @@ -153,6 +154,7 @@ import java.util.List; import java.util.Map; import java.util.Set; +import java.util.concurrent.atomic.AtomicBoolean; public class SystemSchemaTest extends CalciteTestBase { @@ -707,7 +709,9 @@ public void testSegmentsTable() throws Exception new SegmentStatusInCluster(segment2, false, 0, null, false) )); - EasyMock.expect(metadataView.getSegments(EasyMock.anyObject())).andReturn(publishedSegments.iterator()).once(); + EasyMock.expect(metadataView.getSegments(EasyMock.anyObject())) + .andReturn(CloseableIterators.withEmptyBaggage(publishedSegments.iterator())) + .once(); EasyMock.replay(request, responseHolder, responseHandler, metadataView); DataContext dataContext = createDataContext(); @@ -828,7 +832,9 @@ public void testSegmentsTableWithProjection() throws JsonProcessingException new SegmentStatusInCluster(segment2, false, 0, null, false) )); - EasyMock.expect(metadataView.getSegments(EasyMock.anyObject())).andReturn(publishedSegments.iterator()).once(); + EasyMock.expect(metadataView.getSegments(EasyMock.anyObject())) + .andReturn(CloseableIterators.withEmptyBaggage(publishedSegments.iterator())) + .once(); EasyMock.replay(request, responseHolder, responseHandler, metadataView); DataContext dataContext = createDataContext(); @@ -884,6 +890,42 @@ public void testSegmentsTableWithProjection() throws JsonProcessingException ); } + @Test + public void testSegmentsTableReleasesTheCoordinatorResponseWhenClosedEarly() + { + final SegmentsTable segmentsTable = + new SegmentsTable(segmentMetadataCache, metadataView, MAPPER, authMapper, createAuthResult(Users.SUPER)); + final List publishedSegments = List.of( + new SegmentStatusInCluster(segment1, true, 2, null, false), + new SegmentStatusInCluster(segment2, false, 0, null, false) + ); + final AtomicBoolean closed = new AtomicBoolean(false); + + EasyMock.expect(metadataView.getSegments(EasyMock.anyObject())) + .andReturn(CloseableIterators.wrap(publishedSegments.iterator(), () -> closed.set(true))) + .once(); + + EasyMock.replay(request, responseHolder, responseHandler, metadataView); + final Enumerator rows = + segmentsTable.scan(createDataContext(), Collections.emptyList(), null).enumerator(); + Assertions.assertTrue(rows.moveNext()); + rows.close(); + + Assertions.assertTrue(closed.get()); + } + + @Test + public void testSegmentsTableDefersTheCoordinatorFetchUntilRowsAreRead() + { + final SegmentsTable segmentsTable = + new SegmentsTable(segmentMetadataCache, metadataView, MAPPER, authMapper, createAuthResult(Users.SUPER)); + + EasyMock.replay(request, responseHolder, responseHandler, metadataView); + segmentsTable.scan(createDataContext(), Collections.emptyList(), null).enumerator().close(); + + EasyMock.verify(metadataView); + } + @Test public void testServersTable() throws URISyntaxException {