diff --git a/pom.xml b/pom.xml index a2aad9d65..8389fab6c 100644 --- a/pom.xml +++ b/pom.xml @@ -173,6 +173,13 @@ + + com.sematext.hbasewd + hbasewd + 0.1.0 + provided + + org.apache.thrift libthrift diff --git a/src/main/java/com/nhn/pinpoint/common/hbase/HbaseOperations2.java b/src/main/java/com/nhn/pinpoint/common/hbase/HbaseOperations2.java index 702c79c65..f82c90495 100644 --- a/src/main/java/com/nhn/pinpoint/common/hbase/HbaseOperations2.java +++ b/src/main/java/com/nhn/pinpoint/common/hbase/HbaseOperations2.java @@ -2,6 +2,7 @@ package com.nhn.pinpoint.common.hbase; import java.util.List; +import com.sematext.hbase.wd.AbstractRowKeyDistributor; import org.apache.hadoop.hbase.client.Delete; import org.apache.hadoop.hbase.client.Get; import org.apache.hadoop.hbase.client.Increment; @@ -73,6 +74,10 @@ public interface HbaseOperations2 extends HbaseOperations { List> find(String tableName, final List scans, final RowMapper action); + List find(String tableName, final Scan scan, AbstractRowKeyDistributor rowKeyDistributor, final RowMapper action); + + T find(String tableName, final Scan scan, final AbstractRowKeyDistributor rowKeyDistributor, final ResultsExtractor action); + void increment(String tableName, final Increment increment); void incrementColumnValue(String tableName, final byte[] rowName, final byte[] familyName, final byte[] qualifier, final long amount); diff --git a/src/main/java/com/nhn/pinpoint/common/hbase/HbaseTemplate2.java b/src/main/java/com/nhn/pinpoint/common/hbase/HbaseTemplate2.java index 3afbb2fc1..d68e82e23 100644 --- a/src/main/java/com/nhn/pinpoint/common/hbase/HbaseTemplate2.java +++ b/src/main/java/com/nhn/pinpoint/common/hbase/HbaseTemplate2.java @@ -1,5 +1,7 @@ package com.nhn.pinpoint.common.hbase; +import com.sematext.hbase.wd.AbstractRowKeyDistributor; +import com.sematext.hbase.wd.DistributedScanner; import org.apache.hadoop.conf.Configuration; import org.apache.hadoop.hbase.client.*; import org.slf4j.Logger; @@ -9,6 +11,7 @@ import org.springframework.beans.factory.InitializingBean; import org.springframework.data.hadoop.hbase.*; import org.springframework.util.Assert; +import java.io.IOException; import java.util.ArrayList; import java.util.Collections; import java.util.List; @@ -386,6 +389,70 @@ public class HbaseTemplate2 extends HbaseTemplate implements HbaseOperations2, I return find(tableName, scans, new RowMapperResultsExtractor(action)); } + public List find(String tableName, final Scan scan, final AbstractRowKeyDistributor rowKeyDistributor, final RowMapper action) { + final RowMapperResultsExtractor resultsExtractor = new RowMapperResultsExtractor(action); + return execute(tableName, new TableCallback>() { + @Override + public List doInTable(HTableInterface htable) throws Throwable { + ResultScanner scanner = createDistributeScanner(htable, scan, rowKeyDistributor); + try { + return resultsExtractor.extractData(scanner); + } finally { + scanner.close(); + } + } + }); + } + + @Override + public T find(String tableName, final Scan scan, final AbstractRowKeyDistributor rowKeyDistributor, final ResultsExtractor action) { + + return execute(tableName, new TableCallback() { + @Override + public T doInTable(HTableInterface htable) throws Throwable { + boolean debugEnabled = logger.isDebugEnabled(); + long beforeCreateDistributeScan = 0; + if (debugEnabled) { + beforeCreateDistributeScan = System.currentTimeMillis(); + } + ResultScanner scanner = createDistributeScanner(htable, scan, rowKeyDistributor); + if (debugEnabled) { + logger.debug("DistributeScanner createTime:{}", (System.currentTimeMillis() - beforeCreateDistributeScan)); + } + long beforeDistributeScan = 0; + if (debugEnabled) { + beforeDistributeScan = System.currentTimeMillis(); + } + try { + return action.extractData(scanner); + } finally { + scanner.close(); + if (debugEnabled) { + logger.debug("DistributeScanner scanTime:{}", (System.currentTimeMillis() - beforeDistributeScan)); + } + } + } + }); + } + + public ResultScanner createDistributeScanner(HTableInterface htable, Scan originalScan, AbstractRowKeyDistributor rowKeyDistributor) throws IOException { + + Scan[] scans = rowKeyDistributor.getDistributedScans(originalScan); + for(int i = 0; i < scans.length; i++) { + Scan scan = scans[i]; + scan.setId(originalScan.getId() + "-" + i); + // caching만 넣으면 되나? + scan.setCaching(originalScan.getCaching()); + } + + ResultScanner[] scanner = new ResultScanner[scans.length]; + for (int i = 0; i < scans.length; i++) { + scanner[i] = htable.getScanner(scans[i]); + } + + return new DistributedScanner(rowKeyDistributor, scanner); + } + public void increment(String tableName, final Increment increment) { execute(tableName, new TableCallback() { @Override