修复:

1. 第一次MapReduce没有统计结果的问题
2. Summary表中的Rowkey添加上UserId的标识
This commit is contained in:
ascrutae 2016-02-22 14:34:36 +08:00
parent 683ab87924
commit 0e83844686
1 changed files with 10 additions and 9 deletions

View File

@ -3,17 +3,19 @@ package com.ai.cloud.skywalking.analysis.categorize2chain;
import com.ai.cloud.skywalking.analysis.categorize2chain.model.ChainInfo;
import com.ai.cloud.skywalking.analysis.categorize2chain.model.ChainNode;
import com.ai.cloud.skywalking.analysis.categorize2chain.util.HBaseUtil;
import org.apache.hadoop.hbase.client.Put;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import java.io.IOException;
import java.sql.*;
import java.sql.SQLException;
import java.sql.Timestamp;
import java.text.SimpleDateFormat;
import java.util.*;
import java.util.Date;
public class ChainSummary {
private static Logger logger = LoggerFactory.getLogger(ChainSummary.class.getName());
private Map<String, ChainSpecificTimeWindowSummary> loadedChainSpecificTimeWindowSummary;
private Map<String, Timestamp> updateChainInfo;
@ -24,19 +26,17 @@ public class ChainSummary {
public void summary(ChainInfo chainInfo) {
for (ChainNode node : chainInfo.getNodes()) {
String csk = generateChainSummaryKey(chainInfo.getCID(), node.getStartDate());
if (loadedChainSpecificTimeWindowSummary.containsKey(csk)) {
String csk = generateChainSummaryKey(chainInfo, node.getStartDate());
if (!loadedChainSpecificTimeWindowSummary.containsKey(csk)) {
loadedChainSpecificTimeWindowSummary.put(csk, ChainSpecificTimeWindowSummary.load(csk));
}
loadedChainSpecificTimeWindowSummary.get(csk).summaryNodeValue(node);
}
updateChainInfo.put(chainInfo.getCID(), new Timestamp(System.currentTimeMillis()));
}
private String generateChainSummaryKey(String chainToken, long startDate) {
return chainToken + "-" + new SimpleDateFormat("yyyy/MM/dd HH:mm:ss").
private String generateChainSummaryKey(ChainInfo chainInfo, long startDate) {
return chainInfo.getCID() + "-" + chainInfo.getUserId() + "-" + new SimpleDateFormat("yyyy/MM/dd HH:mm:ss").
format(new Date(startDate / (1000 * 60 * 60) * (1000 * 60 * 60)));
}
@ -51,6 +51,7 @@ public class ChainSummary {
private void batchSaveChainSpecificTimeWindowSummary() throws IOException, InterruptedException {
List<Put> puts = new ArrayList<Put>();
logger.info("There are [" + loadedChainSpecificTimeWindowSummary.size() + "] summary data will be storage to HBase");
for (Map.Entry<String, ChainSpecificTimeWindowSummary> entry : loadedChainSpecificTimeWindowSummary.entrySet()) {
Put put = new Put(entry.getKey().getBytes());
entry.getValue().save(put);