1616import com .clickhouse .client .api .insert .InsertResponse ;
1717import com .clickhouse .client .api .insert .InsertSettings ;
1818import com .clickhouse .client .api .internal .ClientStatisticsHolder ;
19+ import com .clickhouse .client .api .internal .ClientUtils ;
1920import com .clickhouse .client .api .internal .CredentialsManager ;
2021import com .clickhouse .client .api .internal .HttpAPIClientHelper ;
2122import com .clickhouse .client .api .internal .MapUtils ;
3637import com .clickhouse .client .api .serde .POJOSerDe ;
3738import com .clickhouse .client .api .transport .Endpoint ;
3839import com .clickhouse .client .api .transport .HttpEndpoint ;
40+ import com .clickhouse .client .api .transport .internal .TransportRequest ;
41+ import com .clickhouse .client .api .transport .internal .TransportResponse ;
3942import com .clickhouse .client .config .ClickHouseClientOption ;
4043import com .clickhouse .data .ClickHouseColumn ;
4144import com .clickhouse .data .ClickHouseDataType ;
4245import com .clickhouse .data .ClickHouseFormat ;
4346import com .google .common .collect .ImmutableList ;
4447import net .jpountz .lz4 .LZ4Factory ;
4548import org .apache .hc .core5 .concurrent .DefaultThreadFactory ;
46- import org .apache .hc .core5 .http .ClassicHttpResponse ;
47- import org .apache .hc .core5 .http .Header ;
4849import org .apache .hc .core5 .http .HttpStatus ;
4950import org .slf4j .Logger ;
5051import org .slf4j .LoggerFactory ;
@@ -1319,43 +1320,37 @@ public CompletableFuture<InsertResponse> insert(String tableName, List<?> data,
13191320 RuntimeException lastException = null ;
13201321 for (int i = 0 ; i <= maxRetries ; i ++) {
13211322 // Execute request
1322- try (ClassicHttpResponse httpResponse =
1323- httpClientHelper .executeRequest (selectedEndpoint , requestSettings .getAllSettings (),
1324- out -> {
1325- out .write ("INSERT INTO " .getBytes ());
1326- out .write (tableName .getBytes ());
1327- out .write (" \n FORMAT " .getBytes ());
1328- out .write (format .name ().getBytes ());
1329- out .write (" \n " .getBytes ());
1330- for (Object obj : data ) {
1331-
1332- for (POJOFieldSerializer serializer : serializersForTable ) {
1333- try {
1334- serializer .serialize (obj , out );
1335- } catch (InvocationTargetException | IllegalAccessException | IOException e ) {
1336- throw new DataSerializationException (obj , serializer , e );
1337- }
1338- }
1323+ TransportRequest transportRequest = httpClientHelper .createRequest (selectedEndpoint , requestSettings .getAllSettings (),
1324+ out -> {
1325+ out .write ("INSERT INTO " .getBytes ());
1326+ out .write (tableName .getBytes ());
1327+ out .write (" \n FORMAT " .getBytes ());
1328+ out .write (format .name ().getBytes ());
1329+ out .write (" \n " .getBytes ());
1330+ for (Object obj : data ) {
1331+
1332+ for (POJOFieldSerializer serializer : serializersForTable ) {
1333+ try {
1334+ serializer .serialize (obj , out );
1335+ } catch (InvocationTargetException | IllegalAccessException | IOException e ) {
1336+ throw new DataSerializationException (obj , serializer , e );
13391337 }
1340- out .close ();
1341- })) {
1342-
1343-
1338+ }
1339+ }
1340+ out .close ();
1341+ });
1342+ try (TransportResponse httpResponse = httpClientHelper .executeRequest (transportRequest )) {
13441343 // Check response
1345- if (httpResponse .getCode () == HttpStatus .SC_SERVICE_UNAVAILABLE ) {
1346- LOG .warn ("Failed to get response. Server returned {}. Retrying. (Duration: {})" , httpResponse .getCode (), durationSince (startTime ));
1344+ if (httpResponse .getStatusCode () == HttpStatus .SC_SERVICE_UNAVAILABLE ) {
1345+ LOG .warn ("Failed to get response. Server returned {}. Retrying. (Duration: {})" , httpResponse .getStatusCode (), durationSince (startTime ));
13471346 selectedEndpoint = getNextAliveNode ();
13481347 continue ;
13491348 }
13501349
13511350 ClientStatisticsHolder clientStats = globalClientStats .remove (operationId );
1352- OperationMetrics metrics = new OperationMetrics (clientStats );
1353- String summary = HttpAPIClientHelper .getHeaderVal (httpResponse .getFirstHeader (ClickHouseHttpProto .HEADER_SRV_SUMMARY ), "{}" );
1354- ProcessParser .parseSummary (summary , metrics );
1355- String queryId = HttpAPIClientHelper .getHeaderVal (httpResponse .getFirstHeader (ClickHouseHttpProto .HEADER_QUERY_ID ), requestSettings .getQueryId (), String ::valueOf );
1356- metrics .operationComplete ();
1357- metrics .setQueryId (queryId );
1358- return new InsertResponse (metrics , HttpAPIClientHelper .collectResponseHeaders (httpResponse ));
1351+ OperationMetrics metrics = completeOperation (httpResponse , clientStats , requestSettings .getQueryId ());
1352+
1353+ return new InsertResponse (httpResponse , metrics );
13591354 } catch (Exception e ) {
13601355 String msg = requestExMsg ("Insert" , (i + 1 ), durationSince (startTime ).toMillis (), requestSettings .getQueryId ());
13611356 lastException = httpClientHelper .wrapException (msg , e , requestSettings .getQueryId ());
@@ -1373,7 +1368,6 @@ public CompletableFuture<InsertResponse> insert(String tableName, List<?> data,
13731368 throw (lastException == null ? new ClientException (errMsg ) : lastException ); };
13741369
13751370 return runAsyncOperation (supplier , requestSettings .getAllSettings ());
1376-
13771371 }
13781372
13791373 /**
@@ -1509,7 +1503,7 @@ public CompletableFuture<InsertResponse> insert(String tableName,
15091503 clientStats .start (ClientMetrics .OP_DURATION );
15101504 final ClientStatisticsHolder finalClientStats = clientStats ;
15111505
1512- Supplier < InsertResponse > responseSupplier ;
1506+
15131507
15141508 final int writeBufferSize = requestSettings .getInputStreamCopyBufferSize () <= 0 ?
15151509 (int ) configuration .get (ClientConfigProperties .CLIENT_NETWORK_BUFFER_SIZE .getKey ()) :
@@ -1533,36 +1527,32 @@ public CompletableFuture<InsertResponse> insert(String tableName,
15331527 if (requestSettings .getQueryId () == null && queryIdGenerator != null ) {
15341528 requestSettings .setQueryId (queryIdGenerator .get ());
15351529 }
1536- responseSupplier = () -> {
1530+
1531+ Supplier <InsertResponse > responseSupplier = () -> {
15371532 long startTime = System .nanoTime ();
15381533 // Selecting some node
15391534 Endpoint selectedEndpoint = getNextAliveNode ();
15401535
15411536 RuntimeException lastException = null ;
15421537 for (int i = 0 ; i <= retries ; i ++) {
15431538 // Execute request
1544- try (ClassicHttpResponse httpResponse =
1545- httpClientHelper .executeRequest (selectedEndpoint , requestSettings .getAllSettings (),
1546- out -> {
1547- writer .onOutput (out );
1548- out .close ();
1549- })) {
1539+ TransportRequest transportRequest = httpClientHelper .createRequest (selectedEndpoint , requestSettings .getAllSettings (),
1540+ out -> {
1541+ writer .onOutput (out );
1542+ out .close ();
1543+ });
15501544
1545+ try (TransportResponse httpResponse = httpClientHelper .executeRequest (transportRequest )) {
15511546
15521547 // Check response
1553- if (httpResponse .getCode () == HttpStatus .SC_SERVICE_UNAVAILABLE ) {
1554- LOG .warn ("Failed to get response. Server returned {}. Retrying. (Duration: {})" , httpResponse .getCode (), durationSince (startTime ));
1548+ if (httpResponse .getStatusCode () == HttpStatus .SC_SERVICE_UNAVAILABLE ) {
1549+ LOG .warn ("Failed to get response. Server returned {}. Retrying. (Duration: {})" , httpResponse .getStatusCode (), durationSince (startTime ));
15551550 selectedEndpoint = getNextAliveNode ();
15561551 continue ;
15571552 }
15581553
1559- OperationMetrics metrics = new OperationMetrics (finalClientStats );
1560- String summary = HttpAPIClientHelper .getHeaderVal (httpResponse .getFirstHeader (ClickHouseHttpProto .HEADER_SRV_SUMMARY ), "{}" );
1561- ProcessParser .parseSummary (summary , metrics );
1562- String queryId = HttpAPIClientHelper .getHeaderVal (httpResponse .getFirstHeader (ClickHouseHttpProto .HEADER_QUERY_ID ), requestSettings .getQueryId (), String ::valueOf );
1563- metrics .operationComplete ();
1564- metrics .setQueryId (queryId );
1565- return new InsertResponse (metrics , HttpAPIClientHelper .collectResponseHeaders (httpResponse ));
1554+ OperationMetrics metrics = completeOperation (httpResponse , finalClientStats , requestSettings .getQueryId ());
1555+ return new InsertResponse (httpResponse , metrics );
15661556 } catch (Exception e ) {
15671557 String msg = requestExMsg ("Insert" , (i + 1 ), durationSince (startTime ).toMillis (), requestSettings .getQueryId ());
15681558 lastException = httpClientHelper .wrapException (msg , e , requestSettings .getQueryId ());
@@ -1668,38 +1658,28 @@ public CompletableFuture<QueryResponse> query(String sqlQuery, Map<String, Objec
16681658 Endpoint selectedEndpoint = getNextAliveNode ();
16691659 RuntimeException lastException = null ;
16701660 for (int i = 0 ; i <= retries ; i ++) {
1671- ClassicHttpResponse httpResponse = null ;
1661+ TransportRequest request = httpClientHelper .createRequest (selectedEndpoint , requestSettings .getAllSettings (), sqlQuery );
1662+ TransportResponse transportResp = null ;
16721663 try {
1673- httpResponse = httpClientHelper .executeRequest (selectedEndpoint ,
1674- requestSettings .getAllSettings (),
1675- sqlQuery );
1664+ transportResp = httpClientHelper .executeRequest (request );
16761665 // Check response
1677- if (httpResponse . getCode () == HttpStatus .SC_SERVICE_UNAVAILABLE ) {
1678- LOG .warn ("Failed to get response. Server returned {}. Retrying. (Duration: {})" , httpResponse . getCode (), durationSince (startTime ));
1666+ if (transportResp . getStatusCode () == HttpStatus .SC_SERVICE_UNAVAILABLE ) {
1667+ LOG .warn ("Failed to get response. Server returned {}. Retrying. (Duration: {})" , transportResp . getStatusCode (), durationSince (startTime ));
16791668 selectedEndpoint = getNextAliveNode ();
1680- HttpAPIClientHelper . closeQuietly ( httpResponse );
1669+ ClientUtils . quiteClose ( transportResp , LOG );
16811670 continue ;
16821671 }
16831672
1684- OperationMetrics metrics = new OperationMetrics (clientStats );
1685- String summary = HttpAPIClientHelper .getHeaderVal (httpResponse
1686- .getFirstHeader (ClickHouseHttpProto .HEADER_SRV_SUMMARY ), "{}" );
1687- ProcessParser .parseSummary (summary , metrics );
1688- String queryId = HttpAPIClientHelper .getHeaderVal (httpResponse
1689- .getFirstHeader (ClickHouseHttpProto .HEADER_QUERY_ID ), requestSettings .getQueryId ());
1690- metrics .setQueryId (queryId );
1691- metrics .operationComplete ();
1692- Header formatHeader = httpResponse .getFirstHeader (ClickHouseHttpProto .HEADER_FORMAT );
1693- ClickHouseFormat responseFormat = requestSettings .getFormat ();
1694- if (formatHeader != null ) {
1695- responseFormat = ClickHouseFormat .valueOf (formatHeader .getValue ());
1673+ OperationMetrics metrics = completeOperation (transportResp , clientStats , requestSettings .getQueryId ());
1674+ ClickHouseFormat responseFormat = transportResp .getDataFormat ();
1675+ if (responseFormat == null ) {
1676+ responseFormat = requestSettings .getFormat ();
16961677 }
16971678
1698- return new QueryResponse (httpResponse , responseFormat , requestSettings , metrics ,
1699- HttpAPIClientHelper .collectResponseHeaders (httpResponse ));
1679+ return new QueryResponse (transportResp , responseFormat , requestSettings , metrics );
17001680
17011681 } catch (Exception e ) {
1702- HttpAPIClientHelper . closeQuietly ( httpResponse );
1682+ ClientUtils . quiteClose ( transportResp , LOG );
17031683 String msg = requestExMsg ("Query" , (i + 1 ), durationSince (startTime ).toMillis (), requestSettings .getQueryId ());
17041684 lastException = httpClientHelper .wrapException (msg , e , requestSettings .getQueryId ());
17051685 if (httpClientHelper .shouldRetry (e , requestSettings .getAllSettings ())) {
@@ -1721,6 +1701,16 @@ public CompletableFuture<QueryResponse> query(String sqlQuery, Map<String, Objec
17211701 return query (sqlQuery , queryParams , null );
17221702 }
17231703
1704+ private OperationMetrics completeOperation (TransportResponse transportResponse , ClientStatisticsHolder clientStats , String originalQueryId ) {
1705+ OperationMetrics metrics = new OperationMetrics (clientStats );
1706+ String summary = transportResponse .getSummaryJson ();
1707+ ProcessParser .parseSummary (summary , metrics );
1708+ String queryId = transportResponse .getQueryId ();
1709+ metrics .setQueryId (queryId == null ? originalQueryId : queryId );
1710+ metrics .operationComplete ();
1711+ return metrics ;
1712+ }
1713+
17241714 /**
17251715 * <p>Queries data in one of descriptive format and creates a reader out of the response stream.</p>
17261716 * <p>Format is selected internally so is ignored when passed in settings. If query contains format
0 commit comments