diff --git a/bin/pherf-cluster.py b/bin/pherf-cluster.py index 37f29a8daa7..04ce50d404a 100755 --- a/bin/pherf-cluster.py +++ b/bin/pherf-cluster.py @@ -75,7 +75,7 @@ stderr=subprocess.PIPE).communicate() java_cmd = java +' -cp "' + hbasecp + os.pathsep + phoenix_utils.pherf_conf_path + os.pathsep + phoenix_utils.hbase_conf_dir + os.pathsep + phoenix_utils.phoenix_pherf_jar + \ - '" -Dlog4j.configuration=file:' + \ + os.pathsep + phoenix_utils.phoenix_thin_client_jar + '" -Dlog4j.configuration=file:' + \ os.path.join(phoenix_utils.current_dir, "log4j.properties") + \ " org.apache.phoenix.pherf.Pherf " + args diff --git a/phoenix-pherf/src/main/java/org/apache/phoenix/pherf/Pherf.java b/phoenix-pherf/src/main/java/org/apache/phoenix/pherf/Pherf.java index 154d6ff7a09..9957d3291fb 100644 --- a/phoenix-pherf/src/main/java/org/apache/phoenix/pherf/Pherf.java +++ b/phoenix-pherf/src/main/java/org/apache/phoenix/pherf/Pherf.java @@ -254,6 +254,7 @@ public void run() throws Exception { // Schema and Data Load if (preLoadData) { logger.info("\nStarting Data Load..."); + System.out.print("Starting Data Load ..."); Workload workload = new WriteWorkload(parser, generateStatistics); try { workloadExecutor.add(workload); diff --git a/phoenix-pherf/src/main/java/org/apache/phoenix/pherf/PherfConstants.java b/phoenix-pherf/src/main/java/org/apache/phoenix/pherf/PherfConstants.java index e7ba056eb5b..7c6693f0323 100644 --- a/phoenix-pherf/src/main/java/org/apache/phoenix/pherf/PherfConstants.java +++ b/phoenix-pherf/src/main/java/org/apache/phoenix/pherf/PherfConstants.java @@ -38,6 +38,7 @@ public static enum CompareType { private static PherfConstants instance = null; private static Properties instanceProperties = null; + public static final String DEFAULT_PRIMARY_KEY = "TENANT_ID"; public static final int DEFAULT_THREAD_POOL_SIZE = 10; public static final int DEFAULT_BATCH_SIZE = 1000; public static final String DEFAULT_DATE_PATTERN = "yyyy-MM-dd HH:mm:ss.SSS"; diff --git a/phoenix-pherf/src/main/java/org/apache/phoenix/pherf/configuration/DataSequence.java b/phoenix-pherf/src/main/java/org/apache/phoenix/pherf/configuration/DataSequence.java index 056a913c1f6..8bce7076014 100644 --- a/phoenix-pherf/src/main/java/org/apache/phoenix/pherf/configuration/DataSequence.java +++ b/phoenix-pherf/src/main/java/org/apache/phoenix/pherf/configuration/DataSequence.java @@ -19,5 +19,5 @@ package org.apache.phoenix.pherf.configuration; public enum DataSequence { - RANDOM, SEQUENTIAL,LIST; + RANDOM, SEQUENTIAL,LIST,SUPERSEQUENTIAL; } \ No newline at end of file diff --git a/phoenix-pherf/src/main/java/org/apache/phoenix/pherf/configuration/DataTypeMapping.java b/phoenix-pherf/src/main/java/org/apache/phoenix/pherf/configuration/DataTypeMapping.java index c266a5743bf..6216b52b301 100644 --- a/phoenix-pherf/src/main/java/org/apache/phoenix/pherf/configuration/DataTypeMapping.java +++ b/phoenix-pherf/src/main/java/org/apache/phoenix/pherf/configuration/DataTypeMapping.java @@ -25,7 +25,8 @@ public enum DataTypeMapping { CHAR("CHAR", Types.CHAR), DECIMAL("DECIMAL", Types.DECIMAL), INTEGER("INTEGER", Types.INTEGER), - DATE("DATE", Types.DATE); + DATE("DATE", Types.DATE), + YCSBKEY("YCSBKEY", Types.VARCHAR); private final String sType; diff --git a/phoenix-pherf/src/main/java/org/apache/phoenix/pherf/configuration/QuerySet.java b/phoenix-pherf/src/main/java/org/apache/phoenix/pherf/configuration/QuerySet.java index 17d415354c4..a83b53ba435 100644 --- a/phoenix-pherf/src/main/java/org/apache/phoenix/pherf/configuration/QuerySet.java +++ b/phoenix-pherf/src/main/java/org/apache/phoenix/pherf/configuration/QuerySet.java @@ -31,7 +31,24 @@ public class QuerySet { private long numberOfExecutions = PherfConstants.DEFAULT_NUMBER_OF_EXECUTIONS; private long executionDurationInMs = PherfConstants.DEFAULT_THREAD_DURATION_IN_MS; private ExecutionType executionType = ExecutionType.SERIAL; - + private boolean randomPointRead = false; + private String primaryKey = PherfConstants.DEFAULT_PRIMARY_KEY; + + @XmlAttribute + public String getPrimaryKey() { + return primaryKey; + } + public void setPrimaryKey(String pk){ + primaryKey=pk; + } + @XmlAttribute + public boolean isRandomPointRead() { + return randomPointRead; + } + public void setRandomPointRead(boolean value) { + randomPointRead = value; + } + /** * List of queries in each query set * @return diff --git a/phoenix-pherf/src/main/java/org/apache/phoenix/pherf/configuration/Scenario.java b/phoenix-pherf/src/main/java/org/apache/phoenix/pherf/configuration/Scenario.java index 200fdc531eb..a26e24b4afc 100644 --- a/phoenix-pherf/src/main/java/org/apache/phoenix/pherf/configuration/Scenario.java +++ b/phoenix-pherf/src/main/java/org/apache/phoenix/pherf/configuration/Scenario.java @@ -27,6 +27,7 @@ import javax.xml.bind.annotation.XmlRootElement; import org.apache.commons.lang.builder.HashCodeBuilder; +import org.apache.commons.lang3.RandomUtils; import org.apache.phoenix.pherf.util.PhoenixUtil; @XmlRootElement(namespace = "org.apache.phoenix.pherf.configuration.DataModel") @@ -126,6 +127,21 @@ public void setDataOverride(DataOverride dataOverride) { * @return */ public List getQuerySet() { + for(QuerySet qs : querySet){ + if(qs.isRandomPointRead()){ + List queryList = new ArrayList(); + String tableName = this.getTableName(); + String primaryKey = qs.getPrimaryKey(); + for(int i = 0; i < 10; i++) { + Query query = new Query(); + long keyInt = RandomUtils.nextLong(100, 10000100); + query.setStatement("select * from " + this.getTableName() + " where " + primaryKey + " = 'user" + keyInt+"'"); + queryList.add(query); + } + + qs.setQuery(queryList); + } + } return querySet; } diff --git a/phoenix-pherf/src/main/java/org/apache/phoenix/pherf/configuration/XMLConfigParser.java b/phoenix-pherf/src/main/java/org/apache/phoenix/pherf/configuration/XMLConfigParser.java index 93dc94cc4a5..828ff359247 100644 --- a/phoenix-pherf/src/main/java/org/apache/phoenix/pherf/configuration/XMLConfigParser.java +++ b/phoenix-pherf/src/main/java/org/apache/phoenix/pherf/configuration/XMLConfigParser.java @@ -168,6 +168,16 @@ private void init(String pattern) throws Exception { } for (Path path : this.paths) { System.out.println("Adding model for path:" + path.toString()); + DataModel dm = XMLConfigParser.readDataModel(path); + List sl = dm.getScenarios(); + for(Scenario s : sl) { + List lqs = s.getQuerySet(); + for(QuerySet qs : lqs) { + System.out.println("Concurrency is " + qs.getConcurrency()); + System.out.println("Exec duration is " + qs.getExecutionDurationInMs()); + System.out.println("RandomPointRead is " + qs.isRandomPointRead()); + } + } this.dataModels.add(XMLConfigParser.readDataModel(path)); } } @@ -176,3 +186,4 @@ private Collection getResources(String pattern) throws Exception { return resourceList.getResourceList(pattern); } } +; \ No newline at end of file diff --git a/phoenix-pherf/src/main/java/org/apache/phoenix/pherf/result/QuerySetResult.java b/phoenix-pherf/src/main/java/org/apache/phoenix/pherf/result/QuerySetResult.java index c2be5a316e9..c109ae74135 100644 --- a/phoenix-pherf/src/main/java/org/apache/phoenix/pherf/result/QuerySetResult.java +++ b/phoenix-pherf/src/main/java/org/apache/phoenix/pherf/result/QuerySetResult.java @@ -32,6 +32,7 @@ public QuerySetResult(QuerySet querySet) { this.setNumberOfExecutions(querySet.getNumberOfExecutions()); this.setExecutionDurationInMs(querySet.getExecutionDurationInMs()); this.setExecutionType(querySet.getExecutionType()); + this.setRandomPointRead(querySet.isRandomPointRead()); } public QuerySetResult() { diff --git a/phoenix-pherf/src/main/java/org/apache/phoenix/pherf/result/ResultManager.java b/phoenix-pherf/src/main/java/org/apache/phoenix/pherf/result/ResultManager.java index 4d4ca4a2f96..d42b0a31a41 100644 --- a/phoenix-pherf/src/main/java/org/apache/phoenix/pherf/result/ResultManager.java +++ b/phoenix-pherf/src/main/java/org/apache/phoenix/pherf/result/ResultManager.java @@ -18,17 +18,20 @@ package org.apache.phoenix.pherf.result; +import org.apache.commons.collections.IteratorUtils; import org.apache.phoenix.pherf.PherfConstants; import org.apache.phoenix.pherf.result.file.ResultFileDetails; import org.apache.phoenix.pherf.result.impl.CSVFileResultHandler; import org.apache.phoenix.pherf.result.impl.ImageResultHandler; import org.apache.phoenix.pherf.result.impl.XMLResultHandler; -import org.apache.phoenix.util.InstanceResolver; + import org.slf4j.Logger; import org.slf4j.LoggerFactory; import java.util.ArrayList; +import java.util.Iterator; import java.util.List; +import java.util.ServiceLoader; public class ResultManager { private static final Logger logger = LoggerFactory.getLogger(ResultManager.class); @@ -71,10 +74,25 @@ public ResultManager(String fileNameSeed) { @SuppressWarnings("unchecked") public ResultManager(String fileNameSeed, boolean writeRuntimeResults) { this(fileNameSeed, writeRuntimeResults ? - InstanceResolver.get(ResultHandler.class, defaultHandlers) : - InstanceResolver.get(ResultHandler.class, minimalHandlers)); + instanceResolverGet(ResultHandler.class, defaultHandlers) : + instanceResolverGet(ResultHandler.class, minimalHandlers)); + } + // Workaround running against phoenix-4.4 not having this method in InstanceResolver. + @SuppressWarnings("unchecked") + public static List instanceResolverGet(Class clazz, List defaultInstances) { + Iterator iterator = ServiceLoader.load(clazz).iterator(); + if (defaultInstances != null) { + defaultInstances.addAll(IteratorUtils.toList(iterator)); + } else { + defaultInstances = IteratorUtils.toList(iterator); + } + + return defaultInstances; + } + + public ResultManager(String fileNameSeed, List resultHandlers) { this.resultHandlers = resultHandlers; util = new ResultUtil(); diff --git a/phoenix-pherf/src/main/java/org/apache/phoenix/pherf/rules/RulesApplier.java b/phoenix-pherf/src/main/java/org/apache/phoenix/pherf/rules/RulesApplier.java index 454050bef1b..e5211ba16f2 100644 --- a/phoenix-pherf/src/main/java/org/apache/phoenix/pherf/rules/RulesApplier.java +++ b/phoenix-pherf/src/main/java/org/apache/phoenix/pherf/rules/RulesApplier.java @@ -139,6 +139,7 @@ public DataValue getDataValue(Column column) throws Exception{ } switch (column.getType()) { + case VARCHAR: // Use the specified data values from configs if they exist if ((column.getDataValues() != null) && (column.getDataValues().size() > 0)) { @@ -157,7 +158,11 @@ public DataValue getDataValue(Column column) throws Exception{ data = pickDataValueFromList(dataValues); } else { Preconditions.checkArgument(length > 0, "length needs to be > 0"); - if (column.getDataSequence() == DataSequence.SEQUENTIAL) { + + if(column.getDataSequence() == DataSequence.SUPERSEQUENTIAL) { + data = getSuperSequentialDataValue(column); + } + else if (column.getDataSequence() == DataSequence.SEQUENTIAL) { data = getSequentialDataValue(column); } else { data = getRandomDataValue(column); @@ -432,4 +437,14 @@ private DataValue getRandomDataValue(Column column) { varchar = StringUtils.left(varchar, column.getLength()); return new DataValue(column.getType(), varchar); } + + private DataValue getSuperSequentialDataValue(Column column) { + DataValue data = null; + long inc = COUNTER.getAndIncrement(); + String strInc = String.valueOf(inc); + String varchar = "user"; + varchar = varchar + strInc; + data = new DataValue(column.getType(), varchar); + return data; + } } diff --git a/phoenix-pherf/src/main/java/org/apache/phoenix/pherf/workload/WriteWorkload.java b/phoenix-pherf/src/main/java/org/apache/phoenix/pherf/workload/WriteWorkload.java index e536eb9e734..e963381a052 100644 --- a/phoenix-pherf/src/main/java/org/apache/phoenix/pherf/workload/WriteWorkload.java +++ b/phoenix-pherf/src/main/java/org/apache/phoenix/pherf/workload/WriteWorkload.java @@ -23,6 +23,7 @@ import java.sql.Date; import java.sql.PreparedStatement; import java.sql.SQLException; +import java.sql.Statement; import java.sql.Types; import java.text.SimpleDateFormat; import java.util.ArrayList; @@ -242,6 +243,7 @@ public Future upsertData(final Scenario scenario, final List colum SimpleDateFormat simpleDateFormat = new SimpleDateFormat("yyyy-MM-dd HH:mm:ss"); Connection connection = null; PreparedStatement stmt = null; + try { connection = pUtil.getConnection(scenario.getTenantId()); long logStartTime = System.currentTimeMillis(); @@ -253,13 +255,22 @@ public Future upsertData(final Scenario scenario, final List colum last = start = System.currentTimeMillis(); String sql = buildSql(columns, tableName); + //System.out.print("going to create statement"); + + //System.out.print("statement created"); stmt = connection.prepareStatement(sql); for (long i = rowCount; (i > 0) && ((System.currentTimeMillis() - logStartTime) < maxDuration); i--) { - stmt = buildStatement(scenario, columns, stmt, simpleDateFormat); - rowsCreated += stmt.executeUpdate(); + + stmt = buildStatement(scenario, columns, stmt, simpleDateFormat); + //System.out.print(stmt); + stmt.addBatch(); + //rowsCreated += stmt.executeUpdate(); if ((i % getBatchSize()) == 0) { + //System.out.print("Executing batch"); + stmt.executeBatch(); connection.commit(); + //stmt.clearBatch(); duration = System.currentTimeMillis() - last; logger.info("Writer (" + Thread.currentThread().getName() + ") committed Batch. Total " + getBatchSize() @@ -281,6 +292,8 @@ public Future upsertData(final Scenario scenario, final List colum } } finally { if (stmt != null) { + System.out.print("Executing batch"); + stmt.executeBatch(); stmt.close(); } if (connection != null) { @@ -307,7 +320,7 @@ private PreparedStatement buildStatement(Scenario scenario, List columns PreparedStatement statement, SimpleDateFormat simpleDateFormat) throws Exception { int count = 1; for (Column column : columns) { - + //System.out.print("Column is " + column.getName()); DataValue dataValue = getRulesApplier().getDataForRule(scenario, column); switch (column.getType()) { case VARCHAR: @@ -318,6 +331,7 @@ private PreparedStatement buildStatement(Scenario scenario, List columns } break; case CHAR: + //System.out.print("Data value "+ dataValue.getValue()); if (dataValue.getValue().equals("")) { statement.setNull(count, Types.CHAR); } else { @@ -354,6 +368,7 @@ private PreparedStatement buildStatement(Scenario scenario, List columns } count++; } + //System.out.print("Returning statement " + statement.toString()); return statement; }