From 9ed26d936b066be1c17f9031444d225a9e94de26 Mon Sep 17 00:00:00 2001 From: Woonduk Kang Date: Wed, 7 Aug 2013 10:03:05 +0000 Subject: [PATCH] =?UTF-8?q?[=EA=B0=95=EC=9A=B4=EB=8D=95]=20[LUCYSUS-1744]?= =?UTF-8?q?=20traceIndex=20=EB=B6=84=EC=82=B0=20hash=EB=A5=BC=20=EC=9C=84?= =?UTF-8?q?=ED=95=B4=20=EB=B6=84=EC=82=B0=20hash=ED=82=A4=EB=A5=BC=20?= =?UTF-8?q?=EC=B6=94=EA=B0=80=ED=95=A8.=2032=EC=82=AC=EC=9D=B4=EC=A6=88?= =?UTF-8?q?=EB=A1=9C=20=EB=90=98=EC=96=B4=20=EC=9E=88=EC=9D=8C.=20?= =?UTF-8?q?=EC=B6=94=EA=B0=80=20=EC=BD=94=EB=93=9C=EC=A0=95=EB=A6=AC=20?= =?UTF-8?q?=EB=B0=8F=20=EB=B6=84=EC=82=B0=ED=82=A4=EB=A5=BC=20=EC=A2=80?= =?UTF-8?q?=EB=8D=94=20=ED=8A=9C=EB=8B=9D=ED=95=B4=EC=95=BC=20=EB=90=A0?= =?UTF-8?q?=EC=A7=80=20=EC=B6=94=EA=B0=80=20=ED=8C=90=EB=8B=A8=ED=95=B4?= =?UTF-8?q?=EC=95=BC=ED=95=A8.?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit git-svn-id: http://svn.bds.nhncorp.com/pe/hippo-commons/trunk@2104 84d0f5b1-2673-498c-a247-62c4ff18d310 --- pom.xml | 7 ++ .../common/hbase/HbaseOperations2.java | 5 ++ .../pinpoint/common/hbase/HbaseTemplate2.java | 67 +++++++++++++++++++ 3 files changed, 79 insertions(+) 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