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 com.google.common.util.concurrent.ListenableFuture;
21  import com.google.common.util.concurrent.SettableFuture;
22  import org.apache.omid.committable.CommitTable;
23  import org.apache.omid.metrics.MetricsRegistry;
24  import org.apache.omid.metrics.Timer;
25  import org.apache.omid.tso.client.CellId;
26  import org.apache.hadoop.hbase.client.Put;
27  import org.apache.hadoop.hbase.util.Bytes;
28  import org.slf4j.Logger;
29  import org.slf4j.LoggerFactory;
30  
31  import java.io.IOException;
32  import java.util.concurrent.ExecutionException;
33  
34  import static org.apache.omid.metrics.MetricsUtils.name;
35  
36  public class HBaseSyncPostCommitter implements PostCommitActions {
37  
38      private static final Logger LOG = LoggerFactory.getLogger(HBaseSyncPostCommitter.class);
39  
40      private final MetricsRegistry metrics;
41      private final CommitTable.Client commitTableClient;
42  
43      private final Timer commitTableUpdateTimer;
44      private final Timer shadowCellsUpdateTimer;
45  
46      public HBaseSyncPostCommitter(MetricsRegistry metrics, CommitTable.Client commitTableClient) {
47          this.metrics = metrics;
48          this.commitTableClient = commitTableClient;
49  
50          this.commitTableUpdateTimer = metrics.timer(name("omid", "tm", "hbase", "commitTableUpdate", "latency"));
51          this.shadowCellsUpdateTimer = metrics.timer(name("omid", "tm", "hbase", "shadowCellsUpdate", "latency"));
52      }
53  
54      @Override
55      public ListenableFuture<Void> updateShadowCells(AbstractTransaction<? extends CellId> transaction) {
56  
57          SettableFuture<Void> updateSCFuture = SettableFuture.create();
58  
59          HBaseTransaction tx = HBaseTransactionManager.enforceHBaseTransactionAsParam(transaction);
60  
61          shadowCellsUpdateTimer.start();
62          try {
63  
64              // Add shadow cells
65              for (HBaseCellId cell : tx.getWriteSet()) {
66                  Put put = new Put(cell.getRow());
67                  put.add(cell.getFamily(),
68                          CellUtils.addShadowCellSuffix(cell.getQualifier(), 0, cell.getQualifier().length),
69                          tx.getStartTimestamp(),
70                          Bytes.toBytes(tx.getCommitTimestamp()));
71                  try {
72                      cell.getTable().put(put);
73                  } catch (IOException e) {
74                      LOG.warn("{}: Error inserting shadow cell {}", tx, cell, e);
75                      updateSCFuture.setException(
76                              new TransactionManagerException(tx + ": Error inserting shadow cell " + cell, e));
77                  }
78              }
79  
80              // Flush affected tables before returning to avoid loss of shadow cells updates when autoflush is disabled
81              try {
82                  tx.flushTables();
83                  updateSCFuture.set(null);
84              } catch (IOException e) {
85                  LOG.warn("{}: Error while flushing writes", tx, e);
86                  updateSCFuture.setException(new TransactionManagerException(tx + ": Error while flushing writes", e));
87              }
88  
89          } finally {
90              shadowCellsUpdateTimer.stop();
91          }
92  
93          return updateSCFuture;
94  
95      }
96  
97      @Override
98      public ListenableFuture<Void> removeCommitTableEntry(AbstractTransaction<? extends CellId> transaction) {
99  
100         SettableFuture<Void> updateSCFuture = SettableFuture.create();
101 
102         HBaseTransaction tx = HBaseTransactionManager.enforceHBaseTransactionAsParam(transaction);
103 
104         commitTableUpdateTimer.start();
105 
106         try {
107             commitTableClient.completeTransaction(tx.getStartTimestamp()).get();
108             updateSCFuture.set(null);
109         } catch (InterruptedException e) {
110             Thread.currentThread().interrupt();
111             LOG.warn("{}: interrupted during commit table entry delete", tx, e);
112             updateSCFuture.setException(
113                     new TransactionManagerException(tx + ": interrupted during commit table entry delete"));
114         } catch (ExecutionException e) {
115             LOG.warn("{}: can't remove commit table entry", tx, e);
116             updateSCFuture.setException(new TransactionManagerException(tx + ": can't remove commit table entry"));
117         } finally {
118             commitTableUpdateTimer.stop();
119         }
120 
121         return updateSCFuture;
122 
123     }
124 
125 }