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


##########
be/src/service/arrow_flight/arrow_flight_batch_reader.cpp:
##########
@@ -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()) {

Review Comment:
   [P2] Bound FE discovery before cancelling the owning query. This loop visits 
the address-sorted running-FE map synchronously; a stale FE before the owner 
can trigger three connection attempts plus backoff on each of two passes, even 
though the constructor is given 2000 ms. A reused Thrift client can also retain 
the default 60000 ms receive timeout because that argument is applied only when 
a client is created. The Flight stream has already closed while the other BEs 
keep producing or block on buffers. Use an owner-directed route or one short 
total deadline, including cached-client RPCs.



##########
fe/fe-core/src/main/java/org/apache/doris/service/arrowflight/results/FlightSqlChannel.java:
##########
@@ -162,8 +166,16 @@ public long getAllocatedMemory() {
         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);

Review Comment:
   [P2] Keep cancellation routes for earlier BE tickets across session reset. 
`executeQueryStatement` calls `reset()` before every new statement, but a 
previous FlightInfo's tickets remain readable directly from BE. If a client 
obtains query A's parallel endpoints, starts query B, then closes one A stream 
early, this invalidates A's result IDs here; `cancelFlightQuery` returns 
NOT_FOUND and A's other BEs can remain blocked until timeout. Keep each route 
until its own cancellation or execution-time expiry, and cancel outstanding 
routes on session teardown. This is distinct from the prior finished-fragment 
thread: the replacement FE route itself is removed here.



##########
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; }

Review Comment:
   [P2] Scope the global metrics hook to this fixture's manager. Each case 
constructs a `ResultBufferMgr`, whose constructor registers 
`result_buffer_block_count` with a callback capturing `this`; its default 
destructor does not unregister it. `TearDown()` only restores 
`ExecEnv::_result_mgr`, so a filtered run of this new suite leaves a hook 
pointing into the first destroyed fixture. A later metrics scrape can access 
freed storage, and registration by later cases uses `emplace` rather than 
replacing that hook. Remove or restore only the hook this fixture owns before 
`_mgr` is destroyed; preserve any earlier manager's registration.



##########
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));

Review Comment:
   [P2] Ensure abort reaches a BE when its light pool rejects cancellation. 
This new fanout calls `cancel_plan_fragment`, which enters the bounded 
`_light_work_pool`; on `try_offer` failure the BE returns a protobuf CANCELLED 
before cancelling the query or its Flight buffers. The FE retains this route, 
but the closing reader makes only two immediate attempts and sets `_closed`, so 
under sustained light-pool pressure no later caller retries and a producer can 
remain blocked until timeout. Give this cancel reliable admission or retain a 
retry owner until every BE acknowledges it. This is a different pool and RPC 
from the earlier Arrow fetch-pool thread.



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