Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,88 @@
package org.apache.phoenix.mapreduce;

import org.apache.hadoop.conf.Configuration;
import org.apache.hadoop.mapreduce.Job;
import org.apache.phoenix.compile.QueryPlan;
import org.apache.phoenix.end2end.ParallelStatsDisabledIT;
import org.apache.phoenix.mapreduce.index.IndexScrutinyTool.SourceTable;
import org.apache.phoenix.mapreduce.util.PhoenixConfigurationUtil;
import org.apache.phoenix.schema.PTable;
import org.apache.phoenix.util.PhoenixRuntime;
import org.apache.phoenix.util.PropertiesUtil;
import org.apache.phoenix.util.SchemaUtil;
import org.junit.Test;

import java.sql.Connection;
import java.sql.DriverManager;
import java.util.Properties;

import static org.apache.phoenix.util.TestUtil.TEST_PROPERTIES;
import static org.junit.Assert.assertEquals;

public class PhoenixServerBuildIndexInputFormatIT extends ParallelStatsDisabledIT {

@Test
public void testQueryPlanWithSource() throws Exception {
PhoenixServerBuildIndexInputFormat inputFormat;
Configuration conf = new Configuration(getUtility().getConfiguration());
String schemaName = generateUniqueName();
String dataTableName = generateUniqueName();
String dataTableFullName = SchemaUtil.getTableName(schemaName, dataTableName);
String indexTableName = generateUniqueName();
String indexTableFullName = SchemaUtil.getTableName(schemaName, indexTableName);
String viewName = generateUniqueName();
String viewFullName = SchemaUtil.getTableName(schemaName, viewName);
String viewIndexName = generateUniqueName();
String viewIndexFullName = SchemaUtil.getTableName(schemaName, viewIndexName);
Properties props = PropertiesUtil.deepCopy(TEST_PROPERTIES);
try (Connection conn = DriverManager.getConnection(getUrl(), props)) {
conn.createStatement().execute("CREATE TABLE " + dataTableFullName
+ " (ID INTEGER NOT NULL PRIMARY KEY, VAL1 INTEGER, VAL2 INTEGER) ");
conn.createStatement().execute(String.format(
"CREATE INDEX %s ON %s (VAL1) INCLUDE (VAL2)", indexTableName, dataTableFullName));
conn.createStatement().execute("CREATE VIEW " + viewFullName +
" AS SELECT * FROM " + dataTableFullName);
conn.commit();

conn.createStatement().execute(String.format(
"CREATE INDEX %s ON %s (VAL1) INCLUDE (VAL2)", viewIndexName, viewFullName));

PhoenixConfigurationUtil.setIndexToolDataTableName(conf, dataTableFullName);
PhoenixConfigurationUtil.setIndexToolIndexTableName(conf, indexTableFullName);
// use data table as source (default)
assertTableSource(conf, conn);

// use index table as source
PhoenixConfigurationUtil.setIndexToolSourceTable(conf, SourceTable.INDEX_TABLE_SOURCE);
assertTableSource(conf, conn);

PhoenixConfigurationUtil.setIndexToolDataTableName(conf, viewFullName);
PhoenixConfigurationUtil.setIndexToolIndexTableName(conf, viewIndexFullName);
PhoenixConfigurationUtil.setIndexToolSourceTable(conf, SourceTable.DATA_TABLE_SOURCE);

assertTableSource(conf, conn);

PhoenixConfigurationUtil.setIndexToolSourceTable(conf, SourceTable.INDEX_TABLE_SOURCE);
assertTableSource(conf, conn);
}
}

private void assertTableSource(Configuration conf, Connection conn) throws Exception {
String dataTableFullName = PhoenixConfigurationUtil.getIndexToolDataTableName(conf);
String indexTableFullName = PhoenixConfigurationUtil.getIndexToolIndexTableName(conf);
SourceTable sourceTable = PhoenixConfigurationUtil.getIndexToolSourceTable(conf);
boolean fromIndex = sourceTable.equals(SourceTable.INDEX_TABLE_SOURCE);
PTable pDataTable = PhoenixRuntime.getTable(conn, dataTableFullName);
PTable pIndexTable = PhoenixRuntime.getTable(conn, indexTableFullName);

PhoenixServerBuildIndexInputFormat inputFormat = new PhoenixServerBuildIndexInputFormat();
QueryPlan queryPlan = inputFormat.getQueryPlan(Job.getInstance(), conf);
PTable actual = queryPlan.getTableRef().getTable();

if (!fromIndex) {
assertEquals(pDataTable, actual);
} else {
assertEquals(pIndexTable, actual);
}
}
}
Original file line numberDiff line numberDiff line change
Expand Up@@ -30,7 +30,9 @@
import org.apache.hadoop.mapreduce.InputFormat;
import org.apache.hadoop.mapreduce.JobContext;
import org.apache.hadoop.mapreduce.lib.db.DBWritable;
import org.apache.phoenix.compile.*;
import org.apache.phoenix.compile.MutationPlan;
import org.apache.phoenix.compile.QueryPlan;
import org.apache.phoenix.compile.ServerBuildIndexCompiler;
import org.apache.phoenix.coprocessor.BaseScannerRegionObserver;
import org.apache.phoenix.coprocessor.IndexRebuildRegionScanner;
import org.apache.phoenix.coprocessor.MetaDataProtocol;
Expand All@@ -42,8 +44,12 @@
import org.apache.phoenix.mapreduce.util.ConnectionUtil;
import org.apache.phoenix.mapreduce.util.PhoenixConfigurationUtil;
import org.apache.phoenix.query.QueryServices;
import org.apache.phoenix.schema.*;
import org.apache.phoenix.util.*;
import org.apache.phoenix.schema.PTable;
import org.apache.phoenix.schema.TableRef;
import org.apache.phoenix.util.ByteUtil;
import org.apache.phoenix.util.EnvironmentEdgeManager;
import org.apache.phoenix.util.PhoenixRuntime;
import org.apache.phoenix.util.ScanUtil;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;

Expand Down