mirror of
https://github.com/wahyd4/pinpoint.git
synced 2026-08-16 08:16:15 +10:00
[강운덕] [LUCYSUS-1744] traceIndex 분산 hash를 위해 분산 hash키를 추가함. 32사이즈로 되어 있음.
추가 코드정리 및 분산키를 좀더 튜닝해야 될지 추가 판단해야함. git-svn-id: http://svn.bds.nhncorp.com/pe/hippo-commons/trunk@2104 84d0f5b1-2673-498c-a247-62c4ff18d310
This commit is contained in:
@@ -173,6 +173,13 @@
|
||||
</exclusions>
|
||||
</dependency>
|
||||
|
||||
<dependency>
|
||||
<groupId>com.sematext.hbasewd</groupId>
|
||||
<artifactId>hbasewd</artifactId>
|
||||
<version>0.1.0</version>
|
||||
<scope>provided</scope>
|
||||
</dependency>
|
||||
|
||||
<dependency>
|
||||
<groupId>org.apache.thrift</groupId>
|
||||
<artifactId>libthrift</artifactId>
|
||||
|
||||
@@ -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 {
|
||||
|
||||
<T> List<List<T>> find(String tableName, final List<Scan> scans, final RowMapper<T> action);
|
||||
|
||||
<T> List<T> find(String tableName, final Scan scan, AbstractRowKeyDistributor rowKeyDistributor, final RowMapper<T> action);
|
||||
|
||||
<T> T find(String tableName, final Scan scan, final AbstractRowKeyDistributor rowKeyDistributor, final ResultsExtractor<T> action);
|
||||
|
||||
void increment(String tableName, final Increment increment);
|
||||
|
||||
void incrementColumnValue(String tableName, final byte[] rowName, final byte[] familyName, final byte[] qualifier, final long amount);
|
||||
|
||||
@@ -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<T>(action));
|
||||
}
|
||||
|
||||
public <T> List<T> find(String tableName, final Scan scan, final AbstractRowKeyDistributor rowKeyDistributor, final RowMapper<T> action) {
|
||||
final RowMapperResultsExtractor<T> resultsExtractor = new RowMapperResultsExtractor<T>(action);
|
||||
return execute(tableName, new TableCallback<List<T>>() {
|
||||
@Override
|
||||
public List<T> doInTable(HTableInterface htable) throws Throwable {
|
||||
ResultScanner scanner = createDistributeScanner(htable, scan, rowKeyDistributor);
|
||||
try {
|
||||
return resultsExtractor.extractData(scanner);
|
||||
} finally {
|
||||
scanner.close();
|
||||
}
|
||||
}
|
||||
});
|
||||
}
|
||||
|
||||
@Override
|
||||
public <T> T find(String tableName, final Scan scan, final AbstractRowKeyDistributor rowKeyDistributor, final ResultsExtractor<T> action) {
|
||||
|
||||
return execute(tableName, new TableCallback<T>() {
|
||||
@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
|
||||
|
||||
Reference in New Issue
Block a user