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
1 change: 1 addition & 0 deletions docs/configuration/index.md
Original file line number Diff line number Diff line change
Expand Up @@ -706,6 +706,7 @@ All Druid components can communicate with each other over HTTP.
|`druid.global.http.connectTimeout`|Connect timeout for the HTTP client used for most direct RPC between Druid services. This covers, among other things, Overlord-to-task and supervisor-to-task calls in the indexing service, Coordinator lookup management, dynamic config sync between services, MSQ tasks reading from data servers, and general Coordinator/Overlord/Broker service clients. Does not affect Broker-to-Historical query dispatch (see `druid.broker.http.connectTimeout`) or request forwarding (see `clientConnectTimeout`).|`PT10S`|
|`druid.global.http.allocator`|Netty memory allocator used by the direct-RPC HTTP client. Accepts `adaptive` (adaptive between `pooled` and `unpooled` based on load), `pooled`, or `unpooled`.|`adaptive`|
|`druid.global.http.poolImplementation`|How the connection pool tracks demand, never exceeding `numConnections` either way. With `adaptive`, a request discards every stale or broken connection it walks past and opens a new one only once none is left, so the pool falls back to the number of connections the traffic actually needs. With `retaining`, the pool holds on to every connection it has opened, replacing a stale or broken one by a fresh one, one for one, so it stays at its high-water mark.|`adaptive`|
|`druid.client.coordinator.maxAttempts`|Maximum number of attempts, including the first one, that a service makes for a request to the Coordinator before giving up on a retryable error. Must be at least 1.|`15`|

### Common endpoints configuration

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 @@ -240,6 +240,7 @@ public <Intermediate, Final> ListenableFuture<Final> go(
private final Object watermarkLock = new Object();
private long suspendWatermark = -1;
private long resumeWatermark = -1;
private boolean returnedToPool = false;

@Override
protected void channelRead0(ChannelHandlerContext ctx, HttpObject msg)
Expand Down Expand Up @@ -274,6 +275,9 @@ protected void channelRead0(ChannelHandlerContext ctx, HttpObject msg)
public long resume(long resumeChunkNum)
{
synchronized (watermarkLock) {
if (returnedToPool) {
return 0;
}
resumeWatermark = Math.max(resumeWatermark, resumeChunkNum);

if (suspendWatermark >= 0 && resumeWatermark >= suspendWatermark) {
Expand All @@ -292,8 +296,12 @@ public long resume(long resumeChunkNum)
@Override
public void abort()
{
log.debug("[%s] Aborted connection at caller's request.", requestDesc);
channel.close();
// On the event loop, so that it is ordered against finishRequest handing the channel back.
if (channel.eventLoop().inEventLoop()) {
closeUnlessReturnedToPool();
} else {
channel.eventLoop().execute(() -> closeUnlessReturnedToPool());
}
}
};

Expand Down Expand Up @@ -388,9 +396,22 @@ private void finishRequest()
}
removeHandlers();
channel.config().setAutoRead(true);
synchronized (watermarkLock) {
returnedToPool = true;
}
channelResourceContainer.returnResource();
}

private void closeUnlessReturnedToPool()
{
synchronized (watermarkLock) {
if (!returnedToPool) {
log.debug("[%s] Aborted connection at caller's request.", requestDesc);
channel.close();
}
}
}

@Override
public void exceptionCaught(ChannelHandlerContext context, Throwable cause)
{
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -37,7 +37,7 @@ public class AppendableByteArrayInputStream extends InputStream

private volatile boolean done = false;
private volatile Throwable throwable;
private volatile int available = 0;
private volatile long available = 0;

private byte[] curr = new byte[]{};
private int currIndex = 0;
Expand All @@ -49,6 +49,9 @@ public void add(byte[] bytesToAdd)
}

synchronized (singleByteReaderDoer) {
if (done) {
return;
}
bytes.addLast(bytesToAdd);
available += bytesToAdd.length;
singleByteReaderDoer.notify();
Expand All @@ -68,10 +71,18 @@ public void exceptionCaught(Throwable t)
synchronized (singleByteReaderDoer) {
done = true;
throwable = t;
bytes.clear();
available = 0;
singleByteReaderDoer.notifyAll();
}
}

@Override
public void close()
{
exceptionCaught(new IOException("Stream closed"));
}

@Override
public int read() throws IOException
{
Expand Down Expand Up @@ -145,7 +156,7 @@ private long scanThroughBytesAndDoSomething(long numToScan, Doer doer) throws IO
break;
}
try {
available -= numPulled;
releaseUnreadBytes(numPulled);
numPulled = 0;
singleByteReaderDoer.wait();
}
Expand Down Expand Up @@ -181,16 +192,24 @@ private long scanThroughBytesAndDoSomething(long numToScan, Doer doer) throws IO
}

synchronized (singleByteReaderDoer) {
available -= numPulled;
releaseUnreadBytes(numPulled);
}

return numScanned;
}

private void releaseUnreadBytes(long numPulled)
{
// exceptionCaught zeroes the count, bytes pulled concurrently included.
if (throwable == null) {
available -= numPulled;
}
}

@Override
public int available()
{
return available;
return (int) Math.min(available, Integer.MAX_VALUE);
}

private interface Doer
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -87,6 +87,10 @@ ClientResponse<IntermediateType> handleChunk(

void exceptionCaught(ClientResponse<IntermediateType> clientResponse, Throwable e);

/**
* Flow control over the connection carrying one response. Once that response has completed, calls have no effect:
* the connection may already be carrying the response of another request.
*/
interface TrafficCop
{
/**
Expand Down
Loading
Loading