Index: ql/src/test/results/clientpositive/smb_mapjoin_11.q.out =================================================================== --- ql/src/test/results/clientpositive/smb_mapjoin_11.q.out (revision 0) +++ ql/src/test/results/clientpositive/smb_mapjoin_11.q.out (revision 0) @@ -0,0 +1,341 @@ +PREHOOK: query: -- This test verifies that the output of a sort merge join on 2 partitions (one on each side of the join) is bucketed + +-- Create two bucketed and sorted tables +CREATE TABLE test_table1 (key INT, value STRING) PARTITIONED BY (ds STRING) CLUSTERED BY (key) SORTED BY (key) INTO 16 BUCKETS +PREHOOK: type: CREATETABLE +POSTHOOK: query: -- This test verifies that the output of a sort merge join on 2 partitions (one on each side of the join) is bucketed + +-- Create two bucketed and sorted tables +CREATE TABLE test_table1 (key INT, value STRING) PARTITIONED BY (ds STRING) CLUSTERED BY (key) SORTED BY (key) INTO 16 BUCKETS +POSTHOOK: type: CREATETABLE +POSTHOOK: Output: default@test_table1 +PREHOOK: query: CREATE TABLE test_table2 (key INT, value STRING) PARTITIONED BY (ds STRING) CLUSTERED BY (key) SORTED BY (key) INTO 16 BUCKETS +PREHOOK: type: CREATETABLE +POSTHOOK: query: CREATE TABLE test_table2 (key INT, value STRING) PARTITIONED BY (ds STRING) CLUSTERED BY (key) SORTED BY (key) INTO 16 BUCKETS +POSTHOOK: type: CREATETABLE +POSTHOOK: Output: default@test_table2 +PREHOOK: query: FROM src +INSERT OVERWRITE TABLE test_table1 PARTITION (ds = '1') SELECT * +INSERT OVERWRITE TABLE test_table2 PARTITION (ds = '1') SELECT * +PREHOOK: type: QUERY +PREHOOK: Input: default@src +PREHOOK: Output: default@test_table1@ds=1 +PREHOOK: Output: default@test_table2@ds=1 +POSTHOOK: query: FROM src +INSERT OVERWRITE TABLE test_table1 PARTITION (ds = '1') SELECT * +INSERT OVERWRITE TABLE test_table2 PARTITION (ds = '1') SELECT * +POSTHOOK: type: QUERY +POSTHOOK: Input: default@src +POSTHOOK: Output: default@test_table1@ds=1 +POSTHOOK: Output: default@test_table2@ds=1 +POSTHOOK: Lineage: test_table1 PARTITION(ds=1).key EXPRESSION [(src)src.FieldSchema(name:key, type:string, comment:default), ] +POSTHOOK: Lineage: test_table1 PARTITION(ds=1).value SIMPLE [(src)src.FieldSchema(name:value, type:string, comment:default), ] +POSTHOOK: Lineage: test_table2 PARTITION(ds=1).key EXPRESSION [(src)src.FieldSchema(name:key, type:string, comment:default), ] +POSTHOOK: Lineage: test_table2 PARTITION(ds=1).value SIMPLE [(src)src.FieldSchema(name:value, type:string, comment:default), ] +PREHOOK: query: -- Create a bucketed table +CREATE TABLE test_table3 (key INT, value STRING) PARTITIONED BY (ds STRING) CLUSTERED BY (key) INTO 16 BUCKETS +PREHOOK: type: CREATETABLE +POSTHOOK: query: -- Create a bucketed table +CREATE TABLE test_table3 (key INT, value STRING) PARTITIONED BY (ds STRING) CLUSTERED BY (key) INTO 16 BUCKETS +POSTHOOK: type: CREATETABLE +POSTHOOK: Output: default@test_table3 +POSTHOOK: Lineage: test_table1 PARTITION(ds=1).key EXPRESSION [(src)src.FieldSchema(name:key, type:string, comment:default), ] +POSTHOOK: Lineage: test_table1 PARTITION(ds=1).value SIMPLE [(src)src.FieldSchema(name:value, type:string, comment:default), ] +POSTHOOK: Lineage: test_table2 PARTITION(ds=1).key EXPRESSION [(src)src.FieldSchema(name:key, type:string, comment:default), ] +POSTHOOK: Lineage: test_table2 PARTITION(ds=1).value SIMPLE [(src)src.FieldSchema(name:value, type:string, comment:default), ] +PREHOOK: query: -- Insert data into the bucketed table by joining the two bucketed and sorted tables, bucketing is not enforced +EXPLAIN EXTENDED +INSERT OVERWRITE TABLE test_table3 PARTITION (ds = '1') SELECT /*+ MAPJOIN(b) */ a.key, b.value FROM test_table1 a JOIN test_table2 b ON a.key = b.key AND a.ds = '1' AND b.ds = '1' +PREHOOK: type: QUERY +POSTHOOK: query: -- Insert data into the bucketed table by joining the two bucketed and sorted tables, bucketing is not enforced +EXPLAIN EXTENDED +INSERT OVERWRITE TABLE test_table3 PARTITION (ds = '1') SELECT /*+ MAPJOIN(b) */ a.key, b.value FROM test_table1 a JOIN test_table2 b ON a.key = b.key AND a.ds = '1' AND b.ds = '1' +POSTHOOK: type: QUERY +POSTHOOK: Lineage: test_table1 PARTITION(ds=1).key EXPRESSION [(src)src.FieldSchema(name:key, type:string, comment:default), ] +POSTHOOK: Lineage: test_table1 PARTITION(ds=1).value SIMPLE [(src)src.FieldSchema(name:value, type:string, comment:default), ] +POSTHOOK: Lineage: test_table2 PARTITION(ds=1).key EXPRESSION [(src)src.FieldSchema(name:key, type:string, comment:default), ] +POSTHOOK: Lineage: test_table2 PARTITION(ds=1).value SIMPLE [(src)src.FieldSchema(name:value, type:string, comment:default), ] +ABSTRACT SYNTAX TREE: + (TOK_QUERY (TOK_FROM (TOK_JOIN (TOK_TABREF (TOK_TABNAME test_table1) a) (TOK_TABREF (TOK_TABNAME test_table2) b) (AND (AND (= (. (TOK_TABLE_OR_COL a) key) (. (TOK_TABLE_OR_COL b) key)) (= (. (TOK_TABLE_OR_COL a) ds) '1')) (= (. (TOK_TABLE_OR_COL b) ds) '1')))) (TOK_INSERT (TOK_DESTINATION (TOK_TAB (TOK_TABNAME test_table3) (TOK_PARTSPEC (TOK_PARTVAL ds '1')))) (TOK_SELECT (TOK_HINTLIST (TOK_HINT TOK_MAPJOIN (TOK_HINTARGLIST b))) (TOK_SELEXPR (. (TOK_TABLE_OR_COL a) key)) (TOK_SELEXPR (. (TOK_TABLE_OR_COL b) value))))) + +STAGE DEPENDENCIES: + Stage-1 is a root stage + Stage-0 depends on stages: Stage-1 + Stage-2 depends on stages: Stage-0 + +STAGE PLANS: + Stage: Stage-1 + Map Reduce + Alias -> Map Operator Tree: + a + TableScan + alias: a + GatherStats: false + Sorted Merge Bucket Map Join Operator + condition map: + Inner Join 0 to 1 + condition expressions: + 0 {key} + 1 {value} + handleSkewJoin: false + keys: + 0 [Column[key]] + 1 [Column[key]] + outputColumnNames: _col0, _col6 + Position of Big Table: 0 + Select Operator + expressions: + expr: _col0 + type: int + expr: _col6 + type: string + outputColumnNames: _col0, _col6 + Select Operator + expressions: + expr: _col0 + type: int + expr: _col6 + type: string + outputColumnNames: _col0, _col1 + File Output Operator + compressed: false + GlobalTableId: 1 +#### A masked pattern was here #### + NumFilesPerFileSink: 1 + Static Partition Specification: ds=1/ +#### A masked pattern was here #### + table: + input format: org.apache.hadoop.mapred.TextInputFormat + output format: org.apache.hadoop.hive.ql.io.HiveIgnoreKeyTextOutputFormat + properties: + bucket_count 16 + bucket_field_name key + columns key,value + columns.types int:string +#### A masked pattern was here #### + name default.test_table3 + partition_columns ds + serialization.ddl struct test_table3 { i32 key, string value} + serialization.format 1 + serialization.lib org.apache.hadoop.hive.serde2.lazy.LazySimpleSerDe +#### A masked pattern was here #### + serde: org.apache.hadoop.hive.serde2.lazy.LazySimpleSerDe + name: default.test_table3 + TotalFiles: 1 + GatherStats: true + MultiFileSpray: false + Needs Tagging: false + Path -> Alias: +#### A masked pattern was here #### + Path -> Partition: +#### A masked pattern was here #### + Partition + base file name: ds=1 + input format: org.apache.hadoop.mapred.TextInputFormat + output format: org.apache.hadoop.hive.ql.io.HiveIgnoreKeyTextOutputFormat + partition values: + ds 1 + properties: + SORTBUCKETCOLSPREFIX TRUE + bucket_count 16 + bucket_field_name key + columns key,value + columns.types int:string +#### A masked pattern was here #### + name default.test_table1 + numFiles 16 + numPartitions 1 + numRows 500 + partition_columns ds + rawDataSize 5312 + serialization.ddl struct test_table1 { i32 key, string value} + serialization.format 1 + serialization.lib org.apache.hadoop.hive.serde2.lazy.LazySimpleSerDe + totalSize 5812 +#### A masked pattern was here #### + serde: org.apache.hadoop.hive.serde2.lazy.LazySimpleSerDe + + input format: org.apache.hadoop.mapred.TextInputFormat + output format: org.apache.hadoop.hive.ql.io.HiveIgnoreKeyTextOutputFormat + properties: + SORTBUCKETCOLSPREFIX TRUE + bucket_count 16 + bucket_field_name key + columns key,value + columns.types int:string +#### A masked pattern was here #### + name default.test_table1 + numFiles 16 + numPartitions 1 + numRows 500 + partition_columns ds + rawDataSize 5312 + serialization.ddl struct test_table1 { i32 key, string value} + serialization.format 1 + serialization.lib org.apache.hadoop.hive.serde2.lazy.LazySimpleSerDe + totalSize 5812 +#### A masked pattern was here #### + serde: org.apache.hadoop.hive.serde2.lazy.LazySimpleSerDe + name: default.test_table1 + name: default.test_table1 + + Stage: Stage-0 + Move Operator + tables: + partition: + ds 1 + replace: true +#### A masked pattern was here #### + table: + input format: org.apache.hadoop.mapred.TextInputFormat + output format: org.apache.hadoop.hive.ql.io.HiveIgnoreKeyTextOutputFormat + properties: + bucket_count 16 + bucket_field_name key + columns key,value + columns.types int:string +#### A masked pattern was here #### + name default.test_table3 + partition_columns ds + serialization.ddl struct test_table3 { i32 key, string value} + serialization.format 1 + serialization.lib org.apache.hadoop.hive.serde2.lazy.LazySimpleSerDe +#### A masked pattern was here #### + serde: org.apache.hadoop.hive.serde2.lazy.LazySimpleSerDe + name: default.test_table3 +#### A masked pattern was here #### + + Stage: Stage-2 + Stats-Aggr Operator +#### A masked pattern was here #### + + +PREHOOK: query: INSERT OVERWRITE TABLE test_table3 PARTITION (ds = '1') SELECT /*+ MAPJOIN(b) */ a.key, b.value FROM test_table1 a JOIN test_table2 b ON a.key = b.key AND a.ds = '1' AND b.ds = '1' +PREHOOK: type: QUERY +PREHOOK: Input: default@test_table1@ds=1 +PREHOOK: Input: default@test_table2@ds=1 +PREHOOK: Output: default@test_table3@ds=1 +POSTHOOK: query: INSERT OVERWRITE TABLE test_table3 PARTITION (ds = '1') SELECT /*+ MAPJOIN(b) */ a.key, b.value FROM test_table1 a JOIN test_table2 b ON a.key = b.key AND a.ds = '1' AND b.ds = '1' +POSTHOOK: type: QUERY +POSTHOOK: Input: default@test_table1@ds=1 +POSTHOOK: Input: default@test_table2@ds=1 +POSTHOOK: Output: default@test_table3@ds=1 +POSTHOOK: Lineage: test_table1 PARTITION(ds=1).key EXPRESSION [(src)src.FieldSchema(name:key, type:string, comment:default), ] +POSTHOOK: Lineage: test_table1 PARTITION(ds=1).value SIMPLE [(src)src.FieldSchema(name:value, type:string, comment:default), ] +POSTHOOK: Lineage: test_table2 PARTITION(ds=1).key EXPRESSION [(src)src.FieldSchema(name:key, type:string, comment:default), ] +POSTHOOK: Lineage: test_table2 PARTITION(ds=1).value SIMPLE [(src)src.FieldSchema(name:value, type:string, comment:default), ] +POSTHOOK: Lineage: test_table3 PARTITION(ds=1).key SIMPLE [(test_table1)a.FieldSchema(name:key, type:int, comment:null), ] +POSTHOOK: Lineage: test_table3 PARTITION(ds=1).value SIMPLE [(test_table2)b.FieldSchema(name:value, type:string, comment:null), ] +PREHOOK: query: -- Read data from a sampled bucket to verify the data is bucketed +SELECT * FROM test_table3 TABLESAMPLE(BUCKET 2 OUT OF 16) where ds = '1' +PREHOOK: type: QUERY +PREHOOK: Input: default@test_table3@ds=1 +#### A masked pattern was here #### +POSTHOOK: query: -- Read data from a sampled bucket to verify the data is bucketed +SELECT * FROM test_table3 TABLESAMPLE(BUCKET 2 OUT OF 16) where ds = '1' +POSTHOOK: type: QUERY +POSTHOOK: Input: default@test_table3@ds=1 +#### A masked pattern was here #### +POSTHOOK: Lineage: test_table1 PARTITION(ds=1).key EXPRESSION [(src)src.FieldSchema(name:key, type:string, comment:default), ] +POSTHOOK: Lineage: test_table1 PARTITION(ds=1).value SIMPLE [(src)src.FieldSchema(name:value, type:string, comment:default), ] +POSTHOOK: Lineage: test_table2 PARTITION(ds=1).key EXPRESSION [(src)src.FieldSchema(name:key, type:string, comment:default), ] +POSTHOOK: Lineage: test_table2 PARTITION(ds=1).value SIMPLE [(src)src.FieldSchema(name:value, type:string, comment:default), ] +POSTHOOK: Lineage: test_table3 PARTITION(ds=1).key SIMPLE [(test_table1)a.FieldSchema(name:key, type:int, comment:null), ] +POSTHOOK: Lineage: test_table3 PARTITION(ds=1).value SIMPLE [(test_table2)b.FieldSchema(name:value, type:string, comment:null), ] +17 val_17 1 +33 val_33 1 +65 val_65 1 +97 val_97 1 +97 val_97 1 +97 val_97 1 +97 val_97 1 +113 val_113 1 +113 val_113 1 +113 val_113 1 +113 val_113 1 +129 val_129 1 +129 val_129 1 +129 val_129 1 +129 val_129 1 +145 val_145 1 +177 val_177 1 +193 val_193 1 +193 val_193 1 +193 val_193 1 +193 val_193 1 +193 val_193 1 +193 val_193 1 +193 val_193 1 +193 val_193 1 +193 val_193 1 +209 val_209 1 +209 val_209 1 +209 val_209 1 +209 val_209 1 +241 val_241 1 +257 val_257 1 +273 val_273 1 +273 val_273 1 +273 val_273 1 +273 val_273 1 +273 val_273 1 +273 val_273 1 +273 val_273 1 +273 val_273 1 +273 val_273 1 +289 val_289 1 +305 val_305 1 +321 val_321 1 +321 val_321 1 +321 val_321 1 +321 val_321 1 +353 val_353 1 +353 val_353 1 +353 val_353 1 +353 val_353 1 +369 val_369 1 +369 val_369 1 +369 val_369 1 +369 val_369 1 +369 val_369 1 +369 val_369 1 +369 val_369 1 +369 val_369 1 +369 val_369 1 +401 val_401 1 +401 val_401 1 +401 val_401 1 +401 val_401 1 +401 val_401 1 +401 val_401 1 +401 val_401 1 +401 val_401 1 +401 val_401 1 +401 val_401 1 +401 val_401 1 +401 val_401 1 +401 val_401 1 +401 val_401 1 +401 val_401 1 +401 val_401 1 +401 val_401 1 +401 val_401 1 +401 val_401 1 +401 val_401 1 +401 val_401 1 +401 val_401 1 +401 val_401 1 +401 val_401 1 +401 val_401 1 +417 val_417 1 +417 val_417 1 +417 val_417 1 +417 val_417 1 +417 val_417 1 +417 val_417 1 +417 val_417 1 +417 val_417 1 +417 val_417 1 +449 val_449 1 +481 val_481 1 +497 val_497 1 Index: ql/src/test/queries/clientpositive/smb_mapjoin_11.q =================================================================== --- ql/src/test/queries/clientpositive/smb_mapjoin_11.q (revision 0) +++ ql/src/test/queries/clientpositive/smb_mapjoin_11.q (revision 0) @@ -0,0 +1,33 @@ +set hive.optimize.bucketmapjoin = true; +set hive.optimize.bucketmapjoin.sortedmerge = true; +set hive.input.format = org.apache.hadoop.hive.ql.io.BucketizedHiveInputFormat; +set hive.enforce.bucketing=true; +set hive.enforce.sorting=true; +set hive.exec.reducers.max = 1; +set hive.merge.mapfiles=false; +set hive.merge.mapredfiles=false; + +-- This test verifies that the output of a sort merge join on 2 partitions (one on each side of the join) is bucketed + +-- Create two bucketed and sorted tables +CREATE TABLE test_table1 (key INT, value STRING) PARTITIONED BY (ds STRING) CLUSTERED BY (key) SORTED BY (key) INTO 16 BUCKETS; +CREATE TABLE test_table2 (key INT, value STRING) PARTITIONED BY (ds STRING) CLUSTERED BY (key) SORTED BY (key) INTO 16 BUCKETS; + +FROM src +INSERT OVERWRITE TABLE test_table1 PARTITION (ds = '1') SELECT * +INSERT OVERWRITE TABLE test_table2 PARTITION (ds = '1') SELECT *; + +set hive.enforce.bucketing=false; +set hive.enforce.sorting=false; + +-- Create a bucketed table +CREATE TABLE test_table3 (key INT, value STRING) PARTITIONED BY (ds STRING) CLUSTERED BY (key) INTO 16 BUCKETS; + +-- Insert data into the bucketed table by joining the two bucketed and sorted tables, bucketing is not enforced +EXPLAIN EXTENDED +INSERT OVERWRITE TABLE test_table3 PARTITION (ds = '1') SELECT /*+ MAPJOIN(b) */ a.key, b.value FROM test_table1 a JOIN test_table2 b ON a.key = b.key AND a.ds = '1' AND b.ds = '1'; + +INSERT OVERWRITE TABLE test_table3 PARTITION (ds = '1') SELECT /*+ MAPJOIN(b) */ a.key, b.value FROM test_table1 a JOIN test_table2 b ON a.key = b.key AND a.ds = '1' AND b.ds = '1'; + +-- Read data from a sampled bucket to verify the data is bucketed +SELECT * FROM test_table3 TABLESAMPLE(BUCKET 2 OUT OF 16) where ds = '1'; Index: ql/src/java/org/apache/hadoop/hive/ql/exec/Utilities.java =================================================================== --- ql/src/java/org/apache/hadoop/hive/ql/exec/Utilities.java (revision 1394297) +++ ql/src/java/org/apache/hadoop/hive/ql/exec/Utilities.java (working copy) @@ -1158,15 +1158,21 @@ * return an integer only - this should match a pure integer as well. {1,3} is used to limit * matching for attempts #'s 0-999. */ - private static Pattern fileNameTaskIdRegex = Pattern.compile("^.*?([0-9]+)(_[0-9]{1,3})?(\\..*)?$"); + private static final Pattern FILE_NAME_TO_TASK_ID_REGEX = Pattern.compile("^.*?([0-9]+)(_[0-9]{1,3})?(\\..*)?$"); /** * This retruns prefix part + taskID for bucket join for partitioned table */ - private static Pattern fileNamePrefixedTaskIdRegex = + private static final Pattern FILE_NAME_PREFIXED_TASK_ID_REGEX = Pattern.compile("^.*?((\\(.*\\))?[0-9]+)(_[0-9]{1,3})?(\\..*)?$"); /** + * This breaks a prefixed bucket number into the prefix and the taskID + */ + private static final Pattern PREFIXED_TASK_ID_REGEX = + Pattern.compile("^(.*?\\(.*\\))?([0-9]+)$"); + + /** * Get the task id from the filename. It is assumed that the filename is derived from the output * of getTaskId * @@ -1174,7 +1180,7 @@ * filename to extract taskid from */ public static String getTaskIdFromFilename(String filename) { - return getIdFromFilename(filename, fileNameTaskIdRegex); + return getIdFromFilename(filename, FILE_NAME_TO_TASK_ID_REGEX); } /** @@ -1185,7 +1191,7 @@ * filename to extract taskid from */ public static String getPrefixedTaskIdFromFilename(String filename) { - return getIdFromFilename(filename, fileNamePrefixedTaskIdRegex); + return getIdFromFilename(filename, FILE_NAME_PREFIXED_TASK_ID_REGEX); } private static String getIdFromFilename(String filename, Pattern pattern) { @@ -1228,14 +1234,41 @@ return replaceTaskId(taskId, String.valueOf(bucketNum)); } + /** + * Returns strBucketNum with enough 0's prefixing the task ID portion of the String to make it + * equal in length to taskId + * + * @param taskId - the taskId used as a template for length + * @param strBucketNum - the bucket number of the output, may or may not be prefixed + * @return + */ private static String replaceTaskId(String taskId, String strBucketNum) { - int bucketNumLen = strBucketNum.length(); + Matcher m = PREFIXED_TASK_ID_REGEX.matcher(strBucketNum); + if (!m.matches()) { + LOG.warn("Unable to determine bucket number from file ID: " + strBucketNum + ". Using " + + "file ID as bucket number."); + return adjustBucketNumLen(strBucketNum, taskId); + } else { + String adjustedBucketNum = adjustBucketNumLen(m.group(2), taskId); + return m.group(1) + adjustedBucketNum; + } + } + + /** + * Adds 0's to the beginning of bucketNum until bucketNum and taskId are the same length. + * + * @param bucketNum - the bucket number, should not be prefixed + * @param taskId - the taskId used as a template for length + * @return + */ + private static String adjustBucketNumLen(String bucketNum, String taskId) { + int bucketNumLen = bucketNum.length(); int taskIdLen = taskId.length(); StringBuffer s = new StringBuffer(); for (int i = 0; i < taskIdLen - bucketNumLen; i++) { s.append("0"); } - return s.toString() + strBucketNum; + return s.toString() + bucketNum; } /**