From 8c29b339700274a9d33ca9183141bd0d28ff2c2c Mon Sep 17 00:00:00 2001 From: Sergey Date: Thu, 30 May 2024 12:19:53 -0700 Subject: [PATCH] Revert "fix response close for getRowsUpdated(#1538)" --- .../com/clickhouse/r2dbc/ClickHouseResult.java | 16 ++++------------ .../clickhouse/r2dbc/ClickHouseResult091.java | 14 ++++---------- 2 files changed, 8 insertions(+), 22 deletions(-) diff --git a/clickhouse-r2dbc/src/main/java/com/clickhouse/r2dbc/ClickHouseResult.java b/clickhouse-r2dbc/src/main/java/com/clickhouse/r2dbc/ClickHouseResult.java index 5c934ab12..8601f1236 100644 --- a/clickhouse-r2dbc/src/main/java/com/clickhouse/r2dbc/ClickHouseResult.java +++ b/clickhouse-r2dbc/src/main/java/com/clickhouse/r2dbc/ClickHouseResult.java @@ -31,18 +31,10 @@ public class ClickHouseResult implements Result { .map(rec -> ClickHousePair.of(resp.getColumns(), rec)))) .map(pair -> new ClickHouseRow(pair.getRight(), pair.getLeft())) .map(RowSegment::new); - - this.updatedCount = Mono.using(() -> response, - resp -> Mono.just(response).map(ClickHouseResponse::getSummary) - .map(ClickHouseResponseSummary::getProgress) - .map(ClickHouseResponseSummary.Progress::getWrittenRows) - .map(UpdateCount::new), - resp -> { - if (!resp.isClosed()) { - resp.close(); - } - }); - + this.updatedCount = Mono.just(response).map(ClickHouseResponse::getSummary) + .map(ClickHouseResponseSummary::getProgress) + .map(ClickHouseResponseSummary.Progress::getWrittenRows) + .map(UpdateCount::new); this.segments = Flux.concat(this.updatedCount, this.rowSegments); } diff --git a/clickhouse-r2dbc/src/main/java/com/clickhouse/r2dbc/ClickHouseResult091.java b/clickhouse-r2dbc/src/main/java/com/clickhouse/r2dbc/ClickHouseResult091.java index 535ea4456..4b67bf889 100644 --- a/clickhouse-r2dbc/src/main/java/com/clickhouse/r2dbc/ClickHouseResult091.java +++ b/clickhouse-r2dbc/src/main/java/com/clickhouse/r2dbc/ClickHouseResult091.java @@ -31,16 +31,10 @@ class ClickHouseResult implements Result { .map(rec -> ClickHousePair.of(resp.getColumns(), rec)))) .map(pair -> new ClickHouseRow(pair.getRight(), pair.getLeft())) .map(RowSegment::new); - this.updatedCount = Mono.using(() -> response, - resp -> Mono.just(response).map(ClickHouseResponse::getSummary) - .map(ClickHouseResponseSummary::getProgress) - .map(ClickHouseResponseSummary.Progress::getWrittenRows) - .map(UpdateCount::new), - resp -> { - if (!resp.isClosed()) { - resp.close(); - } - }); + this.updatedCount = Mono.just(response).map(ClickHouseResponse::getSummary) + .map(ClickHouseResponseSummary::getProgress) + .map(ClickHouseResponseSummary.Progress::getWrittenRows) + .map(UpdateCount::new); this.segments = Flux.concat(this.updatedCount, this.rowSegments); }