Skip to content

Commit fd444b0

Browse files
committed
Added more tests and removed dead code halding 503 in client
1 parent 54cc9c6 commit fd444b0

6 files changed

Lines changed: 468 additions & 39 deletions

File tree

client-v2/src/main/java/com/clickhouse/client/api/Client.java

Lines changed: 8 additions & 31 deletions
Original file line numberDiff line numberDiff line change
@@ -47,7 +47,6 @@
4747
import com.google.common.collect.ImmutableList;
4848
import net.jpountz.lz4.LZ4Factory;
4949
import org.apache.hc.core5.concurrent.DefaultThreadFactory;
50-
import org.apache.hc.core5.http.HttpStatus;
5150
import org.slf4j.Logger;
5251
import org.slf4j.LoggerFactory;
5352

@@ -1398,18 +1397,11 @@ public CompletableFuture<InsertResponse> insert(String tableName, List<?> data,
13981397
}
13991398
out.close();
14001399
});
1401-
try (TransportResponse httpResponse = httpClientHelper.executeRequest(transportRequest)) {
1402-
// Check response
1403-
if (httpResponse.getStatusCode() == HttpStatus.SC_SERVICE_UNAVAILABLE) {
1404-
LOG.warn("Failed to get response. Server returned {}. Retrying. (Duration: {})", httpResponse.getStatusCode(), durationSince(startTime));
1405-
selectedEndpoint = getNextAliveNode();
1406-
continue;
1407-
}
1408-
1400+
try (TransportResponse transportResponse = httpClientHelper.executeRequest(transportRequest)) {
14091401
ClientStatisticsHolder clientStats = globalClientStats.remove(operationId);
1410-
OperationMetrics metrics = completeOperation(httpResponse, clientStats, requestSettings.getQueryId());
1402+
OperationMetrics metrics = completeOperation(transportResponse, clientStats, requestSettings.getQueryId());
14111403

1412-
return new InsertResponse(httpResponse, metrics);
1404+
return new InsertResponse(transportResponse, metrics);
14131405
} catch (Exception e) {
14141406
String msg = requestExMsg("Insert", (i + 1), durationSince(startTime).toMillis(), requestSettings.getQueryId());
14151407
lastException = httpClientHelper.wrapException(msg, e, requestSettings.getQueryId());
@@ -1424,7 +1416,8 @@ public CompletableFuture<InsertResponse> insert(String tableName, List<?> data,
14241416

14251417
String errMsg = requestExMsg("Insert", retries, durationSince(startTime).toMillis(), requestSettings.getQueryId());
14261418
LOG.warn(errMsg);
1427-
throw (lastException == null ? new ClientException(errMsg) : lastException); };
1419+
throw (lastException == null ? new ClientException(errMsg) : lastException);
1420+
};
14281421

14291422
return runAsyncOperation(supplier, requestSettings.getAllSettings());
14301423
}
@@ -1601,17 +1594,9 @@ public CompletableFuture<InsertResponse> insert(String tableName,
16011594
out.close();
16021595
});
16031596

