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.mapreduce;
019
020import static org.junit.jupiter.api.Assertions.assertEquals;
021import static org.junit.jupiter.api.Assertions.assertFalse;
022import static org.junit.jupiter.api.Assertions.assertNotEquals;
023import static org.junit.jupiter.api.Assertions.assertTrue;
024import static org.junit.jupiter.api.Assertions.fail;
025
026import java.io.IOException;
027import java.util.List;
028import java.util.NavigableMap;
029import java.util.TreeMap;
030import org.apache.hadoop.conf.Configuration;
031import org.apache.hadoop.fs.FileSystem;
032import org.apache.hadoop.fs.Path;
033import org.apache.hadoop.hbase.Cell;
034import org.apache.hadoop.hbase.HBaseTestingUtil;
035import org.apache.hadoop.hbase.HConstants;
036import org.apache.hadoop.hbase.KeyValue;
037import org.apache.hadoop.hbase.ServerName;
038import org.apache.hadoop.hbase.TableName;
039import org.apache.hadoop.hbase.client.RegionInfo;
040import org.apache.hadoop.hbase.client.RegionInfoBuilder;
041import org.apache.hadoop.hbase.mapreduce.WALInputFormat.WALKeyRecordReader;
042import org.apache.hadoop.hbase.mapreduce.WALInputFormat.WALRecordReader;
043import org.apache.hadoop.hbase.regionserver.MultiVersionConcurrencyControl;
044import org.apache.hadoop.hbase.testclassification.MapReduceTests;
045import org.apache.hadoop.hbase.testclassification.MediumTests;
046import org.apache.hadoop.hbase.util.Bytes;
047import org.apache.hadoop.hbase.util.CommonFSUtils;
048import org.apache.hadoop.hbase.util.EnvironmentEdgeManager;
049import org.apache.hadoop.hbase.util.Threads;
050import org.apache.hadoop.hbase.wal.AbstractFSWALProvider;
051import org.apache.hadoop.hbase.wal.WAL;
052import org.apache.hadoop.hbase.wal.WALEdit;
053import org.apache.hadoop.hbase.wal.WALEditInternalHelper;
054import org.apache.hadoop.hbase.wal.WALFactory;
055import org.apache.hadoop.hbase.wal.WALKey;
056import org.apache.hadoop.hbase.wal.WALKeyImpl;
057import org.apache.hadoop.mapreduce.InputSplit;
058import org.apache.hadoop.mapreduce.MapReduceTestUtil;
059import org.junit.jupiter.api.AfterAll;
060import org.junit.jupiter.api.BeforeAll;
061import org.junit.jupiter.api.BeforeEach;
062import org.junit.jupiter.api.Tag;
063import org.junit.jupiter.api.Test;
064import org.slf4j.Logger;
065import org.slf4j.LoggerFactory;
066
067/**
068 * JUnit tests for the WALRecordReader
069 */
070@Tag(MapReduceTests.TAG)
071@Tag(MediumTests.TAG)
072public class TestWALRecordReader {
073
074  private static final Logger LOG = LoggerFactory.getLogger(TestWALRecordReader.class);
075  private final static HBaseTestingUtil TEST_UTIL = new HBaseTestingUtil();
076  private static Configuration conf;
077  private static FileSystem fs;
078  private static Path hbaseDir;
079  private static FileSystem walFs;
080  private static Path walRootDir;
081  // visible for TestHLogRecordReader
082  static final TableName tableName = TableName.valueOf(getName());
083  private static final byte[] rowName = tableName.getName();
084  // visible for TestHLogRecordReader
085  static final RegionInfo info = RegionInfoBuilder.newBuilder(tableName).build();
086  private static final byte[] family = Bytes.toBytes("column");
087  private static final byte[] value = Bytes.toBytes("value");
088  private static Path logDir;
089  protected MultiVersionConcurrencyControl mvcc;
090  protected static NavigableMap<byte[], Integer> scopes = new TreeMap<>(Bytes.BYTES_COMPARATOR);
091
092  private static String getName() {
093    return "TestWALRecordReader";
094  }
095
096  private static String getServerName() {
097    ServerName serverName = ServerName.valueOf("TestWALRecordReader", 1, 1);
098    return serverName.toString();
099  }
100
101  @BeforeEach
102  public void setUp() throws Exception {
103    fs.delete(hbaseDir, true);
104    walFs.delete(walRootDir, true);
105    mvcc = new MultiVersionConcurrencyControl();
106  }
107
108  @BeforeAll
109  public static void setUpBeforeClass() throws Exception {
110    // Make block sizes small.
111    conf = TEST_UTIL.getConfiguration();
112    conf.setInt("dfs.blocksize", 1024 * 1024);
113    conf.setInt("dfs.replication", 1);
114    TEST_UTIL.startMiniDFSCluster(1);
115
116    conf = TEST_UTIL.getConfiguration();
117    fs = TEST_UTIL.getDFSCluster().getFileSystem();
118
119    hbaseDir = TEST_UTIL.createRootDir();
120    walRootDir = TEST_UTIL.createWALRootDir();
121    walFs = CommonFSUtils.getWALFileSystem(conf);
122    logDir = new Path(walRootDir, HConstants.HREGION_LOGDIR_NAME);
123  }
124
125  @AfterAll
126  public static void tearDownAfterClass() throws Exception {
127    fs.delete(hbaseDir, true);
128    walFs.delete(walRootDir, true);
129    TEST_UTIL.shutdownMiniCluster();
130  }
131
132  /**
133   * Test partial reads from the WALs based on passed time range.
134   */
135  @Test
136  public void testPartialRead() throws Exception {
137    final WALFactory walfactory = new WALFactory(conf, getName());
138    WAL log = walfactory.getWAL(info);
139    // This test depends on timestamp being millisecond based and the filename of the WAL also
140    // being millisecond based.
141    long ts = EnvironmentEdgeManager.currentTime();
142    WALEdit edit = new WALEdit();
143    WALEditInternalHelper.addExtendedCell(edit,
144      new KeyValue(rowName, family, Bytes.toBytes("1"), ts, value));
145    log.appendData(info, getWalKeyImpl(ts, scopes), edit);
146    edit = new WALEdit();
147    WALEditInternalHelper.addExtendedCell(edit,
148      new KeyValue(rowName, family, Bytes.toBytes("2"), ts + 1, value));
149    log.appendData(info, getWalKeyImpl(ts + 1, scopes), edit);
150    log.sync();
151    Threads.sleep(10);
152    LOG.info("Before 1st WAL roll " + log.toString());
153    log.rollWriter();
154    LOG.info("Past 1st WAL roll " + log.toString());
155
156    Thread.sleep(1);
157    long ts1 = EnvironmentEdgeManager.currentTime();
158
159    edit = new WALEdit();
160    WALEditInternalHelper.addExtendedCell(edit,
161      new KeyValue(rowName, family, Bytes.toBytes("3"), ts1 + 1, value));
162    log.appendData(info, getWalKeyImpl(ts1 + 1, scopes), edit);
163    edit = new WALEdit();
164    WALEditInternalHelper.addExtendedCell(edit,
165      new KeyValue(rowName, family, Bytes.toBytes("4"), ts1 + 2, value));
166    log.appendData(info, getWalKeyImpl(ts1 + 2, scopes), edit);
167    log.sync();
168    log.shutdown();
169    walfactory.shutdown();
170    LOG.info("Closed WAL " + log.toString());
171
172    WALInputFormat input = new WALInputFormat();
173    Configuration jobConf = new Configuration(conf);
174    jobConf.set("mapreduce.input.fileinputformat.inputdir", logDir.toString());
175    jobConf.setLong(WALInputFormat.END_TIME_KEY, ts);
176
177    // Only 1st file is considered, and only its 1st entry is in-range.
178    List<InputSplit> splits = input.getSplits(MapreduceTestingShim.createJobContext(jobConf));
179    assertEquals(1, splits.size());
180    testSplit(splits.get(0), Bytes.toBytes("1"));
181
182    jobConf.setLong(WALInputFormat.END_TIME_KEY, ts1 + 1);
183    splits = input.getSplits(MapreduceTestingShim.createJobContext(jobConf));
184    assertEquals(2, splits.size());
185    // Both entries from first file are in-range.
186    testSplit(splits.get(0), Bytes.toBytes("1"), Bytes.toBytes("2"));
187    // Only the 1st entry from the 2nd file is in-range.
188    testSplit(splits.get(1), Bytes.toBytes("3"));
189
190    jobConf.setLong(WALInputFormat.START_TIME_KEY, ts + 1);
191    jobConf.setLong(WALInputFormat.END_TIME_KEY, ts1 + 1);
192    splits = input.getSplits(MapreduceTestingShim.createJobContext(jobConf));
193    assertEquals(2, splits.size());
194    // The 1st file was created before startTime but stayed open until it rolled, so its 2nd
195    // entry, written at exactly startTime, is in-range.
196    testSplit(splits.get(0), Bytes.toBytes("2"));
197    // Only the 1st entry from the 2nd file is in-range.
198    testSplit(splits.get(1), Bytes.toBytes("3"));
199  }
200
201  /**
202   * Test basic functionality
203   */
204  @Test
205  public void testWALRecordReader() throws Exception {
206    final WALFactory walfactory = new WALFactory(conf, getName());
207    WAL log = walfactory.getWAL(info);
208    byte[] value = Bytes.toBytes("value");
209    WALEdit edit = new WALEdit();
210    WALEditInternalHelper.addExtendedCell(edit, new KeyValue(rowName, family, Bytes.toBytes("1"),
211      EnvironmentEdgeManager.currentTime(), value));
212    long txid =
213      log.appendData(info, getWalKeyImpl(EnvironmentEdgeManager.currentTime(), scopes), edit);
214    log.sync(txid);
215
216    Thread.sleep(1); // make sure 2nd log gets a later timestamp
217    long secondTs = EnvironmentEdgeManager.currentTime();
218    log.rollWriter();
219
220    edit = new WALEdit();
221    WALEditInternalHelper.addExtendedCell(edit, new KeyValue(rowName, family, Bytes.toBytes("2"),
222      EnvironmentEdgeManager.currentTime(), value));
223    txid = log.appendData(info, getWalKeyImpl(EnvironmentEdgeManager.currentTime(), scopes), edit);
224    log.sync(txid);
225    log.shutdown();
226    walfactory.shutdown();
227    long thirdTs = EnvironmentEdgeManager.currentTime();
228
229    // should have 2 log files now
230    WALInputFormat input = new WALInputFormat();
231    Configuration jobConf = new Configuration(conf);
232    jobConf.set("mapreduce.input.fileinputformat.inputdir", logDir.toString());
233
234    // make sure both logs are found
235    List<InputSplit> splits = input.getSplits(MapreduceTestingShim.createJobContext(jobConf));
236    assertEquals(2, splits.size());
237
238    // should return exactly one KV
239    testSplit(splits.get(0), Bytes.toBytes("1"));
240    // same for the 2nd split
241    testSplit(splits.get(1), Bytes.toBytes("2"));
242
243    // now test basic time ranges:
244
245    // set an endtime, the 2nd log file can be ignored completely.
246    jobConf.setLong(WALInputFormat.END_TIME_KEY, secondTs - 1);
247    splits = input.getSplits(MapreduceTestingShim.createJobContext(jobConf));
248    assertEquals(1, splits.size());
249    testSplit(splits.get(0), Bytes.toBytes("1"));
250
251    // now set a start time strictly after the last WAL's modification time
252    jobConf.setLong(WALInputFormat.END_TIME_KEY, Long.MAX_VALUE);
253    jobConf.setLong(WALInputFormat.START_TIME_KEY, thirdTs + 1);
254    splits = input.getSplits(MapreduceTestingShim.createJobContext(jobConf));
255    assertTrue(splits.isEmpty());
256  }
257
258  /**
259   * Test WALRecordReader tolerance to moving WAL from active to archive directory
260   * @throws Exception exception
261   */
262  @Test
263  public void testWALRecordReaderActiveArchiveTolerance() throws Exception {
264    final WALFactory walfactory = new WALFactory(conf, getName());
265    WAL log = walfactory.getWAL(info);
266    byte[] value = Bytes.toBytes("value");
267    WALEdit edit = new WALEdit();
268    WALEditInternalHelper.addExtendedCell(edit, new KeyValue(rowName, family, Bytes.toBytes("1"),
269      EnvironmentEdgeManager.currentTime(), value));
270    long txid =
271      log.appendData(info, getWalKeyImpl(EnvironmentEdgeManager.currentTime(), scopes), edit);
272    log.sync(txid);
273
274    Thread.sleep(10); // make sure 2nd edit gets a later timestamp
275
276    edit = new WALEdit();
277    WALEditInternalHelper.addExtendedCell(edit, new KeyValue(rowName, family, Bytes.toBytes("2"),
278      EnvironmentEdgeManager.currentTime(), value));
279    txid = log.appendData(info, getWalKeyImpl(EnvironmentEdgeManager.currentTime(), scopes), edit);
280    log.sync(txid);
281    log.shutdown();
282
283    // should have 2 log entries now
284    WALInputFormat input = new WALInputFormat();
285    Configuration jobConf = new Configuration(conf);
286    jobConf.set("mapreduce.input.fileinputformat.inputdir", logDir.toString());
287    // make sure log is found
288    List<InputSplit> splits = input.getSplits(MapreduceTestingShim.createJobContext(jobConf));
289    assertEquals(1, splits.size());
290    WALInputFormat.WALSplit split = (WALInputFormat.WALSplit) splits.get(0);
291    LOG.debug("log=" + logDir + " file=" + split.getLogFileName());
292
293    testSplitWithMovingWAL(splits.get(0), Bytes.toBytes("1"), Bytes.toBytes("2"));
294  }
295
296  protected WALKeyImpl getWalKeyImpl(final long time, NavigableMap<byte[], Integer> scopes) {
297    return new WALKeyImpl(info.getEncodedNameAsBytes(), tableName, time, mvcc, scopes);
298  }
299
300  private WALRecordReader<WALKey> getReader() {
301    return new WALKeyRecordReader();
302  }
303
304  /**
305   * Create a new reader from the split, and match the edits against the passed columns.
306   */
307  private void testSplit(InputSplit split, byte[]... columns) throws Exception {
308    WALRecordReader<WALKey> reader = getReader();
309    reader.initialize(split, MapReduceTestUtil.createDummyMapTaskAttemptContext(conf));
310
311    for (byte[] column : columns) {
312      assertTrue(reader.nextKeyValue());
313      Cell cell = reader.getCurrentValue().getCells().get(0);
314      if (
315        !Bytes.equals(column, 0, column.length, cell.getQualifierArray(), cell.getQualifierOffset(),
316          cell.getQualifierLength())
317      ) {
318        fail("expected [" + Bytes.toString(column) + "], actual [" + Bytes.toString(
319          cell.getQualifierArray(), cell.getQualifierOffset(), cell.getQualifierLength()) + "]");
320      }
321    }
322    assertFalse(reader.nextKeyValue());
323    reader.close();
324  }
325
326  /**
327   * Create a new reader from the split, match the edits against the passed columns, moving WAL to
328   * archive in between readings
329   */
330  private void testSplitWithMovingWAL(InputSplit split, byte[] col1, byte[] col2) throws Exception {
331    WALRecordReader<WALKey> reader = getReader();
332    reader.initialize(split, MapReduceTestUtil.createDummyMapTaskAttemptContext(conf));
333
334    assertTrue(reader.nextKeyValue());
335    Cell cell = reader.getCurrentValue().getCells().get(0);
336    if (
337      !Bytes.equals(col1, 0, col1.length, cell.getQualifierArray(), cell.getQualifierOffset(),
338        cell.getQualifierLength())
339    ) {
340      fail("expected [" + Bytes.toString(col1) + "], actual [" + Bytes.toString(
341        cell.getQualifierArray(), cell.getQualifierOffset(), cell.getQualifierLength()) + "]");
342    }
343    // Move log file to archive directory
344    // While WAL record reader is open
345    WALInputFormat.WALSplit split_ = (WALInputFormat.WALSplit) split;
346    Path logFile = new Path(split_.getLogFileName());
347    Path archivedLogDir = getWALArchiveDir(conf);
348    Path archivedLogLocation = new Path(archivedLogDir, logFile.getName());
349    assertNotEquals(split_.getLogFileName(), archivedLogLocation.toString());
350
351    assertTrue(fs.rename(logFile, archivedLogLocation));
352    assertTrue(fs.exists(archivedLogDir));
353    assertFalse(fs.exists(logFile));
354    // TODO: This is not behaving as expected. WALInputFormat#WALKeyRecordReader doesn't open
355    // TODO: the archivedLogLocation to read next key value.
356    assertTrue(reader.nextKeyValue());
357    cell = reader.getCurrentValue().getCells().get(0);
358    if (
359      !Bytes.equals(col2, 0, col2.length, cell.getQualifierArray(), cell.getQualifierOffset(),
360        cell.getQualifierLength())
361    ) {
362      fail("expected [" + Bytes.toString(col2) + "], actual [" + Bytes.toString(
363        cell.getQualifierArray(), cell.getQualifierOffset(), cell.getQualifierLength()) + "]");
364    }
365    reader.close();
366  }
367
368  private Path getWALArchiveDir(Configuration conf) throws IOException {
369    Path rootDir = CommonFSUtils.getWALRootDir(conf);
370    String archiveDir = AbstractFSWALProvider.getWALArchiveDirectoryName(conf, getServerName());
371    return new Path(rootDir, archiveDir);
372  }
373}