FlightSqlQueryCancellation.java

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