github-actions[bot] commented on code in PR #68630:
URL: https://github.com/apache/doris/pull/68630#discussion_r4131669262


##########
be/src/exec/sink/writer/varrow_flight_result_writer.cpp:
##########
@@ -97,15 +97,29 @@ Status 
ArrowFlightResultBlockBuffer::get_schema(std::shared_ptr<arrow::Schema>*
                                              print_id(_fragment_id), _status));
 }
 
-Status ArrowFlightResultBlockBuffer::get_arrow_batch(std::shared_ptr<Block>* 
result) {
+void ArrowFlightResultBlockBuffer::cancel_query(const Status& reason) {
+    cancel(reason);
+    // The result buffer can be keyed by a fragment instance id rather than 
the query id.
+    // Keep only a weak query reference, and cancel outside the buffer/map 
locks.
+    if (auto query_ctx = _query_ctx.lock()) {

Review Comment:
   [P1] Send stream abort to the query coordinator after a local fragment 
finishes. In parallel result mode, BE A can finish its fragment while BE B is 
still producing, leaving A's unread Flight buffer registered. A's fragment is 
then removed and the buffer holds only a weak QueryContext reference. If the 
client reads one batch from A and closes before EOF, `_query_ctx.lock()` can 
fail here; even a surviving context has no active local fragment to report a 
new error to FE. BE B therefore receives no cancellation and can remain active 
or backpressured until timeout. Propagate the abort through a query-level route 
that remains usable after the selected result fragment has finished, and cover 
a multi-BE query with one early-finishing endpoint.



##########
be/src/runtime/result_buffer_mgr.cpp:
##########
@@ -141,6 +141,23 @@ bool ResultBufferMgr::cancel(const TUniqueId& unique_id, 
const Status& reason) {
     return exist;
 }
 
+void ResultBufferMgr::cancel_arrow_flight_query(const TUniqueId& buffer_id, 
const Status& reason) {
+    std::shared_ptr<ArrowFlightResultBlockBuffer> buffer;
+    {
+        std::unique_lock<std::shared_mutex> lock(_buffer_map_lock);
+        auto it = _buffer_map.find(buffer_id);
+        if (it == _buffer_map.end()) {
+            return;
+        }
+        buffer = 
std::dynamic_pointer_cast<ArrowFlightResultBlockBuffer>(it->second);
+        if (!buffer) {
+            return;
+        }
+        _buffer_map.erase(it);

Review Comment:
   [P2] Cancel sibling Flight buffers when aborting a query. In nonparallel 
result-sink mode, each instance has its own buffer ID and Nereids can advertise 
multiple endpoints on one BE. Closing one stream erases only its buffer here, 
then cancels the shared QueryContext. The other sinks close with an error, but 
`ResultBlockBuffer::close` leaves their queued Blocks intact and 
`ResultSinkLocalState::close` returns before scheduling deferred cleanup. Those 
buffers retain unread result memory until the original execution timeout. 
Cancel all Flight buffers for the query (or otherwise promptly clear and remove 
sibling buffers) when one endpoint aborts.



##########
gensrc/proto/internal_service.proto:
##########
@@ -334,6 +334,7 @@ message PFetchDataResult {
 
 message PFetchArrowDataRequest {
     optional PUniqueId finst_id = 1;
+    optional bool cancel = 2 [default = false];

Review Comment:
   [P2] Preserve cancellation when the result BE is older. A new proxy BE can 
forward a Flight stream to an older result BE during a rolling upgrade. That BE 
does not know this newly added optional field, so it handles the `cancel=true` 
request as another ordinary fetch; with no ready batch it parks the RPC, and 
the proxy times out after one second while the producer remains blocked. Please 
use a compatible cancellation route or gate the behavior on result-BE 
capability so closing the stream does not leave the old BE's query running 
until timeout.



##########
be/src/service/arrow_flight/arrow_flight_batch_reader.cpp:
##########
@@ -124,21 +203,37 @@ arrow::Status 
ArrowFlightBatchLocalReader::ReadNextImpl(std::shared_ptr<arrow::R
     return arrow::Status::OK();
 }
 
-arrow::Status 
ArrowFlightBatchLocalReader::ReadNext(std::shared_ptr<arrow::RecordBatch>* out) 
{
-    RETURN_ARROW_STATUS_IF_CATCH_EXCEPTION(ReadNextImpl(out));
-}
-
 ArrowFlightBatchRemoteReader::ArrowFlightBatchRemoteReader(
         const std::shared_ptr<QueryStatement>& statement,
         const std::shared_ptr<PBackendService_Stub>& stub)
         : ArrowFlightBatchReaderBase(statement), _brpc_stub(stub), 
_block(nullptr) {
     _mem_tracker = MemTrackerLimiter::create_shared(
             MemTrackerLimiter::Type::QUERY,
             fmt::format("ArrowFlightBatchRemoteReader#QueryId={}", 
print_id(_statement->query_id)));
+    _cancel_query = [stub, id = statement->query_id](const Status&) {
+        auto request = std::make_shared<PFetchArrowDataRequest>();
+        request->mutable_finst_id()->set_hi(id.hi);
+        request->mutable_finst_id()->set_lo(id.lo);
+        request->set_cancel(true);
+        auto callback = 
DummyBrpcCallback<PFetchArrowDataResult>::create_shared();
+        auto closure = AutoReleaseClosure<
+                PFetchArrowDataRequest,
+                
DummyBrpcCallback<PFetchArrowDataResult>>::create_unique(request, callback);
+        // Bound teardown even if the result BE is unavailable.
+        callback->cntl_->set_timeout_ms(1000);
+        stub->fetch_arrow_data(closure->cntl_.get(), closure->request_.get(),
+                               closure->response_.get(), closure.get());
+        closure.release();
+        callback->join();
+        if (callback->cntl_->Failed()) {

Review Comment:
   [P2] Handle a rejected remote cancellation. `fetch_arrow_data` puts this 
cancel request into the bounded Arrow Flight work pool. If the queue is full, 
`offer_failed` returns a protobuf `CANCELLED` status while the brpc controller 
itself succeeds. This code checks only `cntl_->Failed()`, so the reader remains 
closed even though the result BE never cancelled its buffer or query; a 
backpressured producer can stay alive until timeout. Check the response status 
and provide a cleanup path that still succeeds when this pool rejects work.



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to