RowBinlogTtlDiscovery.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.binlog;

import org.apache.doris.catalog.Database;
import org.apache.doris.catalog.Env;
import org.apache.doris.catalog.MaterializedIndex;
import org.apache.doris.catalog.OlapTable;
import org.apache.doris.catalog.Partition;
import org.apache.doris.catalog.Table;
import org.apache.doris.catalog.Tablet;
import org.apache.doris.cloud.catalog.CloudReplica;
import org.apache.doris.cloud.qe.ComputeGroupException;
import org.apache.doris.cloud.system.CloudSystemInfoService;
import org.apache.doris.common.Config;
import org.apache.doris.common.util.MasterDaemon;
import org.apache.doris.proto.InternalService;
import org.apache.doris.rpc.BackendServiceProxy;
import org.apache.doris.rpc.RpcException;
import org.apache.doris.system.Backend;
import org.apache.doris.thrift.TStatusCode;

import com.google.common.util.concurrent.FutureCallback;
import com.google.common.util.concurrent.Futures;
import com.google.common.util.concurrent.MoreExecutors;
import org.apache.logging.log4j.LogManager;
import org.apache.logging.log4j.Logger;

import java.util.ArrayList;
import java.util.Collections;
import java.util.HashMap;
import java.util.Iterator;
import java.util.List;
import java.util.Map;
import java.util.concurrent.TimeUnit;

/** Bounded catalog discovery for tablets absent from a BE's metadata cache. */
public final class RowBinlogTtlDiscovery extends MasterDaemon {
    private static final Logger LOG = LogManager.getLogger(RowBinlogTtlDiscovery.class);
    private static final int MAX_TABLETS_PER_ROUND = 1024;
    private static final int MAX_TABLETS_PER_BACKEND = 64;
    private Iterator<Database> databases = Collections.emptyIterator();
    private Iterator<Table> tables = Collections.emptyIterator();
    private Iterator<Partition> partitions = Collections.emptyIterator();
    private Iterator<Tablet> tablets = Collections.emptyIterator();
    private Iterator<MaterializedIndex> indexes = Collections.emptyIterator();
    private OlapTable currentTable;

    public RowBinlogTtlDiscovery() {
        super("row-binlog-ttl-discovery", 100);
    }

    @Override
    protected void runAfterCatalogReady() {
        discover();
    }

    public void discover() {
        if (!Config.isCloudMode() || !Config.enable_feature_binlog) {
            setInterval(15000);
            return;
        }
        try {
            CloudSystemInfoService info = (CloudSystemInfoService) Env.getCurrentSystemInfo();
            Map<Long, List<Long>> targets = new HashMap<>();
            long deadline = System.nanoTime() + TimeUnit.MILLISECONDS.toNanos(50);
            if (!databases.hasNext() && !tables.hasNext() && !partitions.hasNext()
                    && !indexes.hasNext() && !tablets.hasNext()) {
                databases = Env.getCurrentInternalCatalog().getDbs().iterator();
            }
            for (int n = 0; n < MAX_TABLETS_PER_ROUND && System.nanoTime() < deadline; n++) {
                Tablet tablet = nextTablet(deadline);
                if (tablet == null) {
                    break;
                }
                CloudReplica replica = (CloudReplica) tablet.getReplicas().get(0);
                boolean batchFull = false;
                // One existing replica owner per compute group. BE read/write separation
                // decides which group may compact; Meta Service arbitrates competing jobs.
                for (String clusterId : info.getCloudClusterIds()) {
                    long backendId;
                    try {
                        backendId = replica.getBackendIdWithClusterId(clusterId);
                    } catch (ComputeGroupException e) {
                        LOG.debug("ROW binlog TTL discovery skipped, compute group={}, tablet={}",
                                clusterId, tablet.getId(), e);
                        continue;
                    }
                    Backend backend = info.getBackend(backendId);
                    if (backend != null && backend.isAlive()) {
                        List<Long> backendTablets = targets.computeIfAbsent(backendId, ignored -> new ArrayList<>());
                        backendTablets.add(tablet.getId());
                        batchFull |= backendTablets.size() >= MAX_TABLETS_PER_BACKEND;
                    }
                }
                // Bound each BE's burst to its prepare queue size, but finish all compute groups
                // for this tablet before yielding so advancing the iterator cannot skip a target.
                if (batchFull) {
                    break;
                }
            }
            for (Map.Entry<Long, List<Long>> target : targets.entrySet()) {
                Backend backend = info.getBackend(target.getKey());
                if (backend == null || !backend.isAlive()) {
                    continue;
                }
                InternalService.PSyncTabletMetaRequest request = InternalService.PSyncTabletMetaRequest.newBuilder()
                        .addAllTabletIds(target.getValue()).setDiscoverRowBinlogTtl(true).build();
                try {
                    Futures.addCallback(BackendServiceProxy.getInstance()
                                    .syncTabletMeta(backend.getBrpcAddress(), request),
                            new FutureCallback<InternalService.PSyncTabletMetaResponse>() {
                                @Override
                                public void onSuccess(InternalService.PSyncTabletMetaResponse response) {
                                    if (!response.hasStatus() || response.getStatus().getStatusCode()
                                            != TStatusCode.OK.getValue() || response.getFailedTablets() > 0) {
                                        LOG.warn("ROW binlog TTL discovery deferred, backend={}, response={}",
                                                target.getKey(), response);
                                    }
                                }

                                @Override
                                public void onFailure(Throwable t) {
                                    LOG.warn("ROW binlog TTL discovery failed, backend={}", target.getKey(), t);
                                }
                            }, MoreExecutors.directExecutor());
                } catch (RpcException e) {
                    LOG.warn("ROW binlog TTL discovery dispatch failed, backend={}", target.getKey(), e);
                }
            }
            // Continue a large sweep promptly, but avoid repeatedly scanning a small catalog.
            setInterval(databases.hasNext() || tables.hasNext() || partitions.hasNext()
                    || indexes.hasNext() || tablets.hasNext() ? 100 : 15000);
        } catch (RuntimeException e) {
            LOG.warn("ROW binlog TTL discovery deferred until next catalog sweep", e);
        }
    }

    private Tablet nextTablet(long deadline) {
        while (System.nanoTime() < deadline) {
            if (tablets.hasNext()) {
                return tablets.next();
            }
            if (indexes.hasNext()) {
                MaterializedIndex index = indexes.next();
                if (index.isRowBinlog()) {
                    // getTablets() is an immutable snapshot; no per-tablet copy or long table lock.
                    tablets = index.getTablets().iterator();
                }
            } else if (partitions.hasNext()) {
                Partition partition = partitions.next();
                currentTable.readLock();
                try {
                    indexes = partition.getMaterializedIndices(
                            MaterializedIndex.IndexExtState.VISIBLE, true).iterator();
                } finally {
                    currentTable.readUnlock();
                }
            } else if (tables.hasNext()) {
                Table table = tables.next();
                if (table instanceof OlapTable && ((OlapTable) table).hasRowBinlogTtl()) {
                    currentTable = (OlapTable) table;
                    currentTable.readLock();
                    try {
                        partitions = new ArrayList<>(currentTable.getPartitions()).iterator();
                    } finally {
                        currentTable.readUnlock();
                    }
                }
            } else if (databases.hasNext()) {
                tables = databases.next().getTables().iterator();
            } else {
                currentTable = null;
                return null;
            }
        }
        return null;
    }
}