Skip to content
Open
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 @@ -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;
Expand Down Expand Up @@ -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;
Expand Down Expand Up @@ -180,12 +181,12 @@ public void setup()
final NoopCoordinatorClient coordinatorClient = new NoopCoordinatorClient()
{
@Override
public ListenableFuture<CloseableIterator<SegmentStatusInCluster>> fetchAllUsedSegmentsWithOvershadowedStatus(
public ListenableFuture<ResourceHolder<Iterator<SegmentStatusInCluster>>> fetchAllUsedSegmentsWithOvershadowedStatus(
Set<String> watchedDataSources,
boolean includeRealtimeSegments
)
{
return Futures.immediateFuture(CloseableIterators.withEmptyBaggage(publishedSegments.iterator()));
return Futures.immediateFuture(StupidResourceHolder.create(publishedSegments.iterator()));
}
};

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand All @@ -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;
Expand Down Expand Up @@ -154,13 +154,14 @@ public void test_fetchUsedSegments()

@Test
@Timeout(20)
public void test_fetchAllUsedSegmentsWithOvershadowedStatus() throws IOException
public void test_fetchAllUsedSegmentsWithOvershadowedStatus()
{
runIndexTask();

try (CloseableIterator<SegmentStatusInCluster> iterator = cluster.callApi().onLeaderCoordinator(
try (ResourceHolder<Iterator<SegmentStatusInCluster>> segments = cluster.callApi().onLeaderCoordinator(
c -> c.fetchAllUsedSegmentsWithOvershadowedStatus(Set.of(dataSource), true))
) {
final Iterator<SegmentStatusInCluster> iterator = segments.get();
Assertions.assertTrue(iterator.hasNext());
SegmentStatusInCluster segmentStatus = iterator.next();
Assertions.assertEquals(dataSource, segmentStatus.getDataSegment().getDataSource());
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand All @@ -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;
Expand Down Expand Up @@ -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.
* <p>
* The caller is responsible for closing the holder.
* <p>
* 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<CloseableIterator<SegmentStatusInCluster>> fetchAllUsedSegmentsWithOvershadowedStatus(
ListenableFuture<ResourceHolder<Iterator<SegmentStatusInCluster>>> fetchAllUsedSegmentsWithOvershadowedStatus(
@Nullable Set<String> watchedDataSources,
boolean includeRealtimeSegments
);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand All @@ -52,13 +52,15 @@
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;
import java.net.URI;
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;
Expand Down Expand Up @@ -321,7 +323,7 @@ public Map<String, LookupExtractorFactoryContainer> fetchLookupsForTierSync(Stri
}

@Override
public ListenableFuture<CloseableIterator<SegmentStatusInCluster>> fetchAllUsedSegmentsWithOvershadowedStatus(
public ListenableFuture<ResourceHolder<Iterator<SegmentStatusInCluster>>> fetchAllUsedSegmentsWithOvershadowedStatus(
@Nullable Set<String> watchedDataSources,
boolean includeRealtimeSegments
)
Expand All @@ -346,11 +348,25 @@ public ListenableFuture<CloseableIterator<SegmentStatusInCluster>> fetchAllUsedS
new InputStreamResponseHandler()
),
inputStream -> {
return new JsonParserIterator<>(
final JsonParserIterator<SegmentStatusInCluster> segments = new JsonParserIterator<>(
jsonMapper.getTypeFactory().constructType(SegmentStatusInCluster.class),
Futures.immediateFuture(inputStream),
jsonMapper
);
return new ResourceHolder<Iterator<SegmentStatusInCluster>>()
{
@Override
public Iterator<SegmentStatusInCluster> get()
{
return segments;
}

@Override
public void close()
{
CloseableUtils.closeAndWrapExceptions(() -> CloseableUtils.closeAll(segments, inputStream));
}
};
}
);
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand All @@ -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;
Expand Down Expand Up @@ -131,7 +132,7 @@ public Map<String, LookupExtractorFactoryContainer> fetchLookupsForTierSync(
}

@Override
public ListenableFuture<CloseableIterator<SegmentStatusInCluster>> fetchAllUsedSegmentsWithOvershadowedStatus(
public ListenableFuture<ResourceHolder<Iterator<SegmentStatusInCluster>>> fetchAllUsedSegmentsWithOvershadowedStatus(
@Nullable Set<String> watchedDataSources,
boolean includeOvershadowed
)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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;
Expand Down Expand Up @@ -575,10 +577,10 @@ public void test_fetchAllUsedSegmentsWithOvershadowedStatus_includeRealtime() th
jsonMapper.writeValueAsBytes(segments)
);

CloseableIterator<SegmentStatusInCluster> iterator = FutureUtils.getUnchecked(
Iterator<SegmentStatusInCluster> iterator = FutureUtils.getUnchecked(
coordinatorClient.fetchAllUsedSegmentsWithOvershadowedStatus(null, true),
true
);
).get();
List<SegmentStatusInCluster> actualSegments = new ArrayList<>();
while (iterator.hasNext()) {
actualSegments.add(iterator.next());
Expand Down Expand Up @@ -606,10 +608,10 @@ public void test_fetchAllUsedSegmentsWithOvershadowedStatus_noParams() throws Js
jsonMapper.writeValueAsBytes(segments)
);

CloseableIterator<SegmentStatusInCluster> iterator = FutureUtils.getUnchecked(
Iterator<SegmentStatusInCluster> iterator = FutureUtils.getUnchecked(
coordinatorClient.fetchAllUsedSegmentsWithOvershadowedStatus(null, false),
true
);
).get();
List<SegmentStatusInCluster> actualSegments = new ArrayList<>();
while (iterator.hasNext()) {
actualSegments.add(iterator.next());
Expand All @@ -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<Iterator<SegmentStatusInCluster>> 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
{
Expand All @@ -635,10 +657,10 @@ public void test_fetchAllUsedSegmentsWithOvershadowedStatus_filterByDataSource()
jsonMapper.writeValueAsBytes(ImmutableList.of(SEGMENT3))
);

CloseableIterator<SegmentStatusInCluster> iterator = FutureUtils.getUnchecked(
Iterator<SegmentStatusInCluster> iterator = FutureUtils.getUnchecked(
coordinatorClient.fetchAllUsedSegmentsWithOvershadowedStatus(Set.of("abc"), true),
true
);
).get();

List<SegmentStatusInCluster> actualSegments = new ArrayList<>();
while (iterator.hasNext()) {
Expand Down Expand Up @@ -666,10 +688,10 @@ public void test_fetchAllUsedSegmentsWithOvershadowedStatus_filterByDataSources(
);

Set<String> dataSources = new LinkedHashSet<>(List.of("xyz", "abc"));
CloseableIterator<SegmentStatusInCluster> iterator = FutureUtils.getUnchecked(
Iterator<SegmentStatusInCluster> iterator = FutureUtils.getUnchecked(
coordinatorClient.fetchAllUsedSegmentsWithOvershadowedStatus(dataSources, true),
true
);
).get();

List<SegmentStatusInCluster> actualSegments = new ArrayList<>();
while (iterator.hasNext()) {
Expand All @@ -696,10 +718,10 @@ public void test_fetchAllUsedSegmentsWithOvershadowedStatus_filterByDataSourceOn
jsonMapper.writeValueAsBytes(List.of(SEGMENT3))
);

CloseableIterator<SegmentStatusInCluster> iterator = FutureUtils.getUnchecked(
Iterator<SegmentStatusInCluster> iterator = FutureUtils.getUnchecked(
coordinatorClient.fetchAllUsedSegmentsWithOvershadowedStatus(Set.of("abc"), false),
true
);
).get();

List<SegmentStatusInCluster> actualSegments = new ArrayList<>();
while (iterator.hasNext()) {
Expand Down
Loading
Loading