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.regionserver.compactions; 019 020import static org.apache.hadoop.hbase.HConstants.MAJOR_COMPACTION_PERIOD; 021import static org.apache.hadoop.hbase.regionserver.CustomTieringMultiFileWriter.CUSTOM_TIERING_TIME_RANGE; 022import static org.apache.hadoop.hbase.regionserver.compactions.CustomCellTieringValueProvider.TIERING_CELL_QUALIFIER; 023import static org.apache.hadoop.hbase.regionserver.compactions.CustomTieredCompactor.TIERING_VALUE_PROVIDER; 024import static org.apache.hadoop.hbase.regionserver.compactions.RowKeyDateTieringValueProvider.TIERING_KEY_DATE_FORMAT; 025import static org.apache.hadoop.hbase.regionserver.compactions.RowKeyDateTieringValueProvider.TIERING_KEY_DATE_PATTERN; 026import static org.junit.jupiter.api.Assertions.assertEquals; 027import static org.junit.jupiter.api.Assertions.assertNotNull; 028import static org.junit.jupiter.api.Assertions.assertNull; 029import static org.junit.jupiter.api.Assertions.assertTrue; 030import static org.junit.jupiter.api.Assertions.fail; 031 032import java.io.IOException; 033import java.text.SimpleDateFormat; 034import java.util.ArrayList; 035import java.util.Date; 036import java.util.List; 037import org.apache.hadoop.hbase.HBaseTestingUtil; 038import org.apache.hadoop.hbase.TableName; 039import org.apache.hadoop.hbase.Waiter; 040import org.apache.hadoop.hbase.client.Admin; 041import org.apache.hadoop.hbase.client.ColumnFamilyDescriptorBuilder; 042import org.apache.hadoop.hbase.client.Connection; 043import org.apache.hadoop.hbase.client.Put; 044import org.apache.hadoop.hbase.client.Table; 045import org.apache.hadoop.hbase.client.TableDescriptorBuilder; 046import org.apache.hadoop.hbase.regionserver.CustomTieredStoreEngine; 047import org.apache.hadoop.hbase.regionserver.HStore; 048import org.apache.hadoop.hbase.regionserver.HStoreFile; 049import org.apache.hadoop.hbase.regionserver.TimeRangeTracker; 050import org.apache.hadoop.hbase.testclassification.RegionServerTests; 051import org.apache.hadoop.hbase.testclassification.SmallTests; 052import org.apache.hadoop.hbase.util.Bytes; 053import org.junit.jupiter.api.AfterEach; 054import org.junit.jupiter.api.BeforeEach; 055import org.junit.jupiter.api.Tag; 056import org.junit.jupiter.api.Test; 057 058@Tag(RegionServerTests.TAG) 059@Tag(SmallTests.TAG) 060public class TestCustomCellTieredCompactor { 061 062 public static final byte[] FAMILY = Bytes.toBytes("cf"); 063 064 protected HBaseTestingUtil utility; 065 066 protected Admin admin; 067 068 @BeforeEach 069 public void setUp() throws Exception { 070 utility = new HBaseTestingUtil(); 071 utility.getConfiguration().setInt("hbase.hfile.compaction.discharger.interval", 10); 072 utility.getConfiguration().setLong(MAJOR_COMPACTION_PERIOD, 10L); 073 utility.startMiniCluster(); 074 } 075 076 @AfterEach 077 public void tearDown() throws Exception { 078 utility.shutdownMiniCluster(); 079 } 080 081 @Test 082 public void testCustomCellTieredCompactor() throws Exception { 083 ColumnFamilyDescriptorBuilder clmBuilder = ColumnFamilyDescriptorBuilder.newBuilder(FAMILY); 084 clmBuilder.setValue("hbase.hstore.engine.class", CustomTieredStoreEngine.class.getName()); 085 clmBuilder.setValue(TIERING_CELL_QUALIFIER, "date"); 086 TableName tableName = TableName.valueOf("testCustomCellTieredCompactor"); 087 TableDescriptorBuilder tblBuilder = TableDescriptorBuilder.newBuilder(tableName); 088 tblBuilder.setColumnFamily(clmBuilder.build()); 089 utility.getAdmin().createTable(tblBuilder.build()); 090 utility.waitTableAvailable(tableName); 091 Connection connection = utility.getConnection(); 092 Table table = connection.getTable(tableName); 093 long recordTime = System.currentTimeMillis(); 094 // write data and flush multiple store files: 095 for (int i = 0; i < 6; i++) { 096 List<Put> puts = new ArrayList<>(2); 097 Put put = new Put(Bytes.toBytes(i)); 098 put.addColumn(FAMILY, Bytes.toBytes("val"), Bytes.toBytes("v" + i)); 099 put.addColumn(FAMILY, Bytes.toBytes("date"), 100 Bytes.toBytes(recordTime - (11L * 366L * 24L * 60L * 60L * 1000L))); 101 puts.add(put); 102 put = new Put(Bytes.toBytes(i + 1000)); 103 put.addColumn(FAMILY, Bytes.toBytes("val"), Bytes.toBytes("v" + (i + 1000))); 104 put.addColumn(FAMILY, Bytes.toBytes("date"), Bytes.toBytes(recordTime)); 105 puts.add(put); 106 table.put(puts); 107 utility.flush(tableName); 108 } 109 table.close(); 110 long firstCompactionTime = System.currentTimeMillis(); 111 utility.getAdmin().majorCompact(tableName); 112 Waiter.waitFor(utility.getConfiguration(), 5000, 113 () -> utility.getMiniHBaseCluster().getMaster().getLastMajorCompactionTimestamp(tableName) 114 > firstCompactionTime); 115 long numHFiles = utility.getNumHFiles(tableName, FAMILY); 116 // The first major compaction would have no means to detect more than one tier, 117 // because without the min/max values available in the file info portion of the selected files 118 // for compaction, CustomCellDateTieredCompactionPolicy has no means 119 // to calculate the proper boundaries. 120 assertEquals(1, numHFiles); 121 utility.getMiniHBaseCluster().getRegions(tableName).get(0).getStore(FAMILY).getStorefiles() 122 .forEach(file -> { 123 byte[] rangeBytes = file.getMetadataValue(CUSTOM_TIERING_TIME_RANGE); 124 assertNotNull(rangeBytes); 125 try { 126 TimeRangeTracker timeRangeTracker = TimeRangeTracker.parseFrom(rangeBytes); 127 assertEquals((recordTime - (11L * 366L * 24L * 60L * 60L * 1000L)), 128 timeRangeTracker.getMin()); 129 assertEquals(recordTime, timeRangeTracker.getMax()); 130 } catch (IOException e) { 131 fail(e.getMessage()); 132 } 133 }); 134 // now do major compaction again, to make sure we write two separate files 135 long secondCompactionTime = System.currentTimeMillis(); 136 utility.getAdmin().majorCompact(tableName); 137 Waiter.waitFor(utility.getConfiguration(), 5000, 138 () -> utility.getMiniHBaseCluster().getMaster().getLastMajorCompactionTimestamp(tableName) 139 > secondCompactionTime); 140 numHFiles = utility.getNumHFiles(tableName, FAMILY); 141 assertEquals(2, numHFiles); 142 utility.getMiniHBaseCluster().getRegions(tableName).get(0).getStore(FAMILY).getStorefiles() 143 .forEach(file -> { 144 byte[] rangeBytes = file.getMetadataValue(CUSTOM_TIERING_TIME_RANGE); 145 assertNotNull(rangeBytes); 146 try { 147 TimeRangeTracker timeRangeTracker = TimeRangeTracker.parseFrom(rangeBytes); 148 assertEquals(timeRangeTracker.getMin(), timeRangeTracker.getMax()); 149 } catch (IOException e) { 150 fail(e.getMessage()); 151 } 152 }); 153 } 154 155 @Test 156 public void testCustomCellTieredCompactorWithRowKeyDateTieringValue() throws Exception { 157 // Restart mini cluster with RowKeyDateTieringValueProvider 158 utility.shutdownMiniCluster(); 159 utility.getConfiguration().set(TIERING_VALUE_PROVIDER, 160 RowKeyDateTieringValueProvider.class.getName()); 161 utility.startMiniCluster(); 162 163 ColumnFamilyDescriptorBuilder clmBuilder = ColumnFamilyDescriptorBuilder.newBuilder(FAMILY); 164 clmBuilder.setValue("hbase.hstore.engine.class", CustomTieredStoreEngine.class.getName()); 165 166 // Table 1: Date at end with format yyyyMMddHHmmssSSS 167 TableName table1Name = TableName.valueOf("testTable1"); 168 TableDescriptorBuilder tbl1Builder = TableDescriptorBuilder.newBuilder(table1Name); 169 tbl1Builder.setColumnFamily(clmBuilder.build()); 170 tbl1Builder.setValue(TIERING_KEY_DATE_PATTERN, "(\\d{17})$"); 171 tbl1Builder.setValue(TIERING_KEY_DATE_FORMAT, "yyyyMMddHHmmssSSS"); 172 utility.getAdmin().createTable(tbl1Builder.build()); 173 utility.waitTableAvailable(table1Name); 174 175 // Table 2: Date at beginning with format yyyy-MM-dd HH:mm:ss 176 TableName table2Name = TableName.valueOf("testTable2"); 177 TableDescriptorBuilder tbl2Builder = TableDescriptorBuilder.newBuilder(table2Name); 178 tbl2Builder.setColumnFamily(clmBuilder.build()); 179 tbl2Builder.setValue(TIERING_KEY_DATE_PATTERN, "^(\\d{4}-\\d{2}-\\d{2} \\d{2}:\\d{2}:\\d{2})"); 180 tbl2Builder.setValue(TIERING_KEY_DATE_FORMAT, "yyyy-MM-dd HH:mm:ss"); 181 utility.getAdmin().createTable(tbl2Builder.build()); 182 utility.waitTableAvailable(table2Name); 183 184 Connection connection = utility.getConnection(); 185 long recordTime = System.currentTimeMillis(); 186 long oldTime = recordTime - (11L * 366L * 24L * 60L * 60L * 1000L); 187 188 SimpleDateFormat sdf1 = new SimpleDateFormat("yyyyMMddHHmmssSSS"); 189 SimpleDateFormat sdf2 = new SimpleDateFormat("yyyy-MM-dd HH:mm:ss"); 190 191 // Write to Table 1 with date at end 192 Table table1 = connection.getTable(table1Name); 193 for (int i = 0; i < 6; i++) { 194 List<Put> puts = new ArrayList<>(2); 195 196 // Old data 197 String oldDate = sdf1.format(new Date(oldTime)); 198 Put put = new Put(Bytes.toBytes("row_" + i + "_" + oldDate)); 199 put.addColumn(FAMILY, Bytes.toBytes("val"), Bytes.toBytes("v" + i)); 200 puts.add(put); 201 202 // Recent data 203 String recentDate = sdf1.format(new Date(recordTime)); 204 put = new Put(Bytes.toBytes("row_" + (i + 1000) + "_" + recentDate)); 205 put.addColumn(FAMILY, Bytes.toBytes("val"), Bytes.toBytes("v" + (i + 1000))); 206 puts.add(put); 207 208 table1.put(puts); 209 utility.flush(table1Name); 210 } 211 table1.close(); 212 213 // Write to Table 2 with date at beginning 214 Table table2 = connection.getTable(table2Name); 215 for (int i = 0; i < 6; i++) { 216 List<Put> puts = new ArrayList<>(2); 217 218 // Old data 219 String oldDate = sdf2.format(new Date(oldTime)); 220 Put put = new Put(Bytes.toBytes(oldDate + "_row_" + i)); 221 put.addColumn(FAMILY, Bytes.toBytes("val"), Bytes.toBytes("v" + i)); 222 puts.add(put); 223 224 // Recent data 225 String recentDate = sdf2.format(new Date(recordTime)); 226 put = new Put(Bytes.toBytes(recentDate + "_row_" + (i + 1000))); 227 put.addColumn(FAMILY, Bytes.toBytes("val"), Bytes.toBytes("v" + (i + 1000))); 228 puts.add(put); 229 230 table2.put(puts); 231 utility.flush(table2Name); 232 } 233 table2.close(); 234 235 // First compaction for Table 1 236 long compactionTime1 = System.currentTimeMillis(); 237 utility.getAdmin().majorCompact(table1Name); 238 Waiter.waitFor(utility.getConfiguration(), 5000, 239 () -> utility.getMiniHBaseCluster().getMaster().getLastMajorCompactionTimestamp(table1Name) 240 > compactionTime1); 241 242 assertEquals(1, utility.getNumHFiles(table1Name, FAMILY)); 243 244 utility.getMiniHBaseCluster().getRegions(table1Name).get(0).getStore(FAMILY).getStorefiles() 245 .forEach(file -> { 246 byte[] rangeBytes = file.getMetadataValue(CUSTOM_TIERING_TIME_RANGE); 247 assertNotNull(rangeBytes); 248 try { 249 TimeRangeTracker timeRangeTracker = TimeRangeTracker.parseFrom(rangeBytes); 250 assertEquals(oldTime, timeRangeTracker.getMin()); 251 assertEquals(recordTime, timeRangeTracker.getMax()); 252 } catch (IOException e) { 253 fail(e.getMessage()); 254 } 255 }); 256 257 // Second compaction for Table 1 258 long secondCompactionTime1 = System.currentTimeMillis(); 259 utility.getAdmin().majorCompact(table1Name); 260 Waiter.waitFor(utility.getConfiguration(), 5000, 261 () -> utility.getMiniHBaseCluster().getMaster().getLastMajorCompactionTimestamp(table1Name) 262 > secondCompactionTime1); 263 264 assertEquals(2, utility.getNumHFiles(table1Name, FAMILY)); 265 266 utility.getMiniHBaseCluster().getRegions(table1Name).get(0).getStore(FAMILY).getStorefiles() 267 .forEach(file -> { 268 byte[] rangeBytes = file.getMetadataValue(CUSTOM_TIERING_TIME_RANGE); 269 assertNotNull(rangeBytes); 270 try { 271 TimeRangeTracker timeRangeTracker = TimeRangeTracker.parseFrom(rangeBytes); 272 assertEquals(timeRangeTracker.getMin(), timeRangeTracker.getMax()); 273 } catch (IOException e) { 274 fail(e.getMessage()); 275 } 276 }); 277 278 // First compaction for Table 2 279 long compactionTime2 = System.currentTimeMillis(); 280 utility.getAdmin().majorCompact(table2Name); 281 Waiter.waitFor(utility.getConfiguration(), 5000, 282 () -> utility.getMiniHBaseCluster().getMaster().getLastMajorCompactionTimestamp(table2Name) 283 > compactionTime2); 284 285 assertEquals(1, utility.getNumHFiles(table2Name, FAMILY)); 286 287 utility.getMiniHBaseCluster().getRegions(table2Name).get(0).getStore(FAMILY).getStorefiles() 288 .forEach(file -> { 289 byte[] rangeBytes = file.getMetadataValue(CUSTOM_TIERING_TIME_RANGE); 290 assertNotNull(rangeBytes); 291 try { 292 TimeRangeTracker timeRangeTracker = TimeRangeTracker.parseFrom(rangeBytes); 293 // Table 2 uses yyyy-MM-dd HH:mm:ss format, so we need to account for second precision 294 // The parsed time will be truncated to second precision (no milliseconds) 295 long expectedOldTime = (oldTime / 1000) * 1000; 296 long expectedRecentTime = (recordTime / 1000) * 1000; 297 assertEquals(expectedOldTime, timeRangeTracker.getMin()); 298 assertEquals(expectedRecentTime, timeRangeTracker.getMax()); 299 } catch (IOException e) { 300 fail(e.getMessage()); 301 } 302 }); 303 304 // Second compaction for Table 2 305 long secondCompactionTime2 = System.currentTimeMillis(); 306 utility.getAdmin().majorCompact(table2Name); 307 Waiter.waitFor(utility.getConfiguration(), 5000, 308 () -> utility.getMiniHBaseCluster().getMaster().getLastMajorCompactionTimestamp(table2Name) 309 > secondCompactionTime2); 310 311 assertEquals(2, utility.getNumHFiles(table2Name, FAMILY)); 312 313 utility.getMiniHBaseCluster().getRegions(table2Name).get(0).getStore(FAMILY).getStorefiles() 314 .forEach(file -> { 315 byte[] rangeBytes = file.getMetadataValue(CUSTOM_TIERING_TIME_RANGE); 316 assertNotNull(rangeBytes); 317 try { 318 TimeRangeTracker timeRangeTracker = TimeRangeTracker.parseFrom(rangeBytes); 319 assertEquals(timeRangeTracker.getMin(), timeRangeTracker.getMax()); 320 } catch (IOException e) { 321 fail(e.getMessage()); 322 } 323 }); 324 } 325 326 @Test 327 public void testShouldPerformMajorCompactionWhenTimeRangeMetadataIsNull() throws Exception { 328 ColumnFamilyDescriptorBuilder clmBuilder = ColumnFamilyDescriptorBuilder.newBuilder(FAMILY); 329 clmBuilder.setValue("hbase.hstore.engine.class", CustomTieredStoreEngine.class.getName()); 330 clmBuilder.setValue(TIERING_CELL_QUALIFIER, "date"); 331 TableName tableName = TableName.valueOf("testShouldCompactWhenNoTimeRangeMetadata"); 332 TableDescriptorBuilder tblBuilder = TableDescriptorBuilder.newBuilder(tableName); 333 tblBuilder.setColumnFamily(clmBuilder.build()); 334 utility.getAdmin().createTable(tblBuilder.build()); 335 utility.waitTableAvailable(tableName); 336 Connection connection = utility.getConnection(); 337 Table table = connection.getTable(tableName); 338 long recordTime = System.currentTimeMillis(); 339 // Write data and flush to create store files without CUSTOM_TIERING_TIME_RANGE metadata 340 for (int i = 0; i < 2; i++) { 341 Put put = new Put(Bytes.toBytes(i)); 342 put.addColumn(FAMILY, Bytes.toBytes("val"), Bytes.toBytes("v" + i)); 343 put.addColumn(FAMILY, Bytes.toBytes("date"), Bytes.toBytes(recordTime)); 344 table.put(put); 345 utility.flush(tableName); 346 } 347 table.close(); 348 349 HStore store = 350 (HStore) utility.getMiniHBaseCluster().getRegions(tableName).get(0).getStore(FAMILY); 351 // Verify that flushed files do not have CUSTOM_TIERING_TIME_RANGE metadata 352 for (HStoreFile sf : store.getStorefiles()) { 353 assertNull(sf.getMetadataValue(CUSTOM_TIERING_TIME_RANGE)); 354 } 355 // shouldPerformMajorCompaction must return true due to null timeRangeBytes 356 assertTrue(store.shouldPerformMajorCompaction()); 357 } 358}