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
5 changes: 5 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -44,6 +44,11 @@

### Bug Fixes

- **[client-v2]** Fixed `Client.cancelTransportRequest(queryId)` being silently dropped when it landed between two
attempts of a retried operation (query, POJO insert and stream insert): the operation issued the next attempt anyway
and could complete successfully. The request of an attempt now stays registered until the whole operation is over,
and the cancellation is checked before every attempt, so a cancelled operation stops instead of sending another
request. (https://github.com/ClickHouse/clickhouse-java/issues/2989)
- **[client-v2]** Fixed LZ4 input streams not closing their underlying HTTP response stream. Closing an LZ4 stream
returned by `QueryResponse.getInputStream()` now releases the wrapped transport stream, including after a partial
read. (https://github.com/ClickHouse/clickhouse-java/issues/2985)
Expand Down
104 changes: 64 additions & 40 deletions client-v2/src/main/java/com/clickhouse/client/api/Client.java
Original file line number Diff line number Diff line change
Expand Up @@ -1465,50 +1465,55 @@ public CompletableFuture<InsertResponse> insert(String tableName, List<?> data,
Endpoint selectedEndpoint = nodeSelector.getEndpoint();
final String queryId = requestSettings.getQueryId();
RuntimeException lastException = null;
for (int i = 0; i <= maxAttempts; i++) {
// Execute request
TransportRequest transportRequest = httpClientHelper.createRequest(selectedEndpoint, requestSettings.getAllSettings(),
out -> {
out.write("INSERT INTO ".getBytes());
out.write(tableName.getBytes());
out.write(" \n FORMAT ".getBytes());
out.write(format.name().getBytes());
out.write(" \n".getBytes());
for (Object obj : data) {

for (POJOFieldSerializer serializer : serializersForTable) {
try {
serializer.serialize(obj, out);
} catch (InvocationTargetException | IllegalAccessException e) {
throw new DataSerializationException(obj, serializer, e);
try {
for (int i = 0; i <= maxAttempts; i++) {
failIfCancelled(queryId, i, lastException);
// Execute request
TransportRequest transportRequest = httpClientHelper.createRequest(selectedEndpoint, requestSettings.getAllSettings(),
out -> {
out.write("INSERT INTO ".getBytes());
out.write(tableName.getBytes());
out.write(" \n FORMAT ".getBytes());
out.write(format.name().getBytes());
out.write(" \n".getBytes());
for (Object obj : data) {

for (POJOFieldSerializer serializer : serializersForTable) {
try {
serializer.serialize(obj, out);
} catch (InvocationTargetException | IllegalAccessException e) {
throw new DataSerializationException(obj, serializer, e);
}
}
}
}
out.close();
});
out.close();
});

registerTransportReq(queryId, transportRequest);
registerTransportReq(queryId, transportRequest);

try (TransportResponse transportResponse = httpClientHelper.executeRequest(transportRequest)) {
ClientStatisticsHolder clientStats = globalClientStats.remove(operationId);
OperationMetrics metrics = completeOperation(transportResponse, clientStats, requestSettings.getQueryId());
try (TransportResponse transportResponse = httpClientHelper.executeRequest(transportRequest)) {
ClientStatisticsHolder clientStats = globalClientStats.remove(operationId);
OperationMetrics metrics = completeOperation(transportResponse, clientStats, requestSettings.getQueryId());

return new InsertResponse(transportResponse, metrics);
} catch (Exception e) {
String msg = requestExMsg("Insert", (i + 1), durationSince(startTime).toMillis(), requestSettings.getQueryId());
lastException = httpClientHelper.wrapException(msg, e, requestSettings.getQueryId());
if (httpClientHelper.shouldRetry(e, requestSettings.getAllSettings()) && requestIsNotCancelled(queryId)) {
if (i < maxAttempts) {
selectedEndpoint = logRetryAndSelectNextNode("Insert", i, maxAttempts, requestSettings.getQueryId(), selectedEndpoint, e);
return new InsertResponse(transportResponse, metrics);
} catch (Exception e) {
String msg = requestExMsg("Insert", (i + 1), durationSince(startTime).toMillis(), requestSettings.getQueryId());
lastException = httpClientHelper.wrapException(msg, e, requestSettings.getQueryId());
if (httpClientHelper.shouldRetry(e, requestSettings.getAllSettings()) && requestIsNotCancelled(queryId)) {
if (i < maxAttempts) {
selectedEndpoint = logRetryAndSelectNextNode("Insert", i, maxAttempts, requestSettings.getQueryId(), selectedEndpoint, e);
} else {
nodeSelector.getNextAliveNode(selectedEndpoint);
}
} else {
nodeSelector.getNextAliveNode(selectedEndpoint);
throw lastException;
}
} else {
throw lastException;
}
} finally {
unregisterTransportReq(queryId);
}
} finally {
// The request of the last attempt stays registered until the operation is over, so a cancellation
// landing between two attempts is not lost.
unregisterTransportReq(queryId);
}

String errMsg = requestExMsg("Insert", maxAttempts + 1, durationSince(startTime).toMillis(), requestSettings.getQueryId());
Expand Down Expand Up @@ -1688,6 +1693,7 @@ public CompletableFuture<InsertResponse> insert(String tableName,
final String queryId = requestSettings.getQueryId();
try {
for (int i = 0; i <= maxAttempts; i++) {
failIfCancelled(queryId, i, lastException);
// Execute request
TransportRequest transportRequest = httpClientHelper.createRequest(selectedEndpoint, requestSettings.getAllSettings(),
out -> {
Expand All @@ -1711,10 +1717,6 @@ public CompletableFuture<InsertResponse> insert(String tableName,
} else {
throw lastException;
}
} finally {
// Insert completes once the request returns; the response exposes no stream to read afterwards,
// so the request is no longer cancellable and can be unregistered.
unregisterTransportReq(requestSettings.getQueryId());
}

if (i < maxAttempts) {
Expand All @@ -1726,6 +1728,8 @@ public CompletableFuture<InsertResponse> insert(String tableName,
}
}
} finally {
// The request of the last attempt stays registered until the operation is over, so a cancellation
// landing between two attempts is not lost.
unregisterTransportReq(queryId);
}

Expand Down Expand Up @@ -1828,6 +1832,7 @@ public CompletableFuture<QueryResponse> query(String sqlQuery, Map<String, Objec
final String queryId = requestSettings.getQueryId();
try {
for (int i = 0; i <= maxAttempts; i++) {
failIfCancelled(queryId, i, lastException);
TransportRequest request = httpClientHelper.createRequest(selectedEndpoint, requestSettings.getAllSettings(), sqlQuery);
registerTransportReq(queryId, request);
TransportResponse transportResp = null;
Expand Down Expand Up @@ -1901,6 +1906,24 @@ private boolean requestIsNotCancelled(String queryId) {
return true;
}

/**
* Stops an operation that was cancelled before its next attempt is issued. A cancellation may have landed
* between two attempts (on the stream insert path for instance from {@link DataStreamWriter#onRetry()}): the
* request of the previous attempt stays registered until the operation is over, so it is seen here and no
* further request is issued.
* Only attempts that follow a failed one are checked: before the first attempt this operation has nothing
* registered yet, so a request found under the same query id belongs to another operation that is still
* running and its cancellation must not stop this one.
* The same exception type and message are used when a cancellation aborts an in-flight request, so a caller
* sees one outcome for a cancelled operation regardless of when the cancellation landed. The failure of the
* last attempt is kept as cause. Called from the retry loop of every operation.
*/
private void failIfCancelled(String queryId, int attempt, RuntimeException lastException) {
if (attempt > 0 && !requestIsNotCancelled(queryId)) {
throw new TransportException("Request was cancelled on client side", lastException, queryId);
}
}

public CompletableFuture<QueryResponse> query(String sqlQuery, Map<String, Object> queryParams) {
return query(sqlQuery, queryParams, null);
}
Expand Down Expand Up @@ -2416,7 +2439,8 @@ public void updateAccessToken(String accessToken) {
* Tries to cancel ongoing request. This method cancels IO operations but doesn't
* kill query on server side. Original queryId should be used to cancel the request.
* This operation cancels only operations on client side and only that still waiting
* for response.
* for response. A cancellation that lands between two attempts of a retried operation is effective too:
* the operation stops instead of issuing another request.
*
* @param queryId - original query id that was passed in operation settings.
*/
Expand Down
Loading
Loading