001/*
002 * Licensed to the Apache Software Foundation (ASF) under one
003 * or more contributor license agreements.  See the NOTICE file
004 * distributed with this work for additional information
005 * regarding copyright ownership.  The ASF licenses this file
006 * to you under the Apache License, Version 2.0 (the
007 * "License"); you may not use this file except in compliance
008 * with the License.  You may obtain a copy of the License at
009 *
010 *     http://www.apache.org/licenses/LICENSE-2.0
011 *
012 * Unless required by applicable law or agreed to in writing, software
013 * distributed under the License is distributed on an "AS IS" BASIS,
014 * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
015 * See the License for the specific language governing permissions and
016 * limitations under the License.
017 */
018package org.apache.hadoop.hbase.master.assignment;
019
020import static org.apache.hadoop.hbase.HConstants.DEFAULT_HBASE_ENABLE_SEPARATE_CHILD_REGIONS;
021
022import java.io.IOException;
023import java.util.ArrayList;
024import java.util.Collections;
025import java.util.List;
026import java.util.ListIterator;
027import java.util.NavigableSet;
028import java.util.stream.Collectors;
029import java.util.stream.IntStream;
030import java.util.stream.Stream;
031import org.apache.commons.lang3.ArrayUtils;
032import org.apache.hadoop.conf.Configuration;
033import org.apache.hadoop.fs.FileSystem;
034import org.apache.hadoop.fs.Path;
035import org.apache.hadoop.hbase.HBaseIOException;
036import org.apache.hadoop.hbase.HConstants;
037import org.apache.hadoop.hbase.ServerName;
038import org.apache.hadoop.hbase.client.AsyncRegionServerAdmin;
039import org.apache.hadoop.hbase.client.RegionInfo;
040import org.apache.hadoop.hbase.client.RegionReplicaUtil;
041import org.apache.hadoop.hbase.favored.FavoredNodesManager;
042import org.apache.hadoop.hbase.master.RegionState;
043import org.apache.hadoop.hbase.master.procedure.MasterProcedureEnv;
044import org.apache.hadoop.hbase.util.CommonFSUtils;
045import org.apache.hadoop.hbase.util.FSUtils;
046import org.apache.hadoop.hbase.util.FutureUtils;
047import org.apache.hadoop.hbase.wal.WALSplitUtil;
048import org.apache.yetus.audience.InterfaceAudience;
049import org.slf4j.Logger;
050import org.slf4j.LoggerFactory;
051
052import org.apache.hadoop.hbase.shaded.protobuf.RequestConverter;
053import org.apache.hadoop.hbase.shaded.protobuf.generated.AdminProtos.GetRegionInfoRequest;
054import org.apache.hadoop.hbase.shaded.protobuf.generated.AdminProtos.GetRegionInfoResponse;
055
056/**
057 * Utility for this assignment package only.
058 */
059@InterfaceAudience.Private
060final class AssignmentManagerUtil {
061  private static final Logger LOG = LoggerFactory.getLogger(AssignmentManagerUtil.class);
062  private static final int DEFAULT_REGION_REPLICA = 1;
063
064  private AssignmentManagerUtil() {
065  }
066
067  /**
068   * Raw call to remote regionserver to get info on a particular region.
069   * @throws IOException Let it out so can report this IOE as reason for failure
070   */
071  static GetRegionInfoResponse getRegionInfoResponse(final MasterProcedureEnv env,
072    final ServerName regionLocation, final RegionInfo hri) throws IOException {
073    return getRegionInfoResponse(env, regionLocation, hri, false);
074  }
075
076  static GetRegionInfoResponse getRegionInfoResponse(final MasterProcedureEnv env,
077    final ServerName regionLocation, final RegionInfo hri, boolean includeBestSplitRow)
078    throws IOException {
079    AsyncRegionServerAdmin admin =
080      env.getMasterServices().getAsyncClusterConnection().getRegionServerAdmin(regionLocation);
081    GetRegionInfoRequest request = null;
082    if (includeBestSplitRow) {
083      request = RequestConverter.buildGetRegionInfoRequest(hri.getRegionName(), false, true);
084    } else {
085      request = RequestConverter.buildGetRegionInfoRequest(hri.getRegionName());
086    }
087    return FutureUtils.get(admin.getRegionInfo(request));
088  }
089
090  private static void lock(List<RegionStateNode> regionNodes) {
091    regionNodes.iterator().forEachRemaining(RegionStateNode::lock);
092  }
093
094  private static void unlock(List<RegionStateNode> regionNodes) {
095    for (ListIterator<RegionStateNode> iter = regionNodes.listIterator(regionNodes.size()); iter
096      .hasPrevious();) {
097      iter.previous().unlock();
098    }
099  }
100
101  static TransitRegionStateProcedure[] createUnassignProceduresForSplitOrMerge(
102    MasterProcedureEnv env, Stream<RegionInfo> regions, int regionReplication) throws IOException {
103    List<RegionStateNode> regionNodes = regions
104      .flatMap(hri -> IntStream.range(0, regionReplication)
105        .mapToObj(i -> RegionReplicaUtil.getRegionInfoForReplica(hri, i)))
106      .map(env.getAssignmentManager().getRegionStates()::getOrCreateRegionStateNode)
107      .collect(Collectors.toList());
108    TransitRegionStateProcedure[] procs = new TransitRegionStateProcedure[regionNodes.size()];
109    boolean rollback = true;
110    int i = 0;
111    // hold the lock at once, and then release it in finally. This is important as SCP may jump in
112    // if we release the lock in the middle when we want to do rollback, and cause problems.
113    lock(regionNodes);
114    try {
115      for (; i < procs.length; i++) {
116        RegionStateNode regionNode = regionNodes.get(i);
117        TransitRegionStateProcedure proc =
118          TransitRegionStateProcedure.unassignSplitMerge(env, regionNode.getRegionInfo());
119        if (regionNode.getProcedure() != null) {
120          throw new HBaseIOException(
121            "The parent region " + regionNode + " is currently in transition, give up");
122        }
123        regionNode.setProcedure(proc);
124        procs[i] = proc;
125      }
126      // all succeeded, set rollback to false
127      rollback = false;
128    } finally {
129      if (rollback) {
130        for (;;) {
131          i--;
132          if (i < 0) {
133            break;
134          }
135          RegionStateNode regionNode = regionNodes.get(i);
136          regionNode.unsetProcedure(procs[i]);
137        }
138      }
139      unlock(regionNodes);
140    }
141    return procs;
142  }
143
144  /**
145   * Create assign procedures for the give regions, according to the {@code regionReplication}.
146   * <p/>
147   * For rolling back, we will submit procedures directly to the {@code ProcedureExecutor}, so it is
148   * possible that we persist the newly scheduled procedures, and then crash before persisting the
149   * rollback state, so when we arrive here the second time, it is possible that some regions have
150   * already been associated with a TRSP.
151   * @param ignoreIfInTransition if true, will skip creating TRSP for the given region if it is
152   *                             already in transition, otherwise we will add an assert that it
153   *                             should not in transition.
154   */
155  private static TransitRegionStateProcedure[] createAssignProcedures(MasterProcedureEnv env,
156    List<RegionInfo> regions, int regionReplication, ServerName targetServer,
157    boolean ignoreIfInTransition) {
158    // create the assign procs only for the primary region using the targetServer
159    TransitRegionStateProcedure[] primaryRegionProcs =
160      regions.stream().map(env.getAssignmentManager().getRegionStates()::getOrCreateRegionStateNode)
161        .map(regionNode -> {
162          TransitRegionStateProcedure proc =
163            TransitRegionStateProcedure.assign(env, regionNode.getRegionInfo(), targetServer);
164          regionNode.lock();
165          try {
166            if (ignoreIfInTransition) {
167              if (regionNode.isTransitionScheduled()) {
168                return null;
169              }
170            } else {
171              // should never fail, as we have the exclusive region lock, and the region is newly
172              // created, or has been successfully closed so should not be on any servers, so SCP
173              // will
174              // not process it either.
175              assert !regionNode.isTransitionScheduled();
176            }
177            regionNode.setProcedure(proc);
178          } finally {
179            regionNode.unlock();
180          }
181          return proc;
182        }).filter(p -> p != null).toArray(TransitRegionStateProcedure[]::new);
183    if (regionReplication == DEFAULT_REGION_REPLICA) {
184      // this is the default case
185      return primaryRegionProcs;
186    }
187    // collect the replica region infos
188    List<RegionInfo> replicaRegionInfos =
189      new ArrayList<RegionInfo>(regions.size() * (regionReplication - 1));
190    for (RegionInfo hri : regions) {
191      // start the index from 1
192      for (int i = 1; i < regionReplication; i++) {
193        RegionInfo ri = RegionReplicaUtil.getRegionInfoForReplica(hri, i);
194        // apply ignoreRITs to replica regions as well.
195        if (
196          !ignoreIfInTransition || !env.getAssignmentManager().getRegionStates()
197            .getOrCreateRegionStateNode(ri).isTransitionScheduled()
198        ) {
199          replicaRegionInfos.add(ri);
200        }
201      }
202    }
203
204    // create round robin procs. Note that we exclude the primary region's target server
205    TransitRegionStateProcedure[] replicaRegionAssignProcs =
206      env.getAssignmentManager().createRoundRobinAssignProcedures(replicaRegionInfos,
207        Collections.singletonList(targetServer));
208    // combine both the procs and return the result
209    return ArrayUtils.addAll(primaryRegionProcs, replicaRegionAssignProcs);
210  }
211
212  /**
213   * Create round robin assign procedures for the given regions, according to the
214   * {@code regionReplication}.
215   * <p/>
216   * For rolling back, we will submit procedures directly to the {@code ProcedureExecutor}, so it is
217   * possible that we persist the newly scheduled procedures, and then crash before persisting the
218   * rollback state, so when we arrive here the second time, it is possible that some regions have
219   * already been associated with a TRSP.
220   * @param ignoreIfInTransition if true, will skip creating TRSP for the given region if it is
221   *                             already in transition, otherwise we will add an assert that it
222   *                             should not in transition.
223   */
224  private static TransitRegionStateProcedure[] createRoundRobinAssignProcedures(
225    MasterProcedureEnv env, List<RegionInfo> regions, int regionReplication,
226    List<ServerName> serversToExclude, boolean ignoreIfInTransition) {
227    List<RegionInfo> regionsAndReplicas = new ArrayList<>(regions);
228    if (regionReplication != DEFAULT_REGION_REPLICA) {
229
230      // collect the replica region infos
231      List<RegionInfo> replicaRegionInfos =
232        new ArrayList<RegionInfo>(regions.size() * (regionReplication - 1));
233      for (RegionInfo hri : regions) {
234        // start the index from 1
235        for (int i = 1; i < regionReplication; i++) {
236          replicaRegionInfos.add(RegionReplicaUtil.getRegionInfoForReplica(hri, i));
237        }
238      }
239      regionsAndReplicas.addAll(replicaRegionInfos);
240    }
241    if (ignoreIfInTransition) {
242      for (RegionInfo region : regionsAndReplicas) {
243        if (
244          env.getAssignmentManager().getRegionStates().getOrCreateRegionStateNode(region)
245            .isTransitionScheduled()
246        ) {
247          return null;
248        }
249      }
250    }
251    // create round robin procs. Note that we exclude the primary region's target server
252    return env.getAssignmentManager().createRoundRobinAssignProcedures(regionsAndReplicas,
253      serversToExclude);
254  }
255
256  static TransitRegionStateProcedure[] createAssignProceduresForSplitDaughters(
257    MasterProcedureEnv env, List<RegionInfo> daughters, int regionReplication,
258    ServerName parentServer) {
259    if (
260      env.getMasterConfiguration().getBoolean(HConstants.HBASE_ENABLE_SEPARATE_CHILD_REGIONS,
261        DEFAULT_HBASE_ENABLE_SEPARATE_CHILD_REGIONS)
262    ) {
263      // keep one daughter on the parent region server
264      TransitRegionStateProcedure[] daughterOne = createAssignProcedures(env,
265        Collections.singletonList(daughters.get(0)), regionReplication, parentServer, false);
266      // round robin assign the other daughter
267      TransitRegionStateProcedure[] daughterTwo =
268        createRoundRobinAssignProcedures(env, Collections.singletonList(daughters.get(1)),
269          regionReplication, Collections.singletonList(parentServer), false);
270      return ArrayUtils.addAll(daughterOne, daughterTwo);
271    }
272    return createAssignProceduresForOpeningNewRegions(env, daughters, regionReplication,
273      parentServer);
274  }
275
276  static TransitRegionStateProcedure[] createAssignProceduresForOpeningNewRegions(
277    MasterProcedureEnv env, List<RegionInfo> regions, int regionReplication,
278    ServerName targetServer) {
279    return createAssignProcedures(env, regions, regionReplication, targetServer, false);
280  }
281
282  static void reopenRegionsForRollback(MasterProcedureEnv env, List<RegionInfo> regions,
283    int regionReplication, ServerName targetServer) {
284    TransitRegionStateProcedure[] procs =
285      createAssignProcedures(env, regions, regionReplication, targetServer, true);
286    if (procs.length > 0) {
287      env.getMasterServices().getMasterProcedureExecutor().submitProcedures(procs);
288    }
289  }
290
291  static void removeNonDefaultReplicas(MasterProcedureEnv env, Stream<RegionInfo> regions,
292    int regionReplication) {
293    // Remove from in-memory states
294    regions.flatMap(hri -> IntStream.range(1, regionReplication)
295      .mapToObj(i -> RegionReplicaUtil.getRegionInfoForReplica(hri, i))).forEach(hri -> {
296        env.getAssignmentManager().getRegionStates().deleteRegion(hri);
297        env.getMasterServices().getServerManager().removeRegion(hri);
298        FavoredNodesManager fnm = env.getMasterServices().getFavoredNodesManager();
299        if (fnm != null) {
300          fnm.deleteFavoredNodesForRegions(Collections.singletonList(hri));
301        }
302      });
303  }
304
305  static void checkClosedRegion(MasterProcedureEnv env, RegionInfo regionInfo) throws IOException {
306    if (!WALSplitUtil.hasRecoveredEdits(env.getMasterConfiguration(), regionInfo)) {
307      return;
308    }
309    // Robustness: corner cases can leave behind recovered.edits whose max seqid is already
310    // covered by the region's durable seqid. Drop those and proceed instead of aborting.
311    if (tryDropStaleRecoveredEdits(env, regionInfo)) {
312      return;
313    }
314    throw new IOException("Recovered.edits are found in Region: " + regionInfo
315      + ", abort split/merge to prevent data loss");
316  }
317
318  /**
319   * Try to remove recovered.edits files that are provably below the region's last flushed seqid.
320   * @return true if, after cleanup, no recovered.edits remain for the region
321   */
322  private static boolean tryDropStaleRecoveredEdits(MasterProcedureEnv env, RegionInfo regionInfo) {
323    long durableSeqId = env.getMasterServices().getServerManager()
324      .getLastFlushedSequenceId(regionInfo.getEncodedNameAsBytes()).getLastFlushedSequenceId();
325    if (durableSeqId <= 0L) {
326      // No authoritative durability info at the master; play safe and let the caller abort.
327      return false;
328    }
329    try {
330      Configuration conf = env.getMasterConfiguration();
331      Path regionWALDir =
332        CommonFSUtils.getWALRegionDir(conf, regionInfo.getTable(), regionInfo.getEncodedName());
333      Path regionDir = FSUtils.getRegionDirFromRootDir(CommonFSUtils.getRootDir(conf), regionInfo);
334      Path wrongRegionWALDir = CommonFSUtils.getWrongWALRegionDir(conf, regionInfo.getTable(),
335        regionInfo.getEncodedName());
336      FileSystem walFs = CommonFSUtils.getWALFileSystem(conf);
337      FileSystem rootFs = CommonFSUtils.getRootDirFileSystem(conf);
338      return dropStaleEditsUnder(walFs, regionWALDir, durableSeqId, regionInfo)
339        && dropStaleEditsUnder(rootFs, regionDir, durableSeqId, regionInfo)
340        && dropStaleEditsUnder(walFs, wrongRegionWALDir, durableSeqId, regionInfo);
341    } catch (IOException e) {
342      LOG.warn("Failed to inspect recovered.edits for {}; falling back to abort", regionInfo, e);
343      return false;
344    }
345  }
346
347  private static boolean dropStaleEditsUnder(FileSystem fs, Path regionDir, long durableSeqId,
348    RegionInfo regionInfo) throws IOException {
349    NavigableSet<Path> files = WALSplitUtil.getSplitEditFilesSorted(fs, regionDir);
350    if (files.isEmpty()) {
351      return true;
352    }
353    for (Path p : files) {
354      // getSplitEditFilesSorted restricts filenames to WALSplitUtil.EDITFILES_NAME_PATTERN
355      // (`-?[0-9]+`), so parseLong cannot throw here.
356      long fileMaxSeqId;
357      try {
358        fileMaxSeqId = Long.parseLong(p.getName());
359      } catch (NumberFormatException e) {
360        LOG.warn("Unable to parse recovered.edits sequence id from {}; falling back to abort", p,
361          e);
362        return false;
363      }
364      if (fileMaxSeqId > durableSeqId) {
365        LOG.info("Recovered.edits {} for {} has maxSeqId={} > durableSeqId={}; needs replay", p,
366          regionInfo, fileMaxSeqId, durableSeqId);
367        return false;
368      }
369    }
370    for (Path p : files) {
371      LOG.info("Removing stale recovered.edits {} for {} (durableSeqId={})", p, regionInfo,
372        durableSeqId);
373      if (!fs.delete(p, false)) {
374        LOG.warn("Failed to delete stale recovered.edits {} for {}", p, regionInfo);
375        return false;
376      }
377    }
378    return true;
379  }
380
381  /**
382   * For splitting, need to test both region info and state, and will return true if either of the
383   * test returns true. Please see the comments in
384   * {@link AssignmentManager#markRegionAsSplit(RegionInfo, ServerName, RegionInfo, RegionInfo)} for
385   * more details on why we need to test two conditions.
386   */
387  static boolean isSplitOrMerged(RegionStateNode regionStateNode) {
388    return regionStateNode.getState() == RegionState.State.SPLIT
389      || regionStateNode.getRegionInfo().isSplit()
390      || regionStateNode.getState() == RegionState.State.MERGED;
391  }
392}