1604-
try (TransportResponse httpResponse = httpClientHelper.executeRequest(transportRequest)) {
1605-
1606-
// Check response
1607-
if (httpResponse.getStatusCode() == HttpStatus.SC_SERVICE_UNAVAILABLE) {
1608-
LOG.warn("Failed to get response. Server returned {}. Retrying. (Duration: {})", httpResponse.getStatusCode(), durationSince(startTime));
1609-
selectedEndpoint = getNextAliveNode();
1610-
continue;
1611-
}
1612-
1613-
OperationMetrics metrics = completeOperation(httpResponse, finalClientStats, requestSettings.getQueryId());
1614-
return new InsertResponse(httpResponse, metrics);
1597+
try (TransportResponse transportResponse = httpClientHelper.executeRequest(transportRequest)) {
1598+
OperationMetrics metrics = completeOperation(transportResponse, finalClientStats, requestSettings.getQueryId());
1599+
return new InsertResponse(transportResponse, metrics);
16151600
} catch (Exception e) {
16161601
String msg = requestExMsg("Insert", (i + 1), durationSince(startTime).toMillis(), requestSettings.getQueryId());
16171602
lastException = httpClientHelper.wrapException(msg, e, requestSettings.getQueryId());
@@ -1722,14 +1707,6 @@ public CompletableFuture<QueryResponse> query(String sqlQuery, Map<String, Objec
17221707
TransportResponse transportResp = null;
17231708
try {
17241709
transportResp = httpClientHelper.executeRequest(request);
1725-
// Check response
1726-
if (transportResp.getStatusCode() == HttpStatus.SC_SERVICE_UNAVAILABLE) {
1727-
LOG.warn("Failed to get response. Server returned {}. Retrying. (Duration: {})", transportResp.getStatusCode(), durationSince(startTime));
1728-
selectedEndpoint = getNextAliveNode();
1729-
ClientUtils.quiteClose(transportResp, LOG);
1730-
continue;
1731-
}
1732-
17331710
OperationMetrics metrics = completeOperation(transportResp, clientStats, requestSettings.getQueryId());
17341711
ClickHouseFormat responseFormat = transportResp.getDataFormat();
17351712
if (responseFormat == null) {

client-v2/src/main/java/com/clickhouse/client/api/query/QueryResponse.java

Lines changed: 6 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -3,15 +3,14 @@
33
import com.clickhouse.client.api.ClientConfigProperties;
44
import com.clickhouse.client.api.ClientException;
55
import com.clickhouse.client.api.http.ClickHouseHttpProto;
6-
import com.clickhouse.client.api.internal.ClientUtils;
76
import com.clickhouse.client.api.metrics.OperationMetrics;
87
import com.clickhouse.client.api.metrics.ServerMetrics;
98
import com.clickhouse.client.api.transport.internal.TransportResponse;
109
import com.clickhouse.data.ClickHouseFormat;
1110

1211
import java.io.InputStream;
13-
import java.time.ZoneId;
1412
import java.util.Map;
13+
import java.util.Objects;
1514
import java.util.TimeZone;
1615

1716
/**
@@ -40,6 +39,7 @@ public class QueryResponse implements AutoCloseable {
4039
private final Map<String, String> responseHeaders;
4140

4241
public QueryResponse(TransportResponse response, ClickHouseFormat format, QuerySettings settings, OperationMetrics operationMetrics) {
42+
Objects.requireNonNull(response, "response is null");
4343
this.transportResponse = response;
4444
this.format = format;
4545
this.operationMetrics = operationMetrics;
@@ -65,12 +65,10 @@ public InputStream getInputStream() {
6565

6666
@Override
6767
public void close() throws Exception {
68-
if (transportResponse != null ) {
69-
try {
70-
transportResponse.close();
71-
} catch (Exception e) {
72-
throw new ClientException("Failed to close response", e);
73-
}
68+
try {
69+
transportResponse.close();
70+
} catch (Exception e) {
71+
throw new ClientException("Failed to close response", e);
7472
}
7573
}
7674

Lines changed: 48 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,48 @@
1+
package com.clickhouse.client.api.internal;
2+
3+
import org.mockito.Mockito;
4+
import org.slf4j.Logger;
5+
import org.testng.Assert;
6+
import org.testng.annotations.Test;
7+
8+
import java.io.Closeable;
9+
import java.io.IOException;
10+
11+
public class ClientUtilsTest {
12+
13+
@Test(groups = {"unit"})
14+
public void testQuiteCloseSwallowsExceptionAndLogs() throws IOException {
15+
Logger log = Mockito.mock(Logger.class);
16+
IOException failure = new IOException("close failed");
17+
Closeable closeable = Mockito.mock(Closeable.class);
18+
Mockito.doThrow(failure).when(closeable).close();
19+
20+
// Should not propagate the exception thrown by close()
21+
ClientUtils.quiteClose(closeable, log);
22+
23+
Mockito.verify(closeable).close();
24+
Mockito.verify(log).warn(Mockito.contains("Failed to close object"), Mockito.eq(failure));
25+
}
26+
27+
@Test(groups = {"unit"})
28+
public void testQuiteCloseClosesSuccessfully() throws IOException {
29+
Logger log = Mockito.mock(Logger.class);
30+
Closeable closeable = Mockito.mock(Closeable.class);
31+
32+
ClientUtils.quiteClose(closeable, log);
33+
34+
Mockito.verify(closeable).close();
35+
Mockito.verifyNoInteractions(log);
36+
}
37+
38+
@Test(groups = {"unit"})
39+
public void testQuiteCloseWithNull() {
40+
Logger log = Mockito.mock(Logger.class);
41+
42+
// Should be a no-op and not throw on a null closeable
43+
ClientUtils.quiteClose(null, log);
44+
45+
Mockito.verifyNoInteractions(log);
46+
Assert.assertTrue(true);
47+
}
48+
}

0 commit comments

Comments
 (0)