Gabriel39 commented on code in PR #68630:
URL: https://github.com/apache/doris/pull/68630#discussion_r4131992141


##########
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:
   Fixed in 59cbfb7f7c2. Stream aborts now reach the owning FE through a 
result-ID route and fan out the existing query cancellation RPC to every 
participating BE. The route retains only IDs and addresses, so it survives both 
local fragment completion and normal FE coordinator unregistration without 
holding the coordinator or its queue slot. Added BE coverage for an expired 
local context and a finished endpoint, FE coverage for multiple backend 
addresses without a live coordinator, and a distributed regression that aborts 
only one endpoint. The updated cluster regression is pending CI.



##########
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:
   Fixed in 59cbfb7f7c2. ResultBufferMgr now indexes all Flight buffers by 
their actual query ID and removes/clears every sibling on query abort, outside 
the map lock and independently of QueryContext lifetime. The existing 
cancel_plan_fragment handler also invokes this cleanup on each updated BE. 
Added a test that expires the context before closing one endpoint and verifies 
that a sibling buffer and its queued Blocks are removed. LIMIT_REACH and 
FINISHED preserve unread results, with dedicated coverage.



##########
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:
   Fixed in 59cbfb7f7c2. Removed cancellation from fetch_arrow_data entirely. 
The FE route uses cancel_plan_fragment, which bypasses the Arrow Flight work 
pool. Both transport and application statuses are checked; failed routes remain 
available for a bounded retry. Added BE tests with the Arrow pool rejecting 
work and an FE application error, plus FE coverage for a rejected backend 
cancellation followed by a successful retry.



##########
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:
   Fixed in 59cbfb7f7c2 by removing the new fetch cancel field. With the 
serving FE and Flight-facing BE upgraded, cancellation now reaches older result 
BEs through their existing cancel_plan_fragment RPC, using the actual query ID 
from FE routing metadata. Added a legacy-handler compatibility test. Older 
result BEs still retain their historical unread-buffer reclamation policy; 
immediate sibling-buffer removal requires upgrading those BEs, as documented in 
the PR description.



-- 
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