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}