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 // getCompactBoundariesForMajor always offers the cutOffTimestamp boundary now, so the first 117 // major compaction already splits the old and recent cells into separate tiers without 118 // relying on CUSTOM_TIERING_TIME_RANGE file metadata. 119 assertEquals(2, numHFiles); 120 utility.getMiniHBaseCluster().getRegions(tableName).get(0).getStore(FAMILY).getStorefiles() 121 .forEach(file -> { 122 byte[] rangeBytes = file.getMetadataValue(CUSTOM_TIERING_TIME_RANGE); 123 assertNotNull(rangeBytes); 124 try { 125 TimeRangeTracker timeRangeTracker = TimeRangeTracker.parseFrom(rangeBytes); 126 assertEquals(timeRangeTracker.getMin(), timeRangeTracker.getMax()); 127 } catch (IOException e) { 128 fail(e.getMessage()); 129 } 130 }); 131 // now do major compaction again, to make sure the two tiers stay separate 132 long secondCompactionTime = System.currentTimeMillis(); 133 utility.getAdmin().majorCompact(tableName); 134 Waiter.waitFor(utility.getConfiguration(), 5000, 135 () -> utility.getMiniHBaseCluster().getMaster().getLastMajorCompactionTimestamp(tableName) 136 > secondCompactionTime); 137 numHFiles = utility.getNumHFiles(tableName, FAMILY); 138 assertEquals(2, numHFiles); 139 utility.getMiniHBaseCluster().getRegions(tableName).get(0).getStore(FAMILY).getStorefiles() 140 .forEach(file -> { 141 byte[] rangeBytes = file.getMetadataValue(CUSTOM_TIERING_TIME_RANGE); 142 assertNotNull(rangeBytes); 143 try { 144 TimeRangeTracker timeRangeTracker = TimeRangeTracker.parseFrom(rangeBytes); 145 assertEquals(timeRangeTracker.getMin(), timeRangeTracker.getMax()); 146 } catch (IOException e) { 147 fail(e.getMessage()); 148 } 149 }); 150 } 151 152 @Test 153 public void testCustomCellTieredCompactorWithRowKeyDateTieringValue() throws Exception { 154 // Restart mini cluster with RowKeyDateTieringValueProvider 155 utility.shutdownMiniCluster(); 156 utility.getConfiguration().set(TIERING_VALUE_PROVIDER, 157 RowKeyDateTieringValueProvider.class.getName()); 158 utility.startMiniCluster(); 159 160 ColumnFamilyDescriptorBuilder clmBuilder = ColumnFamilyDescriptorBuilder.newBuilder(FAMILY); 161 clmBuilder.setValue("hbase.hstore.engine.class", CustomTieredStoreEngine.class.getName()); 162 163 // Table 1: Date at end with format yyyyMMddHHmmssSSS 164 TableName table1Name = TableName.valueOf("testTable1"); 165 TableDescriptorBuilder tbl1Builder = TableDescriptorBuilder.newBuilder(table1Name); 166 tbl1Builder.setColumnFamily(clmBuilder.build()); 167 tbl1Builder.setValue(TIERING_KEY_DATE_PATTERN, "(\\d{17})$"); 168 tbl1Builder.setValue(TIERING_KEY_DATE_FORMAT, "yyyyMMddHHmmssSSS"); 169 utility.getAdmin().createTable(tbl1Builder.build()); 170 utility.waitTableAvailable(table1Name); 171 172 // Table 2: Date at beginning with format yyyy-MM-dd HH:mm:ss 173 TableName table2Name = TableName.valueOf("testTable2"); 174 TableDescriptorBuilder tbl2Builder = TableDescriptorBuilder.newBuilder(table2Name); 175 tbl2Builder.setColumnFamily(clmBuilder.build()); 176 tbl2Builder.setValue(TIERING_KEY_DATE_PATTERN, "^(\\d{4}-\\d{2}-\\d{2} \\d{2}:\\d{2}:\\d{2})"); 177 tbl2Builder.setValue(TIERING_KEY_DATE_FORMAT, "yyyy-MM-dd HH:mm:ss"); 178 utility.getAdmin().createTable(tbl2Builder.build()); 179 utility.waitTableAvailable(table2Name); 180 181 Connection connection = utility.getConnection(); 182 long recordTime = System.currentTimeMillis(); 183 long oldTime = recordTime - (11L * 366L * 24L * 60L * 60L * 1000L); 184 185 SimpleDateFormat sdf1 = new SimpleDateFormat("yyyyMMddHHmmssSSS"); 186 SimpleDateFormat sdf2 = new SimpleDateFormat("yyyy-MM-dd HH:mm:ss"); 187 188 // Write to Table 1 with date at end 189 Table table1 = connection.getTable(table1Name); 190 for (int i = 0; i < 6; i++) { 191 List<Put> puts = new ArrayList<>(2); 192 193 // Old data 194 String oldDate = sdf1.format(new Date(oldTime)); 195 Put put = new Put(Bytes.toBytes("row_" + i + "_" + oldDate)); 196 put.addColumn(FAMILY, Bytes.toBytes("val"), Bytes.toBytes("v" + i)); 197 puts.add(put); 198 199 // Recent data 200 String recentDate = sdf1.format(new Date(recordTime)); 201 put = new Put(Bytes.toBytes("row_" + (i + 1000) + "_" + recentDate)); 202 put.addColumn(FAMILY, Bytes.toBytes("val"), Bytes.toBytes("v" + (i + 1000))); 203 puts.add(put); 204 205 table1.put(puts); 206 utility.flush(table1Name); 207 } 208 table1.close(); 209 210 // Write to Table 2 with date at beginning 211 Table table2 = connection.getTable(table2Name); 212 for (int i = 0; i < 6; i++) { 213 List<Put> puts = new ArrayList<>(2); 214 215 // Old data 216 String oldDate = sdf2.format(new Date(oldTime)); 217 Put put = new Put(Bytes.toBytes(oldDate + "_row_" + i)); 218 put.addColumn(FAMILY, Bytes.toBytes("val"), Bytes.toBytes("v" + i)); 219 puts.add(put); 220 221 // Recent data 222 String recentDate = sdf2.format(new Date(recordTime)); 223 put = new Put(Bytes.toBytes(recentDate + "_row_" + (i + 1000))); 224 put.addColumn(FAMILY, Bytes.toBytes("val"), Bytes.toBytes("v" + (i + 1000))); 225 puts.add(put); 226 227 table2.put(puts); 228 utility.flush(table2Name); 229 } 230 table2.close(); 231 232 // First compaction for Table 1 233 long compactionTime1 = System.currentTimeMillis(); 234 utility.getAdmin().majorCompact(table1Name); 235 Waiter.waitFor(utility.getConfiguration(), 5000, 236 () -> utility.getMiniHBaseCluster().getMaster().getLastMajorCompactionTimestamp(table1Name) 237 > compactionTime1); 238 239 // getCompactBoundariesForMajor always offers the cutOffTimestamp boundary now, so the first 240 // major compaction already splits the old and recent cells into separate tiers. 241 assertEquals(2, utility.getNumHFiles(table1Name, FAMILY)); 242 243 utility.getMiniHBaseCluster().getRegions(table1Name).get(0).getStore(FAMILY).getStorefiles() 244 .forEach(file -> { 245 byte[] rangeBytes = file.getMetadataValue(CUSTOM_TIERING_TIME_RANGE); 246 assertNotNull(rangeBytes); 247 try { 248 TimeRangeTracker timeRangeTracker = TimeRangeTracker.parseFrom(rangeBytes); 249 assertEquals(timeRangeTracker.getMin(), timeRangeTracker.getMax()); 250 } catch (IOException e) { 251 fail(e.getMessage()); 252 } 253 }); 254 255 // Second compaction for Table 1 256 long secondCompactionTime1 = System.currentTimeMillis(); 257 utility.getAdmin().majorCompact(table1Name); 258 Waiter.waitFor(utility.getConfiguration(), 5000, 259 () -> utility.getMiniHBaseCluster().getMaster().getLastMajorCompactionTimestamp(table1Name) 260 > secondCompactionTime1); 261 262 assertEquals(2, utility.getNumHFiles(table1Name, FAMILY)); 263 264 utility.getMiniHBaseCluster().getRegions(table1Name).get(0).getStore(FAMILY).getStorefiles() 265 .forEach(file -> { 266 byte[] rangeBytes = file.getMetadataValue(CUSTOM_TIERING_TIME_RANGE); 267 assertNotNull(rangeBytes); 268 try { 269 TimeRangeTracker timeRangeTracker = TimeRangeTracker.parseFrom(rangeBytes); 270 assertEquals(timeRangeTracker.getMin(), timeRangeTracker.getMax()); 271 } catch (IOException e) { 272 fail(e.getMessage()); 273 } 274 }); 275 276 // First compaction for Table 2 277 long compactionTime2 = System.currentTimeMillis(); 278 utility.getAdmin().majorCompact(table2Name); 279 Waiter.waitFor(utility.getConfiguration(), 5000, 280 () -> utility.getMiniHBaseCluster().getMaster().getLastMajorCompactionTimestamp(table2Name) 281 > compactionTime2); 282 283 // getCompactBoundariesForMajor always offers the cutOffTimestamp boundary now, so the first 284 // major compaction already splits the old and recent cells into separate tiers. 285 assertEquals(2, 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 assertEquals(timeRangeTracker.getMin(), timeRangeTracker.getMax()); 294 } catch (IOException e) { 295 fail(e.getMessage()); 296 } 297 }); 298 299 // Second compaction for Table 2 300 long secondCompactionTime2 = System.currentTimeMillis(); 301 utility.getAdmin().majorCompact(table2Name); 302 Waiter.waitFor(utility.getConfiguration(), 5000, 303 () -> utility.getMiniHBaseCluster().getMaster().getLastMajorCompactionTimestamp(table2Name) 304 > secondCompactionTime2); 305 306 assertEquals(2, utility.getNumHFiles(table2Name, FAMILY)); 307 308 utility.getMiniHBaseCluster().getRegions(table2Name).get(0).getStore(FAMILY).getStorefiles() 309 .forEach(file -> { 310 byte[] rangeBytes = file.getMetadataValue(CUSTOM_TIERING_TIME_RANGE); 311 assertNotNull(rangeBytes); 312 try { 313 TimeRangeTracker timeRangeTracker = TimeRangeTracker.parseFrom(rangeBytes); 314 assertEquals(timeRangeTracker.getMin(), timeRangeTracker.getMax()); 315 } catch (IOException e) { 316 fail(e.getMessage()); 317 } 318 }); 319 } 320 321 @Test 322 public void testShouldPerformMajorCompactionWhenTimeRangeMetadataIsNull() throws Exception { 323 ColumnFamilyDescriptorBuilder clmBuilder = ColumnFamilyDescriptorBuilder.newBuilder(FAMILY); 324 clmBuilder.setValue("hbase.hstore.engine.class", CustomTieredStoreEngine.class.getName()); 325 clmBuilder.setValue(TIERING_CELL_QUALIFIER, "date"); 326 TableName tableName = TableName.valueOf("testShouldCompactWhenNoTimeRangeMetadata"); 327 TableDescriptorBuilder tblBuilder = TableDescriptorBuilder.newBuilder(tableName); 328 tblBuilder.setColumnFamily(clmBuilder.build()); 329 utility.getAdmin().createTable(tblBuilder.build()); 330 utility.waitTableAvailable(tableName); 331 Connection connection = utility.getConnection(); 332 Table table = connection.getTable(tableName); 333 long recordTime = System.currentTimeMillis(); 334 // Write data and flush to create store files without CUSTOM_TIERING_TIME_RANGE metadata 335 for (int i = 0; i < 2; i++) { 336 Put put = new Put(Bytes.toBytes(i)); 337 put.addColumn(FAMILY, Bytes.toBytes("val"), Bytes.toBytes("v" + i)); 338 put.addColumn(FAMILY, Bytes.toBytes("date"), Bytes.toBytes(recordTime)); 339 table.put(put); 340 utility.flush(tableName); 341 } 342 table.close(); 343 344 HStore store = 345 (HStore) utility.getMiniHBaseCluster().getRegions(tableName).get(0).getStore(FAMILY); 346 // Verify that flushed files do not have CUSTOM_TIERING_TIME_RANGE metadata 347 for (HStoreFile sf : store.getStorefiles()) { 348 assertNull(sf.getMetadataValue(CUSTOM_TIERING_TIME_RANGE)); 349 } 350 // shouldPerformMajorCompaction must return true due to null timeRangeBytes 351 assertTrue(store.shouldPerformMajorCompaction()); 352 } 353}