This is an automated email from the ASF dual-hosted git repository.

yiguolei pushed a commit to branch branch-4.1
in repository https://gitbox.apache.org/repos/asf/doris.git


The following commit(s) were added to refs/heads/branch-4.1 by this push:
     new c6536cfde23 [fix](arrow-flight) Cancel queries when result streams 
close early (#68630)
c6536cfde23 is described below

commit c6536cfde23bf8ad39d87a4498e4ea7eb59111a8
Author: Gabriel <[email protected]>
AuthorDate: Tue Sep 29 22:02:05 2026 +0800

    [fix](arrow-flight) Cancel queries when result streams close early (#68630)
    
    ## What problem does this PR solve?
    
    Cancelling or closing an unfinished Arrow Flight stream stops the client
    RPC but can leave result producers blocked until timeout. Propagate
    stream aborts to the whole query and promptly release its Flight
    buffers, while preserving normal EOF and valid unread results after
    LIMIT_REACH or FINISHED.
    
    Local reads use bounded buffer waits with explicit EOF; the reader owns
    cancellation checks. Remote reads interrupt outstanding BRPC waits. A BE
    reports an aborted result ID to the owning FE, which sends the existing
    cancel_plan_fragment RPC to every participating BE. Lightweight
    result-ID routing metadata survives local fragment and coordinator
    completion without retaining the coordinator or its query queue slot;
    session reset and execution-timeout expiry remove those routes. RPC
    transport and application errors are checked, with bounded retries.
    Cancellation bypasses the Arrow read pool.
    
    On updated BEs, a query-to-buffer index removes all sibling Flight
    buffers, even after their local QueryContext expires. Query cancellation
    occurs outside buffer-map locks. A fetch racing with buffer removal
    completes its RPC with an error.
    
    ## Validation
    
    - 18 BE tests passed in 10 consecutive incremental ASAN runs (180
    executions), including real Flight, Thrift and BRPC transport paths,
    expired contexts, sibling cleanup, rejected Arrow work,
    application-error retry, a legacy cancellation handler, and successful
    fragment-stop preservation.
    - Four FE tests passed for cancellation across backend addresses after
    coordinator completion, application-error retry, expiry and session
    cleanup. Changed FE sources and generated Thrift code compiled against
    cached dependencies.
    - FE Checkstyle, clang-format 16 and git diff whitespace checks passed.
    The Groovy regression compiled against the regression framework.
    - The cluster regression uses a distributed table, aborts only one
    endpoint, checks that the query's pipelines disappear from every BE
    within 20 seconds, then checks connection reuse and normal completion.
    Full execution on rebuilt FE/BE binaries remains pending CI; local
    validation used focused incremental builds.
    
    ## Upgrade behavior
    
    Upgrade the serving FE and Flight-facing BE to enable the new
    cancellation route. Result BEs can still be older: the FE uses their
    existing query cancellation RPC rather than an unknown fetch flag. Older
    result BEs retain their historical unread-buffer reclamation behavior
    until upgraded; immediate sibling-buffer reclamation requires the
    updated BE.
---
 .../sink/writer/varrow_flight_result_writer.cpp    |  21 +-
 .../exec/sink/writer/varrow_flight_result_writer.h |  13 +-
 be/src/runtime/result_buffer_mgr.cpp               |  43 ++
 be/src/runtime/result_buffer_mgr.h                 |   5 +
 .../arrow_flight/arrow_flight_batch_reader.cpp     | 154 +++++-
 .../arrow_flight/arrow_flight_batch_reader.h       |  23 +-
 be/src/service/arrow_flight/flight_sql_service.cpp |   8 +-
 be/src/service/internal_service.cpp                |   9 +-
 .../arrow_flight/arrow_flight_cancel_test.cpp      | 567 +++++++++++++++++++++
 .../main/java/org/apache/doris/qe/Coordinator.java |   5 +
 .../org/apache/doris/qe/NereidsCoordinator.java    |   6 +
 .../java/org/apache/doris/qe/StmtExecutor.java     |   6 +
 .../apache/doris/service/FrontendServiceImpl.java  |   6 +
 .../arrowflight/FlightSqlQueryCancellation.java    | 146 ++++++
 .../arrowflight/results/FlightSqlChannel.java      |  14 +-
 .../FlightSqlQueryCancellationTest.java            | 109 ++++
 gensrc/thrift/FrontendService.thrift               |   3 +
 .../test_flight_cancel_cleanup.groovy              | 128 +++++
 18 files changed, 1227 insertions(+), 39 deletions(-)

diff --git a/be/src/exec/sink/writer/varrow_flight_result_writer.cpp 
b/be/src/exec/sink/writer/varrow_flight_result_writer.cpp
index fa20e42f0e7..67ea4b3810d 100644
--- a/be/src/exec/sink/writer/varrow_flight_result_writer.cpp
+++ b/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()) {
+        query_ctx->cancel(reason);
+    }
+}
+
+Status ArrowFlightResultBlockBuffer::get_arrow_batch(std::shared_ptr<Block>* 
result, bool* eos) {
+    *result = nullptr;
+    *eos = false;
     std::unique_lock<std::mutex> l(_lock);
     Defer defer {[&]() { _update_dependency(); }};
     if (!_status.ok()) {
         return _status;
     }
 
-    while (_result_batch_queue.empty() && _status.ok() && !_is_close) {
-        _arrow_data_arrival.wait_for(l, std::chrono::milliseconds(20));
+    if (!_arrow_data_arrival.wait_for(l, std::chrono::milliseconds(20), [&] {
+            return !_result_batch_queue.empty() || !_status.ok() || _is_close;
+        })) {
+        // Let the reader check RPC cancellation without treating an empty 
wait as EOF.
+        return Status::OK();
     }
 
     if (!_status.ok()) {
@@ -126,6 +140,7 @@ Status 
ArrowFlightResultBlockBuffer::get_arrow_batch(std::shared_ptr<Block>* res
 
     // normal path end
     if (_is_close) {
+        *eos = true;
         if (!_status.ok()) {
             return _status;
         }
diff --git a/be/src/exec/sink/writer/varrow_flight_result_writer.h 
b/be/src/exec/sink/writer/varrow_flight_result_writer.h
index 2d0420ed777..a82598658e7 100644
--- a/be/src/exec/sink/writer/varrow_flight_result_writer.h
+++ b/be/src/exec/sink/writer/varrow_flight_result_writer.h
@@ -20,6 +20,7 @@
 #include "common/status.h"
 #include "exec/sink/writer/result_writer.h"
 #include "exprs/vexpr_fwd.h"
+#include "runtime/query_context.h"
 #include "runtime/result_block_buffer.h"
 #include "runtime/runtime_profile.h"
 
@@ -64,13 +65,19 @@ public:
             : ResultBlockBuffer<GetArrowResultBatchCtx>(id, state, 
buffer_size),
               _arrow_schema(schema),
               _profile("ResultBlockBuffer " + print_id(_fragment_id)),
-              _timezone_obj(state->timezone_obj()) {
+              _timezone_obj(state->timezone_obj()),
+              _query_id(state->query_id()),
+              _query_ctx(state->get_query_ctx() ? 
state->get_query_ctx()->weak_from_this()
+                                                : std::weak_ptr<QueryContext> 
{}) {
         _serialize_batch_ns_timer = ADD_TIMER(&_profile, 
"SerializeBatchNsTime");
         _uncompressed_bytes_counter = ADD_COUNTER(&_profile, 
"UncompressedBytes", TUnit::BYTES);
         _compressed_bytes_counter = ADD_COUNTER(&_profile, "CompressedBytes", 
TUnit::BYTES);
     }
     ~ArrowFlightResultBlockBuffer() override = default;
-    Status get_arrow_batch(std::shared_ptr<Block>* result);
+    // A timed wait may return OK with no block and eos=false; the caller must 
retry.
+    Status get_arrow_batch(std::shared_ptr<Block>* result, bool* eos);
+    void cancel_query(const Status& reason);
+    const TUniqueId& query_id() const { return _query_id; }
     void get_timezone(cctz::time_zone& timezone_obj) { timezone_obj = 
_timezone_obj; }
     Status get_schema(std::shared_ptr<arrow::Schema>* arrow_schema);
 
@@ -83,6 +90,8 @@ private:
     RuntimeProfile::Counter* _uncompressed_bytes_counter = nullptr;
     RuntimeProfile::Counter* _compressed_bytes_counter = nullptr;
     cctz::time_zone _timezone_obj;
+    const TUniqueId _query_id;
+    std::weak_ptr<QueryContext> _query_ctx;
 };
 
 class VArrowFlightResultWriter final : public ResultWriter {
diff --git a/be/src/runtime/result_buffer_mgr.cpp 
b/be/src/runtime/result_buffer_mgr.cpp
index 0e53128c0b3..c4a9a16ad7e 100644
--- a/be/src/runtime/result_buffer_mgr.cpp
+++ b/be/src/runtime/result_buffer_mgr.cpp
@@ -93,6 +93,9 @@ Status ResultBufferMgr::create_sender(const TUniqueId& 
unique_id, int buffer_siz
     {
         std::unique_lock<std::shared_mutex> wlock(_buffer_map_lock);
         _buffer_map.insert(std::make_pair(unique_id, control_block));
+        if (arrow_flight) {
+            _arrow_flight_query_buffers[state->query_id()].insert(unique_id);
+        }
         // ResultBlockBufferBase should destroy after max_timeout
         // for exceed max_timeout FE will return timeout to client
         // otherwise in some case may block all fragment handle threads
@@ -135,12 +138,52 @@ bool ResultBufferMgr::cancel(const TUniqueId& unique_id, 
const Status& reason) {
 
     auto exist = _buffer_map.end() != iter;
     if (exist) {
+        if (auto arrow_buffer =
+                    
std::dynamic_pointer_cast<ArrowFlightResultBlockBuffer>(iter->second)) {
+            auto query = 
_arrow_flight_query_buffers.find(arrow_buffer->query_id());
+            if (query != _arrow_flight_query_buffers.end()) {
+                query->second.erase(unique_id);
+                if (query->second.empty()) {
+                    _arrow_flight_query_buffers.erase(query);
+                }
+            }
+        }
         iter->second->cancel(reason);
         _buffer_map.erase(iter);
     }
     return exist;
 }
 
+void ResultBufferMgr::cancel_arrow_flight_query(const TUniqueId& buffer_id, 
const Status& reason) {
+    std::shared_ptr<ArrowFlightResultBlockBuffer> buffer;
+    if (find_buffer(buffer_id, buffer).ok()) {
+        cancel_arrow_flight_buffers(buffer->query_id(), reason);
+    }
+}
+
+void ResultBufferMgr::cancel_arrow_flight_buffers(const TUniqueId& query_id, 
const Status& reason) {
+    std::vector<std::shared_ptr<ArrowFlightResultBlockBuffer>> buffers;
+    {
+        std::unique_lock<std::shared_mutex> lock(_buffer_map_lock);
+        auto query = _arrow_flight_query_buffers.find(query_id);
+        if (query == _arrow_flight_query_buffers.end()) {
+            return;
+        }
+        for (const auto& buffer_id : query->second) {
+            auto it = _buffer_map.find(buffer_id);
+            DCHECK(it != _buffer_map.end());
+            
buffers.push_back(std::static_pointer_cast<ArrowFlightResultBlockBuffer>(it->second));
+            _buffer_map.erase(it);
+        }
+        _arrow_flight_query_buffers.erase(query);
+    }
+    // Clear all sibling buffers even if their local query context has already 
expired.
+    // Never cancel a query while holding the result-buffer map lock.
+    for (const auto& buffer : buffers) {
+        buffer->cancel_query(reason);
+    }
+}
+
 void ResultBufferMgr::cancel_at_time(time_t cancel_time, const TUniqueId& 
unique_id) {
     std::lock_guard<std::mutex> l(_timeout_lock);
     auto iter = _timeout_map.find(cancel_time);
diff --git a/be/src/runtime/result_buffer_mgr.h 
b/be/src/runtime/result_buffer_mgr.h
index 96cb6a5b7e6..e0170c88176 100644
--- a/be/src/runtime/result_buffer_mgr.h
+++ b/be/src/runtime/result_buffer_mgr.h
@@ -27,6 +27,7 @@
 #include <mutex>
 #include <shared_mutex>
 #include <unordered_map>
+#include <unordered_set>
 #include <vector>
 
 #include "common/status.h"
@@ -71,6 +72,9 @@ public:
     // cancel
     bool cancel(const TUniqueId& unique_id, const Status& reason);
 
+    void cancel_arrow_flight_query(const TUniqueId& buffer_id, const Status& 
reason);
+    void cancel_arrow_flight_buffers(const TUniqueId& query_id, const Status& 
reason);
+
     // cancel one query at a future time.
     void cancel_at_time(time_t cancel_time, const TUniqueId& unique_id);
 
@@ -89,6 +93,7 @@ private:
     std::shared_mutex _buffer_map_lock;
     // buffer block map
     BufferMap _buffer_map;
+    std::unordered_map<TUniqueId, std::unordered_set<TUniqueId>> 
_arrow_flight_query_buffers;
 
     // lock for timeout map
     std::mutex _timeout_lock;
diff --git a/be/src/service/arrow_flight/arrow_flight_batch_reader.cpp 
b/be/src/service/arrow_flight/arrow_flight_batch_reader.cpp
index 73761281799..7e11fe9a423 100644
--- a/be/src/service/arrow_flight/arrow_flight_batch_reader.cpp
+++ b/be/src/service/arrow_flight/arrow_flight_batch_reader.cpp
@@ -23,6 +23,8 @@
 #include <arrow/type.h>
 #include <gen_cpp/internal_service.pb.h>
 
+#include <condition_variable>
+
 #include "core/block/block.h"
 #include "format/arrow/arrow_block_convertor.h"
 #include "format/arrow/arrow_row_batch.h"
@@ -35,9 +37,43 @@
 #include "service/backend_options.h"
 #include "util/brpc_client_cache.h"
 #include "util/brpc_closure.h"
+#include "util/client_cache.h"
 
 namespace doris::flight {
 
+namespace {
+// Poll cancellation on the RPC thread: no background thread may retain 
ServerCallContext.
+template <typename Response>
+class FlightReadCallback : public DummyBrpcCallback<Response> {
+public:
+    void call() override {
+        std::lock_guard lock(_mutex);
+        _done = true;
+        _cv.notify_all();
+    }
+
+    void wait(const std::function<bool()>& is_cancelled) {
+        std::unique_lock lock(_mutex);
+        while (!_done) {
+            if (is_cancelled()) {
+                lock.unlock();
+                brpc::StartCancel(this->call_id_);
+                this->join();
+                return;
+            }
+            _cv.wait_for(lock, std::chrono::milliseconds(20));
+        }
+        lock.unlock();
+        this->join();
+    }
+
+private:
+    std::mutex _mutex;
+    std::condition_variable _cv;
+    bool _done = false;
+};
+} // namespace
+
 ArrowFlightBatchReaderBase::ArrowFlightBatchReaderBase(
         const std::shared_ptr<QueryStatement>& statement)
         : _statement(statement) {}
@@ -55,13 +91,79 @@ arrow::Status 
ArrowFlightBatchReaderBase::_return_invalid_status(const std::stri
     return arrow::Status::Invalid(status_msg);
 }
 
+bool ArrowFlightBatchReaderBase::is_cancelled() const {
+    return _closed.load() || (_is_cancelled && _is_cancelled());
+}
+
+void ArrowFlightBatchReaderBase::close(const Status& reason) {
+    if (_closed.exchange(true) || _eof.load()) {
+        return;
+    }
+    auto* env = ExecEnv::GetInstance();
+    if (_statement->result_addr.hostname == BackendOptions::get_localhost() &&
+        _statement->result_addr.port == config::brpc_port) {
+        env->result_mgr()->cancel_arrow_flight_query(_statement->query_id, 
reason);
+    }
+    // A result endpoint can outlive its local fragment. Find the owning FE by 
buffer ID,
+    // then use its query-wide cancellation route, including for older result 
BEs.
+    for (const auto& [address, info] : env->get_running_frontends()) {
+        Status status;
+        for (int attempt = 0; attempt < 2; ++attempt) {
+            FrontendServiceConnection client(env->frontend_client_cache(), 
address, 2000, &status);
+            if (status.ok()) {
+                try {
+                    TStatus result;
+                    client->cancelFlightQuery(result, _statement->query_id);
+                    status = Status::create(result);
+                    if (status.ok()) {
+                        return;
+                    }
+                    if (result.status_code == TStatusCode::NOT_FOUND) {
+                        break;
+                    }
+                } catch (const apache::thrift::TException& e) {
+                    status = Status::RpcError("Flight cancellation failed: 
{}", e.what());
+                    // Discard a transport that may contain an incomplete 
response.
+                    static_cast<void>(client.reopen(1000));
+                }
+            }
+            LOG(WARNING) << "Failed to cancel Flight result " << 
print_id(_statement->query_id)
+                         << " through FE " << address << ": " << status;
+        }
+    }
+}
+
+arrow::Status ArrowFlightBatchReaderBase::Close() {
+    close(Status::Cancelled("Arrow Flight result stream closed before EOF"));
+    return arrow::Status::OK();
+}
+
+arrow::Status 
ArrowFlightBatchReaderBase::ReadNext(std::shared_ptr<arrow::RecordBatch>* out) {
+    *out = nullptr;
+    if (is_cancelled()) {
+        (void)Close();
+        return arrow::Status::Cancelled("Arrow Flight fetch cancelled");
+    }
+    auto status = [&]() -> arrow::Status {
+        RETURN_ARROW_STATUS_IF_CATCH_EXCEPTION(ReadNextImpl(out));
+    }();
+    if (!status.ok()) {
+        close(to_doris_status(status));
+    } else if (!*out) {
+        _eof = true;
+    }
+    return status;
+}
+
 ArrowFlightBatchReaderBase::~ArrowFlightBatchReaderBase() {
+    // Transport errors can destroy a stream without another ReadNext/Close 
call.
+    (void)Close();
     LOG(INFO) << fmt::format(
             "ArrowFlightBatchReader finished, packet_seq={}, 
result_addr={}:{}, finistId={}, "
             "convert_arrow_batch_timer={}, deserialize_block_timer={}, 
peak_memory_usage={}",
             _packet_seq, _statement->result_addr.hostname, 
_statement->result_addr.port,
             print_id(_statement->query_id), _convert_arrow_batch_timer, 
_deserialize_block_timer,
-            _mem_tracker->peak_consumption());
+            _mem_tracker ? _mem_tracker->peak_consumption() : 0);
 }
 
 ArrowFlightBatchLocalReader::ArrowFlightBatchLocalReader(
@@ -74,7 +176,7 @@ ArrowFlightBatchLocalReader::ArrowFlightBatchLocalReader(
 }
 
 arrow::Result<std::shared_ptr<ArrowFlightBatchLocalReader>> 
ArrowFlightBatchLocalReader::Create(
-        const std::shared_ptr<QueryStatement>& statement) {
+        const std::shared_ptr<QueryStatement>& statement, 
std::function<bool()> is_cancelled) {
     DCHECK(statement->result_addr.hostname == BackendOptions::get_localhost());
     std::shared_ptr<ArrowFlightResultBlockBuffer> arrow_buffer;
     RETURN_ARROW_STATUS_IF_ERROR(
@@ -86,6 +188,7 @@ arrow::Result<std::shared_ptr<ArrowFlightBatchLocalReader>> 
ArrowFlightBatchLoca
     std::shared_ptr<MemTrackerLimiter> mem_tracker = 
arrow_buffer->mem_tracker();
     std::shared_ptr<ArrowFlightBatchLocalReader> result(
             new ArrowFlightBatchLocalReader(statement, schema, mem_tracker));
+    result->_is_cancelled = std::move(is_cancelled);
     arrow_buffer->get_timezone(result->_timezone_obj);
     return result;
 }
@@ -99,10 +202,16 @@ arrow::Status 
ArrowFlightBatchLocalReader::ReadNextImpl(std::shared_ptr<arrow::R
     RETURN_ARROW_STATUS_IF_ERROR(
             ExecEnv::GetInstance()->result_mgr()->find_buffer(tid, 
arrow_buffer));
     std::shared_ptr<Block> result;
-    auto st = arrow_buffer->get_arrow_batch(&result);
-    st.prepend("ArrowFlightBatchLocalReader fetch arrow data failed");
-    ARROW_RETURN_NOT_OK(to_arrow_status(st));
-    if (result == nullptr) {
+    bool eos = false;
+    while (!result && !eos) {
+        if (is_cancelled()) {
+            return arrow::Status::Cancelled("Arrow Flight fetch cancelled");
+        }
+        auto st = arrow_buffer->get_arrow_batch(&result, &eos);
+        st.prepend("ArrowFlightBatchLocalReader fetch arrow data failed");
+        ARROW_RETURN_NOT_OK(to_arrow_status(st));
+    }
+    if (eos) {
         // eof, normal path end
         return arrow::Status::OK();
     }
@@ -110,8 +219,8 @@ arrow::Status 
ArrowFlightBatchLocalReader::ReadNextImpl(std::shared_ptr<arrow::R
     {
         // convert one batch
         SCOPED_ATOMIC_TIMER(&_convert_arrow_batch_timer);
-        st = ArrowFlightArrowBlockConvertor(_schema, _timezone_obj)
-                     .convert_to_arrow(*result, arrow::default_memory_pool(), 
out);
+        auto st = ArrowFlightArrowBlockConvertor(_schema, _timezone_obj)
+                          .convert_to_arrow(*result, 
arrow::default_memory_pool(), out);
         st.prepend("ArrowFlightBatchLocalReader convert block to arrow batch 
failed");
         ARROW_RETURN_NOT_OK(to_arrow_status(st));
     }
@@ -124,10 +233,6 @@ 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)
@@ -138,7 +243,7 @@ ArrowFlightBatchRemoteReader::ArrowFlightBatchRemoteReader(
 }
 
 arrow::Result<std::shared_ptr<ArrowFlightBatchRemoteReader>> 
ArrowFlightBatchRemoteReader::Create(
-        const std::shared_ptr<QueryStatement>& statement) {
+        const std::shared_ptr<QueryStatement>& statement, 
std::function<bool()> is_cancelled) {
     std::shared_ptr<PBackendService_Stub> stub =
             ExecEnv::GetInstance()->brpc_internal_client_cache()->get_client(
                     statement->result_addr);
@@ -153,6 +258,7 @@ 
arrow::Result<std::shared_ptr<ArrowFlightBatchRemoteReader>> ArrowFlightBatchRem
 
     std::shared_ptr<ArrowFlightBatchRemoteReader> result(
             new ArrowFlightBatchRemoteReader(statement, stub));
+    result->_is_cancelled = std::move(is_cancelled);
     ARROW_RETURN_NOT_OK(result->init_schema());
     return result;
 }
@@ -163,17 +269,20 @@ arrow::Status 
ArrowFlightBatchRemoteReader::_fetch_schema() {
     auto* pfinst_id = request->mutable_finst_id();
     pfinst_id->set_hi(_statement->query_id.hi);
     pfinst_id->set_lo(_statement->query_id.lo);
-    auto callback = 
DummyBrpcCallback<PFetchArrowFlightSchemaResult>::create_shared();
+    auto callback = 
std::make_shared<FlightReadCallback<PFetchArrowFlightSchemaResult>>();
     auto closure = AutoReleaseClosure<
             PFetchArrowFlightSchemaRequest,
-            
DummyBrpcCallback<PFetchArrowFlightSchemaResult>>::create_unique(request, 
callback);
+            
FlightReadCallback<PFetchArrowFlightSchemaResult>>::create_unique(request, 
callback);
     
callback->cntl_->set_timeout_ms(config::arrow_flight_reader_brpc_controller_timeout_ms);
     callback->cntl_->ignore_eovercrowded();
 
     _brpc_stub->fetch_arrow_flight_schema(closure->cntl_.get(), 
closure->request_.get(),
                                           closure->response_.get(), 
closure.get());
     closure.release();
-    callback->join();
+    callback->wait([this] { return is_cancelled(); });
+    if (is_cancelled()) {
+        return arrow::Status::Cancelled("Arrow Flight fetch cancelled");
+    }
 
     if (callback->cntl_->Failed()) {
         if (!ExecEnv::GetInstance()->brpc_internal_client_cache()->available(
@@ -212,17 +321,20 @@ arrow::Status ArrowFlightBatchRemoteReader::_fetch_data() 
{
         auto* pfinst_id = request->mutable_finst_id();
         pfinst_id->set_hi(_statement->query_id.hi);
         pfinst_id->set_lo(_statement->query_id.lo);
-        auto callback = 
DummyBrpcCallback<PFetchArrowDataResult>::create_shared();
+        auto callback = 
std::make_shared<FlightReadCallback<PFetchArrowDataResult>>();
         auto closure = AutoReleaseClosure<
                 PFetchArrowDataRequest,
-                
DummyBrpcCallback<PFetchArrowDataResult>>::create_unique(request, callback);
+                
FlightReadCallback<PFetchArrowDataResult>>::create_unique(request, callback);
         
callback->cntl_->set_timeout_ms(config::arrow_flight_reader_brpc_controller_timeout_ms);
         callback->cntl_->ignore_eovercrowded();
 
         _brpc_stub->fetch_arrow_data(closure->cntl_.get(), 
closure->request_.get(),
                                      closure->response_.get(), closure.get());
         closure.release();
-        callback->join();
+        callback->wait([this] { return is_cancelled(); });
+        if (is_cancelled()) {
+            return arrow::Status::Cancelled("Arrow Flight fetch cancelled");
+        }
 
         if (callback->cntl_->Failed()) {
             if 
(!ExecEnv::GetInstance()->brpc_internal_client_cache()->available(
@@ -317,8 +429,4 @@ arrow::Status 
ArrowFlightBatchRemoteReader::ReadNextImpl(std::shared_ptr<arrow::
     return arrow::Status::OK();
 }
 
-arrow::Status 
ArrowFlightBatchRemoteReader::ReadNext(std::shared_ptr<arrow::RecordBatch>* 
out) {
-    RETURN_ARROW_STATUS_IF_CATCH_EXCEPTION(ReadNextImpl(out));
-}
-
 } // namespace doris::flight
diff --git a/be/src/service/arrow_flight/arrow_flight_batch_reader.h 
b/be/src/service/arrow_flight/arrow_flight_batch_reader.h
index f72a35a8479..d8dc6618ebe 100644
--- a/be/src/service/arrow_flight/arrow_flight_batch_reader.h
+++ b/be/src/service/arrow_flight/arrow_flight_batch_reader.h
@@ -20,6 +20,7 @@
 #include <cctz/time_zone.h>
 #include <gen_cpp/Types_types.h>
 
+#include <functional>
 #include <memory>
 #include <utility>
 
@@ -48,10 +49,19 @@ class ArrowFlightBatchReaderBase : public 
arrow::RecordBatchReader {
 public:
     // RecordBatchReader force override
     [[nodiscard]] std::shared_ptr<arrow::Schema> schema() const override;
+    arrow::Status Close() override;
+    arrow::Status ReadNext(std::shared_ptr<arrow::RecordBatch>* out) override;
 
 protected:
     ArrowFlightBatchReaderBase(const std::shared_ptr<QueryStatement>& 
statement);
     ~ArrowFlightBatchReaderBase() override;
+    virtual arrow::Status ReadNextImpl(std::shared_ptr<arrow::RecordBatch>* 
out) = 0;
+    bool is_cancelled() const;
+    void close(const Status& reason);
+    std::function<bool()> _is_cancelled;
+    std::atomic<bool> _closed {false};
+    std::atomic<bool> _eof {false};
+
     arrow::Status _return_invalid_status(const std::string& msg);
 
     std::shared_ptr<QueryStatement> _statement;
@@ -67,34 +77,33 @@ protected:
 class ArrowFlightBatchLocalReader : public ArrowFlightBatchReaderBase {
 public:
     static arrow::Result<std::shared_ptr<ArrowFlightBatchLocalReader>> Create(
-            const std::shared_ptr<QueryStatement>& statement);
-
-    arrow::Status ReadNext(std::shared_ptr<arrow::RecordBatch>* out) override;
+            const std::shared_ptr<QueryStatement>& statement,
+            std::function<bool()> is_cancelled = {});
 
 private:
     ArrowFlightBatchLocalReader(const std::shared_ptr<QueryStatement>& 
statement,
                                 const std::shared_ptr<arrow::Schema>& schema,
                                 const std::shared_ptr<MemTrackerLimiter>& 
mem_tracker);
 
-    arrow::Status ReadNextImpl(std::shared_ptr<arrow::RecordBatch>* out);
+    arrow::Status ReadNextImpl(std::shared_ptr<arrow::RecordBatch>* out) 
override;
 };
 
 class ArrowFlightBatchRemoteReader : public ArrowFlightBatchReaderBase {
 public:
     static arrow::Result<std::shared_ptr<ArrowFlightBatchRemoteReader>> Create(
-            const std::shared_ptr<QueryStatement>& statement);
+            const std::shared_ptr<QueryStatement>& statement,
+            std::function<bool()> is_cancelled = {});
 
     // create arrow RecordBatchReader must initialize the schema.
     // so when creating arrow RecordBatchReader, fetch result data once,
     // which will return Block and some necessary information, and extract 
arrow schema from Block.
     arrow::Status init_schema();
-    arrow::Status ReadNext(std::shared_ptr<arrow::RecordBatch>* out) override;
 
 private:
     ArrowFlightBatchRemoteReader(const std::shared_ptr<QueryStatement>& 
statement,
                                  const std::shared_ptr<PBackendService_Stub>& 
stub);
 
-    arrow::Status ReadNextImpl(std::shared_ptr<arrow::RecordBatch>* out);
+    arrow::Status ReadNextImpl(std::shared_ptr<arrow::RecordBatch>* out) 
override;
     arrow::Status _fetch_schema();
     arrow::Status _fetch_data();
 
diff --git a/be/src/service/arrow_flight/flight_sql_service.cpp 
b/be/src/service/arrow_flight/flight_sql_service.cpp
index 26129559b2f..34acc0d6008 100644
--- a/be/src/service/arrow_flight/flight_sql_service.cpp
+++ b/be/src/service/arrow_flight/flight_sql_service.cpp
@@ -87,11 +87,15 @@ public:
         if (statement->result_addr.hostname == BackendOptions::get_localhost() 
&&
             statement->result_addr.port == config::brpc_port) {
             std::shared_ptr<ArrowFlightBatchLocalReader> reader;
-            ARROW_ASSIGN_OR_RAISE(reader, 
ArrowFlightBatchLocalReader::Create(statement));
+            ARROW_ASSIGN_OR_RAISE(
+                    reader, ArrowFlightBatchLocalReader::Create(
+                                    statement, [&context] { return 
context.is_cancelled(); }));
             return std::make_unique<arrow::flight::RecordBatchStream>(reader);
         } else {
             std::shared_ptr<ArrowFlightBatchRemoteReader> reader;
-            ARROW_ASSIGN_OR_RAISE(reader, 
ArrowFlightBatchRemoteReader::Create(statement));
+            ARROW_ASSIGN_OR_RAISE(
+                    reader, ArrowFlightBatchRemoteReader::Create(
+                                    statement, [&context] { return 
context.is_cancelled(); }));
             return std::make_unique<arrow::flight::RecordBatchStream>(reader);
         }
     }
diff --git a/be/src/service/internal_service.cpp 
b/be/src/service/internal_service.cpp
index 30a042f82f4..887140a97aa 100644
--- a/be/src/service/internal_service.cpp
+++ b/be/src/service/internal_service.cpp
@@ -661,6 +661,11 @@ void 
PInternalService::cancel_plan_fragment(google::protobuf::RpcController* /*c
         LOG(INFO) << fmt::format("Cancel query {}, reason: {}", 
print_id(query_id),
                                  actual_cancel_status.to_string());
         _exec_env->fragment_mgr()->cancel_query(query_id, 
actual_cancel_status);
+        // LIMIT_REACH/FINISHED can leave valid unread results; only an abort 
may discard them.
+        if (!actual_cancel_status.ok() && 
!actual_cancel_status.is<ErrorCode::LIMIT_REACH>() &&
+            !actual_cancel_status.is<ErrorCode::FINISHED>()) {
+            _exec_env->result_mgr()->cancel_arrow_flight_buffers(query_id, 
actual_cancel_status);
+        }
 
         // TODO: the logic seems useless, cancel only return Status::OK. 
remove it
         st.to_protobuf(result->mutable_status());
@@ -694,12 +699,14 @@ void 
PInternalService::fetch_arrow_data(google::protobuf::RpcController* control
                                         PFetchArrowDataResult* result,
                                         google::protobuf::Closure* done) {
     bool ret = _arrow_flight_work_pool.try_offer([request, result, done]() {
-        auto ctx = GetArrowResultBatchCtx::create_shared(result, done);
         TUniqueId unique_id = UniqueId(request->finst_id()).to_thrift(); // 
query_id or instance_id
+        auto ctx = GetArrowResultBatchCtx::create_shared(result, done);
         std::shared_ptr<ArrowFlightResultBlockBuffer> arrow_buffer;
         auto st = ExecEnv::GetInstance()->result_mgr()->find_buffer(unique_id, 
arrow_buffer);
         if (!st.ok()) {
             LOG(WARNING) << "Result buffer not found! Query ID: " << 
print_id(unique_id);
+            // A cancellation can remove the buffer before an already queued 
fetch runs.
+            ctx->on_failure(st);
             return;
         }
         if (st = arrow_buffer->get_batch(ctx); !st.ok()) {
diff --git a/be/test/service/arrow_flight/arrow_flight_cancel_test.cpp 
b/be/test/service/arrow_flight/arrow_flight_cancel_test.cpp
new file mode 100644
index 00000000000..2e8d343bc70
--- /dev/null
+++ b/be/test/service/arrow_flight/arrow_flight_cancel_test.cpp
@@ -0,0 +1,567 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements.  See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership.  The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License.  You may obtain a copy of the License at
+//
+//   http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied.  See the License for the
+// specific language governing permissions and limitations
+// under the License.
+
+#include <arrow/flight/client.h>
+#include <arrow/flight/sql/server.h>
+#include <brpc/server.h>
+#include <gtest/gtest.h>
+#include <thrift/protocol/TBinaryProtocol.h>
+#include <thrift/server/TThreadedServer.h>
+#include <thrift/transport/TBufferTransports.h>
+#include <thrift/transport/TServerSocket.h>
+
+#include <future>
+#include <thread>
+
+#include "exec/pipeline/dependency.h"
+#include "load/channel/load_stream_mgr.h"
+#include "runtime/result_buffer_mgr.h"
+#include "service/arrow_flight/arrow_flight_batch_reader.h"
+#include "service/arrow_flight/flight_sql_service.h"
+#include "service/backend_options.h"
+#include "service/internal_service.h"
+#include "testutil/column_helper.h"
+#include "testutil/mock/mock_runtime_state.h"
+#include "util/brpc_client_cache.h"
+#include "util/client_cache.h"
+#include "util/dns_cache.h"
+
+namespace doris::flight {
+
+class ArrowFlightCancelTest : public testing::Test {
+protected:
+    void SetUp() override {
+        _previous_mgr = std::exchange(ExecEnv::GetInstance()->_result_mgr, 
&_mgr);
+        _state._batch_size = 1;
+        _id.hi = 1;
+        _id.lo = 2;
+        auto schema = arrow::schema({arrow::field("value", arrow::int64())});
+        std::shared_ptr<ResultBlockBufferBase> buffer;
+        ASSERT_TRUE(_mgr.create_sender(_id, 16, &buffer, &_state, true, 
schema).ok());
+        _buffer = 
std::dynamic_pointer_cast<ArrowFlightResultBlockBuffer>(buffer);
+        _dep = Dependency::create_shared(0, 0, "Result", true);
+        _buffer->set_dependency(_state.fragment_instance_id(), _dep);
+        TNetworkAddress address;
+        address.hostname = BackendOptions::get_localhost();
+        address.port = config::brpc_port;
+        _statement = std::make_shared<QueryStatement>(_id, address, "select 
value");
+    }
+    void TearDown() override { ExecEnv::GetInstance()->_result_mgr = 
_previous_mgr; }
+    void fill_buffer() {
+        auto block = 
std::make_shared<Block>(ColumnHelper::create_block<DataTypeInt64>({1, 2}));
+        ASSERT_TRUE(_buffer->add_batch(&_state, block).ok());
+        ASSERT_FALSE(_dep->ready());
+    }
+    bool registered() {
+        std::shared_ptr<ArrowFlightResultBlockBuffer> found;
+        return _mgr.find_buffer(_id, found).ok();
+    }
+    MockRuntimeState _state;
+    ResultBufferMgr _mgr;
+    ResultBufferMgr* _previous_mgr = nullptr;
+    TUniqueId _id;
+    std::shared_ptr<ArrowFlightResultBlockBuffer> _buffer;
+    std::shared_ptr<Dependency> _dep;
+    std::shared_ptr<QueryStatement> _statement;
+};
+
+TEST_F(ArrowFlightCancelTest, EarlyCloseReleasesBackpressuredBuffer) {
+    fill_buffer();
+    auto result = ArrowFlightBatchLocalReader::Create(_statement);
+    ASSERT_TRUE(result.ok()) << result.status();
+    auto reader = *result;
+    ASSERT_TRUE(reader->Close().ok());
+    EXPECT_FALSE(registered());
+    EXPECT_TRUE(_dep->ready());
+    EXPECT_TRUE(_state.get_query_ctx()->is_cancelled());
+    std::shared_ptr<Block> block;
+    bool eos = false;
+    EXPECT_FALSE(_buffer->get_arrow_batch(&block, &eos).ok());
+    EXPECT_TRUE(reader->Close().ok());
+}
+
+TEST_F(ArrowFlightCancelTest, DestructionReleasesAbandonedBuffer) {
+    fill_buffer();
+    {
+        auto reader = ArrowFlightBatchLocalReader::Create(_statement);
+        ASSERT_TRUE(reader.ok()) << reader.status();
+    }
+    EXPECT_FALSE(registered());
+    EXPECT_TRUE(_dep->ready());
+    EXPECT_TRUE(_state.get_query_ctx()->is_cancelled());
+}
+
+TEST_F(ArrowFlightCancelTest, CloseReleasesSiblingBuffersAfterContextExpires) {
+    TUniqueId sibling_id = _id;
+    ++sibling_id.lo;
+    std::shared_ptr<ResultBlockBufferBase> sibling;
+    ASSERT_TRUE(_mgr.create_sender(sibling_id, 16, &sibling, &_state, true,
+                                   arrow::schema({arrow::field("value", 
arrow::int64())}))
+                        .ok());
+    auto arrow_sibling = 
std::dynamic_pointer_cast<ArrowFlightResultBlockBuffer>(sibling);
+    auto block = 
std::make_shared<Block>(ColumnHelper::create_block<DataTypeInt64>({1, 2}));
+    ASSERT_TRUE(arrow_sibling->add_batch(&_state, block).ok());
+    auto reader = ArrowFlightBatchLocalReader::Create(_statement);
+    ASSERT_TRUE(reader.ok()) << reader.status();
+    _state._query_ctx_uptr.reset();
+    _state._query_ctx = nullptr;
+    ASSERT_TRUE((*reader)->Close().ok());
+    EXPECT_FALSE(registered());
+    std::shared_ptr<ArrowFlightResultBlockBuffer> found;
+    EXPECT_FALSE(_mgr.find_buffer(sibling_id, found).ok());
+    EXPECT_TRUE(arrow_sibling->_result_batch_queue.empty());
+}
+
+TEST_F(ArrowFlightCancelTest, NormalEofIsNotCancellation) {
+    bool fully_closed = false;
+    ASSERT_TRUE(_buffer->close(_state.fragment_instance_id(), Status::OK(), 0, 
fully_closed).ok());
+    auto result = ArrowFlightBatchLocalReader::Create(_statement);
+    ASSERT_TRUE(result.ok()) << result.status();
+    auto reader = *result;
+    std::shared_ptr<arrow::RecordBatch> batch;
+    ASSERT_TRUE(reader->ReadNext(&batch).ok());
+    EXPECT_EQ(batch, nullptr);
+    EXPECT_TRUE(reader->Close().ok());
+    std::shared_ptr<Block> block;
+    bool eos = false;
+    EXPECT_TRUE(_buffer->get_arrow_batch(&block, &eos).ok());
+    EXPECT_TRUE(eos);
+    EXPECT_FALSE(_state.get_query_ctx()->is_cancelled());
+}
+
+TEST_F(ArrowFlightCancelTest, CancellationInterruptsEmptyBufferWait) {
+    std::atomic<bool> cancelled = false;
+    std::promise<void> checked;
+    std::atomic<int> checks = 0;
+    auto result = ArrowFlightBatchLocalReader::Create(_statement, [&] {
+        if (checks.fetch_add(1) == 3) {
+            checked.set_value();
+        }
+        return cancelled.load();
+    });
+    ASSERT_TRUE(result.ok()) << result.status();
+    auto reader = *result;
+    auto fetch = std::async(std::launch::async, [&] {
+        std::shared_ptr<arrow::RecordBatch> batch;
+        return reader->ReadNext(&batch);
+    });
+    EXPECT_EQ(checked.get_future().wait_for(std::chrono::seconds(2)), 
std::future_status::ready);
+    cancelled = true;
+    const auto ready = fetch.wait_for(std::chrono::seconds(2));
+    // Ensure a broken implementation fails instead of hanging the test runner.
+    if (ready != std::future_status::ready) {
+        _buffer->cancel(Status::Cancelled("test cleanup"));
+    }
+    EXPECT_EQ(ready, std::future_status::ready);
+    EXPECT_FALSE(fetch.get().ok());
+    EXPECT_FALSE(registered());
+    EXPECT_TRUE(_state.get_query_ctx()->is_cancelled());
+}
+
+TEST_F(ArrowFlightCancelTest, EmptyBufferReadReturnsForRetry) {
+    auto fetch = std::async(std::launch::async, [&] {
+        auto block = std::make_shared<Block>();
+        bool eos = true;
+        auto status = _buffer->get_arrow_batch(&block, &eos);
+        EXPECT_EQ(block, nullptr);
+        EXPECT_FALSE(eos);
+        return status;
+    });
+    const auto ready = fetch.wait_for(std::chrono::seconds(2));
+    // Bound the test even if an empty buffer incorrectly waits until query 
completion.
+    if (ready != std::future_status::ready) {
+        _buffer->cancel(Status::Cancelled("test cleanup"));
+    }
+    EXPECT_EQ(ready, std::future_status::ready);
+    EXPECT_TRUE(fetch.get().ok());
+}
+
+TEST_F(ArrowFlightCancelTest, EmptyWaitsDoNotFinishReader) {
+    std::promise<void> retried;
+    std::atomic<int> checks = 0;
+    auto result = ArrowFlightBatchLocalReader::Create(_statement, [&] {
+        if (checks.fetch_add(1) == 3) {
+            retried.set_value();
+        }
+        return false;
+    });
+    ASSERT_TRUE(result.ok()) << result.status();
+    auto reader = *result;
+    std::shared_ptr<arrow::RecordBatch> batch;
+    auto fetch = std::async(std::launch::async, [&] { return 
reader->ReadNext(&batch); });
+    EXPECT_EQ(retried.get_future().wait_for(std::chrono::seconds(2)), 
std::future_status::ready);
+    EXPECT_EQ(fetch.wait_for(std::chrono::seconds(0)), 
std::future_status::timeout);
+    auto block = 
std::make_shared<Block>(ColumnHelper::create_block<DataTypeInt64>({1, 2}));
+    EXPECT_TRUE(_buffer->add_batch(&_state, block).ok());
+    bool fully_closed = false;
+    EXPECT_TRUE(_buffer->close(_state.fragment_instance_id(), Status::OK(), 2, 
fully_closed).ok());
+    const auto ready = fetch.wait_for(std::chrono::seconds(2));
+    if (ready != std::future_status::ready) {
+        _buffer->cancel(Status::Cancelled("test cleanup"));
+    }
+    EXPECT_EQ(ready, std::future_status::ready);
+    ASSERT_TRUE(fetch.get().ok());
+    ASSERT_NE(batch, nullptr);
+    EXPECT_EQ(batch->num_rows(), 2);
+    ASSERT_TRUE(reader->ReadNext(&batch).ok());
+    EXPECT_EQ(batch, nullptr);
+    EXPECT_TRUE(reader->Close().ok());
+    EXPECT_FALSE(_state.get_query_ctx()->is_cancelled());
+}
+
+TEST_F(ArrowFlightCancelTest, ConversionFailureCancelsQuery) {
+    auto block = 
std::make_shared<Block>(ColumnHelper::create_block<DataTypeString>({"bad"}));
+    ASSERT_TRUE(_buffer->add_batch(&_state, block).ok());
+    auto result = ArrowFlightBatchLocalReader::Create(_statement);
+    ASSERT_TRUE(result.ok()) << result.status();
+    std::shared_ptr<arrow::RecordBatch> batch;
+    EXPECT_FALSE((*result)->ReadNext(&batch).ok());
+    EXPECT_FALSE(registered());
+    EXPECT_TRUE(_state.get_query_ctx()->is_cancelled());
+}
+
+TEST_F(ArrowFlightCancelTest, RealFlightRpcCancellationCleansUpQuery) {
+    auto created = FlightSqlServer::create();
+    ASSERT_TRUE(created.ok()) << created.status();
+    auto server = *created;
+    ASSERT_TRUE(server->init(0).ok());
+    auto location = arrow::flight::Location::ForGrpcTcp("127.0.0.1", 
server->port());
+    ASSERT_TRUE(location.ok()) << location.status();
+    auto client_options = arrow::flight::FlightClientOptions::Defaults();
+    client_options.generic_options.emplace_back("grpc.enable_http_proxy", 0);
+    auto connected = arrow::flight::FlightClient::Connect(*location, 
client_options);
+    ASSERT_TRUE(connected.ok()) << connected.status();
+    auto client = std::move(*connected);
+    auto handle = arrow::flight::sql::CreateStatementQueryTicket(
+            print_id(_id) + "&" + BackendOptions::get_localhost() + "&" +
+            std::to_string(config::brpc_port) + "&");
+    ASSERT_TRUE(handle.ok()) << handle.status();
+    arrow::flight::FlightCallOptions options;
+    options.timeout = std::chrono::seconds(5);
+    auto fetched = client->DoGet(options, arrow::flight::Ticket {*handle});
+    ASSERT_TRUE(fetched.ok()) << fetched.status();
+    auto reader = std::move(*fetched);
+    reader->Cancel();
+    for (int i = 0; i < 200 && (registered() || 
!_state.get_query_ctx()->is_cancelled()); ++i) {
+        std::this_thread::sleep_for(std::chrono::milliseconds(10));
+    }
+    EXPECT_FALSE(registered());
+    EXPECT_TRUE(_state.get_query_ctx()->is_cancelled());
+    _buffer->cancel(Status::Cancelled("test cleanup"));
+    reader.reset();
+    EXPECT_TRUE(client->Close().ok());
+    EXPECT_TRUE(server->join().ok());
+}
+
+class FlightCancelCoordinator : public FrontendServiceNull {
+public:
+    void cancelFlightQuery(TStatus& result, const TUniqueId& result_id) 
override {
+        if (++calls <= rejections) {
+            result.__set_status_code(TStatusCode::CANCELLED);
+            return;
+        }
+        if (result_id != buffer_id) {
+            result.__set_status_code(TStatusCode::NOT_FOUND);
+            return;
+        }
+        PCancelPlanFragmentRequest request;
+        request.mutable_finst_id()->set_hi(result_id.hi);
+        request.mutable_finst_id()->set_lo(result_id.lo);
+        request.mutable_query_id()->set_hi(query_id.hi);
+        request.mutable_query_id()->set_lo(query_id.lo);
+        Status::Cancelled("Flight stream 
aborted").to_protobuf(request.mutable_cancel_status());
+        PCancelPlanFragmentResult response;
+        brpc::Controller controller;
+        controller.set_timeout_ms(1000);
+        stub->cancel_plan_fragment(&controller, &request, &response, nullptr);
+        if (controller.Failed()) {
+            Status::RpcError(controller.ErrorText()).to_thrift(&result);
+        } else {
+            Status::create(response.status()).to_thrift(&result);
+        }
+    }
+    TUniqueId buffer_id;
+    TUniqueId query_id;
+    std::shared_ptr<PBackendService_Stub> stub;
+    std::atomic<int> calls {0};
+    int rejections = 0;
+};
+
+class FlightCancelServerReady : public 
apache::thrift::server::TServerEventHandler {
+public:
+    void preServe() override { ready.set_value(); }
+    std::promise<void> ready;
+};
+
+class FlightCancelResultService : public PInternalService {
+public:
+    explicit FlightCancelResultService(ExecEnv* env) : PInternalService(env) {}
+
+    void cancel_plan_fragment(google::protobuf::RpcController* controller,
+                              const PCancelPlanFragmentRequest* request,
+                              PCancelPlanFragmentResult* result,
+                              google::protobuf::Closure* done) override {
+        if (!legacy_cancel) {
+            PInternalService::cancel_plan_fragment(controller, request, 
result, done);
+            return;
+        }
+        brpc::ClosureGuard guard(done);
+        // Emulate the pre-change handler: only the existing query 
cancellation RPC is understood.
+        
_exec_env->fragment_mgr()->cancel_query(UniqueId(request->query_id()).to_thrift(),
+                                                
Status::create(request->cancel_status()));
+        Status::OK().to_protobuf(result->mutable_status());
+    }
+    bool legacy_cancel = false;
+};
+
+class ArrowFlightRemoteCancelTest : public ArrowFlightCancelTest {
+protected:
+    void SetUp() override {
+        ArrowFlightCancelTest::SetUp();
+        static DNSCache dns_cache;
+        _previous_dns_cache = 
std::exchange(ExecEnv::GetInstance()->_dns_cache, &dns_cache);
+        _previous_cache = 
std::exchange(ExecEnv::GetInstance()->_internal_client_cache, &_cache);
+        auto heavy = std::exchange(config::brpc_heavy_work_pool_threads, 1);
+        auto light = std::exchange(config::brpc_light_work_pool_threads, 1);
+        auto flight = 
std::exchange(config::brpc_arrow_flight_work_pool_threads, 2);
+        _previous_load_mgr = 
std::move(ExecEnv::GetInstance()->_load_stream_mgr);
+        ExecEnv::GetInstance()->_load_stream_mgr = 
std::make_unique<LoadStreamMgr>(1);
+        _service = 
std::make_unique<FlightCancelResultService>(ExecEnv::GetInstance());
+        config::brpc_heavy_work_pool_threads = heavy;
+        config::brpc_light_work_pool_threads = light;
+        config::brpc_arrow_flight_work_pool_threads = flight;
+        ASSERT_EQ(_server.AddService(_service.get(), 
brpc::SERVER_DOESNT_OWN_SERVICE), 0);
+        ASSERT_EQ(_server.Start("127.0.0.1:0", nullptr), 0);
+        _statement->result_addr.hostname = "127.0.0.1";
+        _statement->result_addr.port = _server.listen_address().port;
+        _fragment_mgr =
+                std::make_unique<MockFragmentManager>(_cancel_status, 
ExecEnv::GetInstance());
+        _fragment_mgr->stop();
+        _previous_fragment_mgr =
+                std::exchange(ExecEnv::GetInstance()->_fragment_mgr, 
_fragment_mgr.get());
+        _frontend_cache = 
std::make_unique<ClientCache<FrontendServiceClient>>();
+        _previous_frontend_cache = 
std::exchange(ExecEnv::GetInstance()->_frontend_client_cache,
+                                                 _frontend_cache.get());
+        _coordinator = std::make_shared<FlightCancelCoordinator>();
+        _coordinator->buffer_id = _id;
+        _coordinator->query_id = _state.query_id();
+        _coordinator->stub = _cache.get_client(_statement->result_addr);
+        auto processor = 
std::make_shared<FrontendServiceProcessor>(_coordinator);
+        auto socket = 
std::make_shared<apache::thrift::transport::TServerSocket>("127.0.0.1", 0);
+        _fe_server = std::make_unique<apache::thrift::server::TThreadedServer>(
+                processor, socket,
+                
std::make_shared<apache::thrift::transport::TBufferedTransportFactory>(),
+                
std::make_shared<apache::thrift::protocol::TBinaryProtocolFactory>());
+        auto ready = std::make_shared<FlightCancelServerReady>();
+        _fe_server->setServerEventHandler(ready);
+        _fe_thread = std::thread([this] { _fe_server->serve(); });
+        ASSERT_EQ(ready->ready.get_future().wait_for(std::chrono::seconds(5)),
+                  std::future_status::ready);
+        TFrontendInfo frontend;
+        frontend.coordinator_address.hostname = "127.0.0.1";
+        frontend.coordinator_address.port = socket->getPort();
+        _previous_frontends = ExecEnv::GetInstance()->_frontends;
+        ExecEnv::GetInstance()->update_frontends({frontend});
+    }
+    void TearDown() override {
+        _frontend_cache.reset();
+        if (_fe_server) {
+            _fe_server->stop();
+            _fe_thread.join();
+        }
+        ExecEnv::GetInstance()->_dns_cache = _previous_dns_cache;
+        ExecEnv::GetInstance()->_frontend_client_cache = 
_previous_frontend_cache;
+        ExecEnv::GetInstance()->_frontends = _previous_frontends;
+        ExecEnv::GetInstance()->_fragment_mgr = _previous_fragment_mgr;
+        _buffer->cancel(Status::Cancelled("test cleanup"));
+        _server.Stop(0);
+        _server.Join();
+        _service.reset();
+        ExecEnv::GetInstance()->_load_stream_mgr = 
std::move(_previous_load_mgr);
+        ExecEnv::GetInstance()->_internal_client_cache = _previous_cache;
+        ArrowFlightCancelTest::TearDown();
+    }
+    BrpcClientCache<PBackendService_Stub> _cache;
+    BrpcClientCache<PBackendService_Stub>* _previous_cache = nullptr;
+    std::unique_ptr<LoadStreamMgr> _previous_load_mgr;
+    std::unique_ptr<FlightCancelResultService> _service;
+    brpc::Server _server;
+    DNSCache* _previous_dns_cache = nullptr;
+    Status _cancel_status;
+    std::unique_ptr<MockFragmentManager> _fragment_mgr;
+    FragmentMgr* _previous_fragment_mgr = nullptr;
+    std::unique_ptr<ClientCache<FrontendServiceClient>> _frontend_cache;
+    ClientCache<FrontendServiceClient>* _previous_frontend_cache = nullptr;
+    std::shared_ptr<FlightCancelCoordinator> _coordinator;
+    std::unique_ptr<apache::thrift::server::TThreadedServer> _fe_server;
+    std::thread _fe_thread;
+    std::map<TNetworkAddress, FrontendInfo> _previous_frontends;
+};
+
+TEST_F(ArrowFlightRemoteCancelTest, CloseReachesResultBackend) {
+    fill_buffer();
+    auto result = ArrowFlightBatchRemoteReader::Create(_statement);
+    ASSERT_TRUE(result.ok()) << result.status();
+    EXPECT_TRUE((*result)->Close().ok());
+    EXPECT_FALSE(registered());
+    EXPECT_TRUE(_dep->ready());
+    EXPECT_TRUE(_state.get_query_ctx()->is_cancelled());
+    EXPECT_TRUE((*result)->Close().ok());
+}
+
+TEST_F(ArrowFlightRemoteCancelTest, CancelBypassesRejectedArrowWork) {
+    fill_buffer();
+    auto reader = ArrowFlightBatchRemoteReader::Create(_statement);
+    ASSERT_TRUE(reader.ok()) << reader.status();
+    _service->_arrow_flight_work_pool.shutdown();
+    EXPECT_FALSE(_service->_arrow_flight_work_pool.try_offer([] {}));
+    ASSERT_TRUE((*reader)->Close().ok());
+    EXPECT_FALSE(registered());
+    EXPECT_TRUE(_dep->ready());
+    EXPECT_TRUE(_state.get_query_ctx()->is_cancelled());
+    EXPECT_EQ(_coordinator->calls, 1);
+}
+
+TEST_F(ArrowFlightRemoteCancelTest, RetryCoordinatorApplicationError) {
+    fill_buffer();
+    auto reader = ArrowFlightBatchRemoteReader::Create(_statement);
+    ASSERT_TRUE(reader.ok()) << reader.status();
+    _coordinator->rejections = 1;
+    ASSERT_TRUE((*reader)->Close().ok());
+    EXPECT_EQ(_coordinator->calls, 2);
+    EXPECT_FALSE(registered());
+    EXPECT_TRUE(_state.get_query_ctx()->is_cancelled());
+}
+
+TEST_F(ArrowFlightRemoteCancelTest, 
OlderResultBackendReceivesQueryCancellation) {
+    fill_buffer();
+    auto reader = ArrowFlightBatchRemoteReader::Create(_statement);
+    ASSERT_TRUE(reader.ok()) << reader.status();
+    _service->legacy_cancel = true;
+    ASSERT_TRUE((*reader)->Close().ok());
+    EXPECT_EQ(_coordinator->calls, 1);
+    EXPECT_FALSE(_cancel_status.ok());
+    // Older result BEs still own their historical buffer cleanup policy.
+    EXPECT_TRUE(registered());
+}
+
+TEST_F(ArrowFlightRemoteCancelTest, FinishedLocalEndpointStillReportsAbort) {
+    auto statement = std::make_shared<QueryStatement>(*_statement);
+    statement->result_addr.hostname = BackendOptions::get_localhost();
+    statement->result_addr.port = config::brpc_port;
+    fill_buffer();
+    bool fully_closed = false;
+    ASSERT_TRUE(_buffer->close(_state.fragment_instance_id(), Status::OK(), 2, 
fully_closed).ok());
+    auto reader = ArrowFlightBatchLocalReader::Create(statement);
+    ASSERT_TRUE(reader.ok()) << reader.status();
+    std::shared_ptr<arrow::RecordBatch> batch;
+    ASSERT_TRUE((*reader)->ReadNext(&batch).ok());
+    ASSERT_NE(batch, nullptr);
+    _state._query_ctx_uptr.reset();
+    _state._query_ctx = nullptr;
+    ASSERT_TRUE((*reader)->Close().ok());
+    EXPECT_EQ(_coordinator->calls, 1);
+    EXPECT_FALSE(_cancel_status.ok());
+    EXPECT_FALSE(registered());
+}
+
+TEST_F(ArrowFlightRemoteCancelTest, CancellationInterruptsPendingBrpcFetch) {
+    std::atomic<bool> cancelled = false;
+    auto result =
+            ArrowFlightBatchRemoteReader::Create(_statement, [&] { return 
cancelled.load(); });
+    ASSERT_TRUE(result.ok()) << result.status();
+    auto fetch = std::async(std::launch::async, [&] {
+        std::shared_ptr<arrow::RecordBatch> batch;
+        return (*result)->ReadNext(&batch);
+    });
+    // Wait until the server has queued the fetch; cancelling before that 
misses the blocking path.
+    bool waiting = false;
+    for (int i = 0; i < 200; ++i) {
+        {
+            std::lock_guard lock(_buffer->_lock);
+            waiting = !_buffer->_waiting_rpc.empty();
+        }
+        if (waiting) {
+            break;
+        }
+        std::this_thread::sleep_for(std::chrono::milliseconds(10));
+    }
+    EXPECT_TRUE(waiting);
+    cancelled = true;
+    auto ready = fetch.wait_for(std::chrono::seconds(3));
+    if (ready != std::future_status::ready) {
+        _buffer->cancel(Status::Cancelled("test cleanup"));
+    }
+    EXPECT_EQ(ready, std::future_status::ready);
+    EXPECT_FALSE(fetch.get().ok());
+    EXPECT_FALSE(registered());
+    EXPECT_TRUE(_state.get_query_ctx()->is_cancelled());
+}
+
+TEST_F(ArrowFlightRemoteCancelTest, FetchAfterCancellationStillCompletesRpc) {
+    auto reader = ArrowFlightBatchRemoteReader::Create(_statement);
+    ASSERT_TRUE(reader.ok()) << reader.status();
+    ASSERT_TRUE((*reader)->Close().ok());
+    auto stub = _cache.get_client(_statement->result_addr);
+    ASSERT_NE(stub, nullptr);
+    PFetchArrowDataRequest request;
+    request.mutable_finst_id()->set_hi(_id.hi);
+    request.mutable_finst_id()->set_lo(_id.lo);
+    PFetchArrowDataResult response;
+    brpc::Controller controller;
+    controller.set_timeout_ms(1000);
+    stub->fetch_arrow_data(&controller, &request, &response, nullptr);
+    EXPECT_FALSE(controller.Failed()) << controller.ErrorText();
+    EXPECT_TRUE(response.has_status());
+    EXPECT_FALSE(Status::create(response.status()).ok());
+}
+
+TEST_F(ArrowFlightRemoteCancelTest, 
SuccessfulFragmentStopPreservesUnreadResults) {
+    fill_buffer();
+    auto stub = _cache.get_client(_statement->result_addr);
+    for (const auto& reason : {Status::Error<ErrorCode::LIMIT_REACH>("limit 
reached"),
+                               
Status::Error<ErrorCode::FINISHED>("finished")}) {
+        PCancelPlanFragmentRequest request;
+        request.mutable_finst_id()->set_hi(_id.hi);
+        request.mutable_finst_id()->set_lo(_id.lo);
+        request.mutable_query_id()->set_hi(_state.query_id().hi);
+        request.mutable_query_id()->set_lo(_state.query_id().lo);
+        reason.to_protobuf(request.mutable_cancel_status());
+        PCancelPlanFragmentResult response;
+        brpc::Controller controller;
+        controller.set_timeout_ms(1000);
+        stub->cancel_plan_fragment(&controller, &request, &response, nullptr);
+        ASSERT_FALSE(controller.Failed()) << controller.ErrorText();
+        EXPECT_TRUE(Status::create(response.status()).ok());
+        EXPECT_TRUE(registered());
+        EXPECT_FALSE(_buffer->_result_batch_queue.empty());
+    }
+}
+
+TEST_F(ArrowFlightRemoteCancelTest, NormalEofDoesNotCancelQuery) {
+    bool fully_closed = false;
+    ASSERT_TRUE(_buffer->close(_state.fragment_instance_id(), Status::OK(), 0, 
fully_closed).ok());
+    auto result = ArrowFlightBatchRemoteReader::Create(_statement);
+    ASSERT_TRUE(result.ok()) << result.status();
+    std::shared_ptr<arrow::RecordBatch> batch;
+    EXPECT_TRUE((*result)->ReadNext(&batch).ok());
+    EXPECT_EQ(batch, nullptr);
+    EXPECT_TRUE((*result)->Close().ok());
+    EXPECT_FALSE(_state.get_query_ctx()->is_cancelled());
+}
+
+} // namespace doris::flight
diff --git a/fe/fe-core/src/main/java/org/apache/doris/qe/Coordinator.java 
b/fe/fe-core/src/main/java/org/apache/doris/qe/Coordinator.java
index c582dccf7f9..c2019d05d26 100644
--- a/fe/fe-core/src/main/java/org/apache/doris/qe/Coordinator.java
+++ b/fe/fe-core/src/main/java/org/apache/doris/qe/Coordinator.java
@@ -3828,6 +3828,11 @@ public class Coordinator implements CoordInterface {
         return result;
     }
 
+    public List<TNetworkAddress> getBackendBrpcAddresses() {
+        return beToPipelineExecCtxs.values().stream().map(ctx -> 
ctx.brpcAddr.deepCopy())
+                .collect(Collectors.toList());
+    }
+
     @Override
     public List<TNetworkAddress> getInvolvedBackends() {
         List<TNetworkAddress> backendAddresses = Lists.newArrayList();
diff --git 
a/fe/fe-core/src/main/java/org/apache/doris/qe/NereidsCoordinator.java 
b/fe/fe-core/src/main/java/org/apache/doris/qe/NereidsCoordinator.java
index 146df08108c..2dc02abc697 100644
--- a/fe/fe-core/src/main/java/org/apache/doris/qe/NereidsCoordinator.java
+++ b/fe/fe-core/src/main/java/org/apache/doris/qe/NereidsCoordinator.java
@@ -359,6 +359,12 @@ public class NereidsCoordinator extends Coordinator {
         return 
coordinatorContext.asLoadProcessor().loadContext.getErrorTabletInfos();
     }
 
+    @Override
+    public List<TNetworkAddress> getBackendBrpcAddresses() {
+        return executionTask.getChildrenTasks().values().stream()
+                .map(task -> 
task.getBackend().getBrpcAddress().deepCopy()).collect(Collectors.toList());
+    }
+
     @Override
     public List<TNetworkAddress> getInvolvedBackends() {
         return 
Utils.fastToImmutableList(coordinatorContext.backends.get().keySet());
diff --git a/fe/fe-core/src/main/java/org/apache/doris/qe/StmtExecutor.java 
b/fe/fe-core/src/main/java/org/apache/doris/qe/StmtExecutor.java
index b349cf9160d..30a943bba1a 100644
--- a/fe/fe-core/src/main/java/org/apache/doris/qe/StmtExecutor.java
+++ b/fe/fe-core/src/main/java/org/apache/doris/qe/StmtExecutor.java
@@ -1456,6 +1456,12 @@ public class StmtExecutor {
             if (context.getConnectType().equals(ConnectType.ARROW_FLIGHT_SQL)) 
{
                 Preconditions.checkState(!context.isReturnResultFromLocal());
                 profile.getSummaryProfile().setTempStartTime();
+                if (coordBase == coord) {
+                    
context.getFlightSqlChannel().registerRemoteQuery(coord.getQueryId(),
+                            context.getFlightSqlEndpointsLocations().stream()
+                                    .map(endpoint -> 
endpoint.getFinstId()).distinct().collect(Collectors.toList()),
+                            coord.getBackendBrpcAddresses(), 
coord.getQueryOptions().getExecutionTimeout());
+                }
                 // The client pulls the results from the BE later (DoGet). 
Only an external-table
                 // scan in batch mode still needs the coordinator after this 
point: the BE fetches
                 // its splits lazily from the split source the coordinator 
holds, so closing the
diff --git 
a/fe/fe-core/src/main/java/org/apache/doris/service/FrontendServiceImpl.java 
b/fe/fe-core/src/main/java/org/apache/doris/service/FrontendServiceImpl.java
index 54101ad5640..14f38a680e4 100644
--- a/fe/fe-core/src/main/java/org/apache/doris/service/FrontendServiceImpl.java
+++ b/fe/fe-core/src/main/java/org/apache/doris/service/FrontendServiceImpl.java
@@ -141,6 +141,7 @@ import org.apache.doris.qe.VariableMgr;
 import org.apache.doris.resource.BackendSelection;
 import org.apache.doris.resource.BackendSelectionManager;
 import org.apache.doris.service.arrowflight.FlightSqlConnectProcessor;
+import org.apache.doris.service.arrowflight.FlightSqlQueryCancellation;
 import org.apache.doris.statistics.AnalysisManager;
 import org.apache.doris.statistics.ColStatsData;
 import org.apache.doris.statistics.ColumnStatistic;
@@ -1092,6 +1093,11 @@ public class FrontendServiceImpl implements 
FrontendService.Iface {
         return QeProcessorImpl.INSTANCE.reportExecStatus(params, 
getClientAddr());
     }
 
+    @Override
+    public TStatus cancelFlightQuery(TUniqueId resultId) throws TException {
+        return FlightSqlQueryCancellation.INSTANCE.cancel(resultId);
+    }
+
     @Override
     public TFetchSplitBatchResult fetchSplitBatch(TFetchSplitBatchRequest 
request) throws TException {
         TFetchSplitBatchResult result = new TFetchSplitBatchResult();
diff --git 
a/fe/fe-core/src/main/java/org/apache/doris/service/arrowflight/FlightSqlQueryCancellation.java
 
b/fe/fe-core/src/main/java/org/apache/doris/service/arrowflight/FlightSqlQueryCancellation.java
new file mode 100644
index 00000000000..4fb006e774c
--- /dev/null
+++ 
b/fe/fe-core/src/main/java/org/apache/doris/service/arrowflight/FlightSqlQueryCancellation.java
@@ -0,0 +1,146 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements.  See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership.  The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License.  You may obtain a copy of the License at
+//
+//   http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied.  See the License for the
+// specific language governing permissions and limitations
+// under the License.
+
+package org.apache.doris.service.arrowflight;
+
+import org.apache.doris.common.Status;
+import org.apache.doris.proto.InternalService.PCancelPlanFragmentResult;
+import org.apache.doris.rpc.BackendServiceProxy;
+import org.apache.doris.thrift.TNetworkAddress;
+import org.apache.doris.thrift.TStatus;
+import org.apache.doris.thrift.TStatusCode;
+import org.apache.doris.thrift.TUniqueId;
+
+import com.github.benmanes.caffeine.cache.Cache;
+import com.github.benmanes.caffeine.cache.Caffeine;
+import com.github.benmanes.caffeine.cache.Expiry;
+import com.github.benmanes.caffeine.cache.Scheduler;
+import com.github.benmanes.caffeine.cache.Ticker;
+import org.apache.logging.log4j.LogManager;
+import org.apache.logging.log4j.Logger;
+
+import java.util.ArrayList;
+import java.util.List;
+import java.util.concurrent.Future;
+import java.util.concurrent.TimeUnit;
+import java.util.stream.Collectors;
+
+public final class FlightSqlQueryCancellation {
+    private static final Logger LOG = 
LogManager.getLogger(FlightSqlQueryCancellation.class);
+    public static final FlightSqlQueryCancellation INSTANCE = new 
FlightSqlQueryCancellation(Ticker.systemTicker());
+    private final Cache<TUniqueId, Route> results;
+
+    FlightSqlQueryCancellation(Ticker ticker) {
+        results = 
Caffeine.newBuilder().ticker(ticker).scheduler(Scheduler.systemScheduler())
+                .expireAfter(new Expiry<TUniqueId, Route>() {
+                    @Override
+                    public long expireAfterCreate(TUniqueId key, Route route, 
long now) {
+                        return route.ttlNanos;
+                    }
+
+                    @Override
+                    public long expireAfterUpdate(TUniqueId key, Route route, 
long now, long duration) {
+                        return route.ttlNanos;
+                    }
+
+                    @Override
+                    public long expireAfterRead(TUniqueId key, Route route, 
long now, long duration) {
+                        return duration;
+                    }
+                }).build();
+    }
+
+    public void register(TUniqueId queryId, List<TUniqueId> resultIds,
+            List<TNetworkAddress> backends, int timeoutSeconds) {
+        if (backends.isEmpty()) {
+            throw new IllegalArgumentException("Flight query has no 
cancellation backends");
+        }
+        // Keep only cancellation addresses, not the coordinator, scan state, 
or query queue slot.
+        // Ordinary Flight coordinators are unregistered when GetFlightInfo 
returns.
+        Route route = new Route(queryId.deepCopy(),
+                
resultIds.stream().map(TUniqueId::deepCopy).distinct().collect(Collectors.toList()),
+                
backends.stream().map(TNetworkAddress::deepCopy).distinct().collect(Collectors.toList()),
+                TimeUnit.SECONDS.toNanos(Math.max(0L, timeoutSeconds) + 5));
+        for (TUniqueId resultId : route.resultIds) {
+            results.put(resultId, route);
+        }
+    }
+
+    public void unregister(List<TUniqueId> resultIds) {
+        results.invalidateAll(resultIds);
+    }
+
+    public TStatus cancel(TUniqueId resultId) {
+        Route route = results.getIfPresent(resultId);
+        if (route == null) {
+            return new TStatus(TStatusCode.NOT_FOUND);
+        }
+        Status reason = new Status(TStatusCode.CANCELLED, "Arrow Flight stream 
closed before EOF");
+        TStatus status = new TStatus(TStatusCode.OK);
+        List<Future<PCancelPlanFragmentResult>> futures = new ArrayList<>();
+        long deadline = System.nanoTime() + TimeUnit.SECONDS.toNanos(1);
+        for (TNetworkAddress backend : route.backends) {
+            try {
+                // The existing RPC is understood by older result BEs and 
bypasses the Flight read pool.
+                futures.add(BackendServiceProxy.getInstance()
+                        .cancelPipelineXPlanFragmentAsync(backend, 
route.queryId, reason));
+            } catch (Exception e) {
+                LOG.warn("Failed to send Flight cancellation for {} to {}", 
route.queryId, backend, e);
+                status = new TStatus(TStatusCode.INTERNAL_ERROR);
+            }
+        }
+        for (Future<PCancelPlanFragmentResult> future : futures) {
+            try {
+                PCancelPlanFragmentResult response = future.get(Math.max(0L, 
deadline - System.nanoTime()),
+                        TimeUnit.NANOSECONDS);
+                if (!response.hasStatus()) {
+                    status = new TStatus(TStatusCode.INTERNAL_ERROR);
+                } else if (response.getStatus().getStatusCode() != 
TStatusCode.OK.getValue()) {
+                    status = new TStatus(new 
Status(response.getStatus()).getErrorCode());
+                }
+            } catch (InterruptedException e) {
+                Thread.currentThread().interrupt();
+                return new TStatus(TStatusCode.CANCELLED);
+            } catch (Exception e) {
+                LOG.warn("Failed to complete Flight cancellation for {}", 
route.queryId, e);
+                status = new TStatus(TStatusCode.INTERNAL_ERROR);
+            }
+        }
+        if (status.getStatusCode() == TStatusCode.OK) {
+            // An unsuccessful attempt must remain routable for the caller's 
bounded retry.
+            for (TUniqueId id : route.resultIds) {
+                results.asMap().remove(id, route);
+            }
+        }
+        return status;
+    }
+
+    private static final class Route {
+        private final TUniqueId queryId;
+        private final List<TUniqueId> resultIds;
+        private final List<TNetworkAddress> backends;
+        private final long ttlNanos;
+
+        private Route(TUniqueId queryId, List<TUniqueId> resultIds,
+                List<TNetworkAddress> backends, long ttlNanos) {
+            this.queryId = queryId;
+            this.resultIds = resultIds;
+            this.backends = backends;
+            this.ttlNanos = ttlNanos;
+        }
+    }
+}
diff --git 
a/fe/fe-core/src/main/java/org/apache/doris/service/arrowflight/results/FlightSqlChannel.java
 
b/fe/fe-core/src/main/java/org/apache/doris/service/arrowflight/results/FlightSqlChannel.java
index 2781994dfa1..6fe8369bbe4 100644
--- 
a/fe/fe-core/src/main/java/org/apache/doris/service/arrowflight/results/FlightSqlChannel.java
+++ 
b/fe/fe-core/src/main/java/org/apache/doris/service/arrowflight/results/FlightSqlChannel.java
@@ -21,6 +21,9 @@ import org.apache.doris.catalog.Column;
 import org.apache.doris.common.FeConstants;
 import org.apache.doris.qe.ResultSet;
 import org.apache.doris.qe.ResultSetMetaData;
+import org.apache.doris.service.arrowflight.FlightSqlQueryCancellation;
+import org.apache.doris.thrift.TNetworkAddress;
+import org.apache.doris.thrift.TUniqueId;
 
 import com.google.common.cache.Cache;
 import com.google.common.cache.CacheBuilder;
@@ -44,6 +47,7 @@ import java.util.concurrent.TimeUnit;
 public class FlightSqlChannel {
     private final Cache<String, FlightSqlResultCacheEntry> resultCache;
     private final BufferAllocator allocator;
+    private final List<TUniqueId> remoteResultIds = new ArrayList<>();
 
     public FlightSqlChannel() {
         // The Stmt result is not picked up by the Client within 10 minutes 
and will be deleted.
@@ -162,8 +166,16 @@ public class FlightSqlChannel {
         return allocator.getAllocatedMemory();
     }
 
-    public void reset() {
+    public synchronized void registerRemoteQuery(TUniqueId queryId, 
List<TUniqueId> resultIds,
+            List<TNetworkAddress> backends, int timeoutSeconds) {
+        FlightSqlQueryCancellation.INSTANCE.register(queryId, resultIds, 
backends, timeoutSeconds);
+        remoteResultIds.addAll(resultIds);
+    }
+
+    public synchronized void reset() {
         resultCache.invalidateAll();
+        FlightSqlQueryCancellation.INSTANCE.unregister(remoteResultIds);
+        remoteResultIds.clear();
     }
 
     public void close() {
diff --git 
a/fe/fe-core/src/test/java/org/apache/doris/service/arrowflight/FlightSqlQueryCancellationTest.java
 
b/fe/fe-core/src/test/java/org/apache/doris/service/arrowflight/FlightSqlQueryCancellationTest.java
new file mode 100644
index 00000000000..b6a8d9362cf
--- /dev/null
+++ 
b/fe/fe-core/src/test/java/org/apache/doris/service/arrowflight/FlightSqlQueryCancellationTest.java
@@ -0,0 +1,109 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements.  See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership.  The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License.  You may obtain a copy of the License at
+//
+//   http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied.  See the License for the
+// specific language governing permissions and limitations
+// under the License.
+
+package org.apache.doris.service.arrowflight;
+
+import org.apache.doris.proto.InternalService.PCancelPlanFragmentResult;
+import org.apache.doris.proto.Types.PStatus;
+import org.apache.doris.rpc.BackendServiceProxy;
+import org.apache.doris.service.arrowflight.results.FlightSqlChannel;
+import org.apache.doris.thrift.TNetworkAddress;
+import org.apache.doris.thrift.TStatusCode;
+import org.apache.doris.thrift.TUniqueId;
+
+import com.google.common.collect.ImmutableList;
+import com.google.common.util.concurrent.Futures;
+import org.junit.Assert;
+import org.junit.Test;
+import org.mockito.MockedStatic;
+import org.mockito.Mockito;
+
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.atomic.AtomicLong;
+
+public class FlightSqlQueryCancellationTest {
+    private static final TUniqueId QUERY = new TUniqueId(10L, 20L);
+    private static final TUniqueId FIRST = new TUniqueId(10L, 21L);
+    private static final TUniqueId SECOND = new TUniqueId(10L, 22L);
+    private static final TNetworkAddress FINISHED_BE = new 
TNetworkAddress("127.0.0.1", 8001);
+    private static final TNetworkAddress ACTIVE_BE = new 
TNetworkAddress("127.0.0.1", 8002);
+
+    private static PCancelPlanFragmentResult response(TStatusCode status) {
+        return PCancelPlanFragmentResult.newBuilder()
+                
.setStatus(PStatus.newBuilder().setStatusCode(status.getValue())).build();
+    }
+
+    @Test
+    public void 
cancelReachesAllBackendsWithoutLiveCoordinatorOrResultFragment() throws 
Exception {
+        FlightSqlQueryCancellation routes = new FlightSqlQueryCancellation(() 
-> 0L);
+        // Use Guava collection factories because FE tests compile against the 
Java 8 API.
+        // The selected endpoint and its coordinator have finished; only 
routing metadata survives.
+        routes.register(QUERY, ImmutableList.of(FIRST, SECOND), 
ImmutableList.of(FINISHED_BE, ACTIVE_BE), 120);
+        BackendServiceProxy proxy = Mockito.mock(BackendServiceProxy.class);
+        Mockito.when(proxy.cancelPipelineXPlanFragmentAsync(Mockito.any(), 
Mockito.eq(QUERY), Mockito.any()))
+                .thenReturn(Futures.immediateFuture(response(TStatusCode.OK)));
+        try (MockedStatic<BackendServiceProxy> mocked = 
Mockito.mockStatic(BackendServiceProxy.class)) {
+            mocked.when(BackendServiceProxy::getInstance).thenReturn(proxy);
+            Assert.assertEquals(TStatusCode.OK, 
routes.cancel(FIRST).getStatusCode());
+            
Mockito.verify(proxy).cancelPipelineXPlanFragmentAsync(Mockito.eq(FINISHED_BE),
+                    Mockito.eq(QUERY), Mockito.any());
+            
Mockito.verify(proxy).cancelPipelineXPlanFragmentAsync(Mockito.eq(ACTIVE_BE),
+                    Mockito.eq(QUERY), Mockito.any());
+        }
+        Assert.assertEquals(TStatusCode.NOT_FOUND, 
routes.cancel(SECOND).getStatusCode());
+    }
+
+    @Test
+    public void applicationErrorKeepsRouteForRetry() throws Exception {
+        FlightSqlQueryCancellation routes = new FlightSqlQueryCancellation(() 
-> 0L);
+        routes.register(QUERY, ImmutableList.of(FIRST), 
ImmutableList.of(ACTIVE_BE), 120);
+        BackendServiceProxy proxy = Mockito.mock(BackendServiceProxy.class);
+        Mockito.when(proxy.cancelPipelineXPlanFragmentAsync(Mockito.any(), 
Mockito.eq(QUERY), Mockito.any()))
+                
.thenReturn(Futures.immediateFuture(response(TStatusCode.CANCELLED)),
+                        Futures.immediateFuture(response(TStatusCode.OK)));
+        try (MockedStatic<BackendServiceProxy> mocked = 
Mockito.mockStatic(BackendServiceProxy.class)) {
+            mocked.when(BackendServiceProxy::getInstance).thenReturn(proxy);
+            Assert.assertEquals(TStatusCode.CANCELLED, 
routes.cancel(FIRST).getStatusCode());
+            Assert.assertEquals(TStatusCode.OK, 
routes.cancel(FIRST).getStatusCode());
+        }
+    }
+
+    @Test
+    public void routesExpireAfterExecutionTimeout() {
+        AtomicLong clock = new AtomicLong();
+        FlightSqlQueryCancellation routes = new 
FlightSqlQueryCancellation(clock::get);
+        routes.register(QUERY, ImmutableList.of(FIRST, SECOND), 
ImmutableList.of(ACTIVE_BE), 120);
+        clock.set(TimeUnit.SECONDS.toNanos(126));
+        Assert.assertEquals(TStatusCode.NOT_FOUND, 
routes.cancel(FIRST).getStatusCode());
+        Assert.assertEquals(TStatusCode.NOT_FOUND, 
routes.cancel(SECOND).getStatusCode());
+    }
+
+    @Test
+    public void channelCleanupRemovesItsRoutes() {
+        FlightSqlChannel channel = new FlightSqlChannel();
+        try {
+            channel.registerRemoteQuery(QUERY, ImmutableList.of(FIRST, 
SECOND), ImmutableList.of(ACTIVE_BE), 120);
+            channel.reset();
+            Assert.assertEquals(TStatusCode.NOT_FOUND,
+                    
FlightSqlQueryCancellation.INSTANCE.cancel(FIRST).getStatusCode());
+            Assert.assertEquals(TStatusCode.NOT_FOUND,
+                    
FlightSqlQueryCancellation.INSTANCE.cancel(SECOND).getStatusCode());
+        } finally {
+            channel.close();
+        }
+    }
+}
diff --git a/gensrc/thrift/FrontendService.thrift 
b/gensrc/thrift/FrontendService.thrift
index 24889633d3d..4e753eb859b 100644
--- a/gensrc/thrift/FrontendService.thrift
+++ b/gensrc/thrift/FrontendService.thrift
@@ -1934,4 +1934,7 @@ service FrontendService {
     TGetOlapTableMetaResult getOlapTableMeta(1: TGetOlapTableMetaRequest 
request)
 
     Status.TStatus syncCloudTabletStats(1: TSyncCloudTabletStatsRequest 
request)
+
+    // Result buffers can outlive their local fragments; abort through the 
owning FE's query-level route.
+    Status.TStatus cancelFlightQuery(1: Types.TUniqueId result_id)
 }
diff --git 
a/regression-test/suites/arrow_flight_sql_p0/test_flight_cancel_cleanup.groovy 
b/regression-test/suites/arrow_flight_sql_p0/test_flight_cancel_cleanup.groovy
new file mode 100644
index 00000000000..aff57aeac4d
--- /dev/null
+++ 
b/regression-test/suites/arrow_flight_sql_p0/test_flight_cancel_cleanup.groovy
@@ -0,0 +1,128 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements.  See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership.  The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License.  You may obtain a copy of the License at
+//
+//   http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied.  See the License for the
+// specific language governing permissions and limitations
+// under the License.
+
+import org.apache.arrow.driver.jdbc.shaded.com.google.protobuf.Any
+import org.apache.arrow.driver.jdbc.shaded.org.apache.arrow.flight.FlightClient
+import org.apache.arrow.driver.jdbc.shaded.org.apache.arrow.flight.Location
+import 
org.apache.arrow.driver.jdbc.shaded.org.apache.arrow.flight.sql.FlightSqlClient
+import 
org.apache.arrow.driver.jdbc.shaded.org.apache.arrow.flight.sql.impl.FlightSql
+import 
org.apache.arrow.driver.jdbc.shaded.org.apache.arrow.memory.RootAllocator
+
+suite("test_flight_cancel_cleanup", "arrow_flight_sql") {
+    def frontend = jdbc_sql_return_maparray("SHOW FRONTENDS").find {
+        it.IsMaster.toString().equalsIgnoreCase("true") && 
it.Alive.toString().equalsIgnoreCase("true")
+    }
+    assertNotNull(frontend)
+    assertTrue(frontend.ArrowFlightSqlPort.toString().toInteger() > 0)
+    def backends = jdbc_sql_return_maparray("SHOW BACKENDS").findAll {
+        it.Alive.toString().equalsIgnoreCase("true")
+    }
+    assertTrue(!backends.isEmpty())
+    def database = jdbc_sql("SELECT DATABASE()")[0][0]
+    def table = "${database}.flight_cancel_cleanup_source"
+    def allocator = new RootAllocator(Long.MAX_VALUE)
+    def feClient = FlightClient.builder(allocator,
+            Location.forGrpcInsecure(frontend.Host.toString(), 
frontend.ArrowFlightSqlPort.toString().toInteger())).build()
+    def client = new FlightSqlClient(feClient)
+    def auth
+    def consume = { String query ->
+        def info = client.execute(query, auth)
+        def rows = 0
+        info.endpoints.each { endpoint ->
+            FlightClient.builder(allocator, 
endpoint.locations[0]).build().withCloseable { beClient ->
+                beClient.getStream(endpoint.ticket, auth).withCloseable { 
stream ->
+                    while (stream.next()) {
+                        rows += stream.root.rowCount
+                    }
+                }
+            }
+        }
+        rows
+    }
+    try {
+        auth = 
feClient.authenticateBasicToken(context.config.otherConfigs.get("extArrowFlightSqlUser"),
+                
context.config.otherConfigs.get("extArrowFlightSqlPassword")).get()
+        consume("SET enable_sql_cache=false")
+        // The ticket is a query id in parallel mode, so inspect this query 
instead of global task counts.
+        consume("SET enable_parallel_result_sink=true")
+        consume("SET query_timeout=120")
+        jdbc_sql("DROP TABLE IF EXISTS ${table}")
+        jdbc_sql("CREATE TABLE ${table} (id BIGINT) DISTRIBUTED BY HASH(id) 
BUCKETS 8 " +
+                "PROPERTIES(\"replication_num\"=\"1\")")
+        jdbc_sql("INSERT INTO ${table} SELECT number FROM 
numbers(\"number\"=\"1024\")")
+        def tabletBackends = jdbc_sql_return_maparray("SHOW TABLETS FROM 
${table}")
+                .collect { it.BackendId }.unique().size()
+        [false, true].each { explicitCancel ->
+            def info = client.execute("SELECT id, n FROM ${table} " +
+                    "LATERAL VIEW explode_numbers(100000000) expanded AS n", 
auth)
+            if (tabletBackends > 1) {
+                assertTrue(info.endpoints.size() > 1, "Expected distributed 
Flight result endpoints")
+            }
+            def ids = info.endpoints.collect { endpoint ->
+                
Any.parseFrom(endpoint.ticket.bytes).unpack(FlightSql.TicketStatementQuery.class)
+                        .statementHandle.toStringUtf8().split("&")[0]
+            }.unique()
+            // Abort only one endpoint; the cancellation must reach every 
participating BE.
+            [info.endpoints[0]].each { endpoint ->
+                FlightClient.builder(allocator, 
endpoint.locations[0]).build().withCloseable { beClient ->
+                    beClient.getStream(endpoint.ticket, auth).withCloseable { 
stream ->
+                        assertTrue(stream.next())
+                        assertTrue(stream.root.rowCount > 0)
+                        if (explicitCancel) {
+                            stream.cancel("test cancellation", null)
+                        }
+                        // Closing after one batch must also cancel the 
unfinished producer.
+                    }
+                }
+            }
+            def deadline = System.nanoTime() + 
java.util.concurrent.TimeUnit.SECONDS.toNanos(20)
+            def remaining = []
+            while (true) {
+                remaining = []
+                backends.each { backend ->
+                    ids.each { id ->
+                        def body = new 
URL("http://${backend.Host}:${backend.HttpPort}/api/query_pipeline_tasks/${id}";)
+                                .getText(connectTimeout: 3000, readTimeout: 
3000)
+                        if (!body.contains("not found")) {
+                            remaining.add("${id}: ${body}")
+                        }
+                    }
+                }
+                if (remaining.isEmpty() || System.nanoTime() >= deadline) {
+                    break
+                }
+                sleep(100)
+            }
+            assertTrue(remaining.isEmpty(), "Cancelled Flight query retained 
pipelines: ${remaining}")
+            assertEquals(1, consume("SELECT 1 AS healthy"))
+        }
+        assertEquals(10, consume("SELECT number FROM 
numbers(\"number\"=\"10\")"))
+    } finally {
+        try {
+            if (auth != null) {
+                feClient.closeSession(new 
org.apache.arrow.driver.jdbc.shaded.org.apache.arrow.flight.CloseSessionRequest(),
 auth)
+            }
+        } finally {
+            try {
+                client.close()
+                allocator.close()
+            } finally {
+                jdbc_sql("DROP TABLE IF EXISTS ${table}")
+            }
+        }
+    }
+}


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

Reply via email to