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}