View Javadoc

1   /*
2    * Licensed to the Apache Software Foundation (ASF) under one
3    * or more contributor license agreements.  See the NOTICE file
4    * distributed with this work for additional information
5    * regarding copyright ownership.  The ASF licenses this file
6    * to you under the Apache License, Version 2.0 (the
7    * "License"); you may not use this file except in compliance
8    * with the License.  You may obtain a copy of the License at
9    *
10   *   http://www.apache.org/licenses/LICENSE-2.0
11   *
12   * Unless required by applicable law or agreed to in writing, software
13   * distributed under the License is distributed on an "AS IS" BASIS,
14   * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
15   * See the License for the specific language governing permissions and
16   * limitations under the License.
17   */
18  package org.apache.omid.transaction;
19  
20  import static org.apache.omid.metrics.MetricsUtils.name;
21  
22  import java.io.IOException;
23  import java.util.concurrent.ExecutionException;
24  
25  import org.apache.hadoop.hbase.client.Put;
26  import org.apache.hadoop.hbase.util.Bytes;
27  import org.apache.omid.committable.CommitTable;
28  import org.apache.omid.metrics.MetricsRegistry;
29  import org.apache.omid.metrics.Timer;
30  import org.apache.omid.tso.client.CellId;
31  import org.slf4j.Logger;
32  import org.slf4j.LoggerFactory;
33  
34  import com.google.common.util.concurrent.ListenableFuture;
35  import com.google.common.util.concurrent.SettableFuture;
36  
37  public class HBaseSyncPostCommitter implements PostCommitActions {
38  
39      private static final Logger LOG = LoggerFactory.getLogger(HBaseSyncPostCommitter.class);
40  
41      private final MetricsRegistry metrics;
42      private final CommitTable.Client commitTableClient;
43  
44      private final Timer commitTableUpdateTimer;
45      private final Timer shadowCellsUpdateTimer;
46  
47      public HBaseSyncPostCommitter(MetricsRegistry metrics, CommitTable.Client commitTableClient) {
48          this.metrics = metrics;
49          this.commitTableClient = commitTableClient;
50  
51          this.commitTableUpdateTimer = metrics.timer(name("omid", "tm", "hbase", "commitTableUpdate", "latency"));
52          this.shadowCellsUpdateTimer = metrics.timer(name("omid", "tm", "hbase", "shadowCellsUpdate", "latency"));
53      }
54  
55      private void addShadowCell(HBaseCellId cell, HBaseTransaction tx, SettableFuture<Void> updateSCFuture) {
56          Put put = new Put(cell.getRow());
57          put.addColumn(cell.getFamily(),
58                  CellUtils.addShadowCellSuffixPrefix(cell.getQualifier(), 0, cell.getQualifier().length),
59                  cell.getTimestamp(),
60                  Bytes.toBytes(tx.getCommitTimestamp()));
61          try {
62              cell.getTable().getHTable().put(put);
63          } catch (IOException e) {
64              LOG.warn("{}: Error inserting shadow cell {}", tx, cell, e);
65              updateSCFuture.setException(
66                      new TransactionManagerException(tx + ": Error inserting shadow cell " + cell, e));
67          }
68      }
69  
70      @Override
71      public ListenableFuture<Void> updateShadowCells(AbstractTransaction<? extends CellId> transaction) {
72  
73          SettableFuture<Void> updateSCFuture = SettableFuture.create();
74  
75          HBaseTransaction tx = HBaseTransactionManager.enforceHBaseTransactionAsParam(transaction);
76  
77          shadowCellsUpdateTimer.start();
78          try {
79  
80              // Add shadow cells
81              for (HBaseCellId cell : tx.getWriteSet()) {
82                  addShadowCell(cell, tx, updateSCFuture);
83              }
84  
85              for (HBaseCellId cell : tx.getConflictFreeWriteSet()) {
86                  addShadowCell(cell, tx, updateSCFuture);
87              }
88  
89              // Flush affected tables before returning to avoid loss of shadow cells updates when autoflush is disabled
90              try {
91                  tx.flushTables();
92                  updateSCFuture.set(null);
93              } catch (IOException e) {
94                  LOG.warn("{}: Error while flushing writes", tx, e);
95                  updateSCFuture.setException(new TransactionManagerException(tx + ": Error while flushing writes", e));
96              }
97  
98          } finally {
99              shadowCellsUpdateTimer.stop();
100         }
101 
102         return updateSCFuture;
103 
104     }
105 
106     @Override
107     public ListenableFuture<Void> removeCommitTableEntry(AbstractTransaction<? extends CellId> transaction) {
108 
109         SettableFuture<Void> updateSCFuture = SettableFuture.create();
110 
111         HBaseTransaction tx = HBaseTransactionManager.enforceHBaseTransactionAsParam(transaction);
112 
113         commitTableUpdateTimer.start();
114 
115         try {
116             commitTableClient.completeTransaction(tx.getStartTimestamp()).get();
117             updateSCFuture.set(null);
118         } catch (InterruptedException e) {
119             Thread.currentThread().interrupt();
120             LOG.warn("{}: interrupted during commit table entry delete", tx, e);
121             updateSCFuture.setException(
122                     new TransactionManagerException(tx + ": interrupted during commit table entry delete"));
123         } catch (ExecutionException e) {
124             LOG.warn("{}: can't remove commit table entry", tx, e);
125             updateSCFuture.setException(new TransactionManagerException(tx + ": can't remove commit table entry"));
126         } finally {
127             commitTableUpdateTimer.stop();
128         }
129 
130         return updateSCFuture;
131 
132     }
133 
134 }