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]