Showing posts with label Hive. Show all posts
Showing posts with label Hive. Show all posts

Monday, 1 August 2016

Hive Loads Avro Data with Schema Evolution

BACKWARD Compatibility:

If a schema is evolved in a backward compatible way, we can always use the latest schema to query all the data uniformly. For example, removing fields is backward compatible change to a schema, since when we encounter records written with the old schema that contain these fields we can just ignore them. Adding a field with a default value is also backward compatible. 

Let's say we have two version of Employee schema as below.

Schema v1:

{
  "type": "record",
  "name": "Employee",
  "fields": [
      {"name": "email", "type": "string"},
      {"name": "name", "type": "string"},
      {"name": "age", "type": "int"}
  ]

}


Schema v2:


{
  "type": "record",
  "name": "Employee",
  "fields": [
      {"name": "email", "type": "string"},
      {"name": "name", "type": "string"},
      {"name": "yrs", "type": "int", "aliases": ["age"]},
      {"name": "gender", "type": ["null",string"], "default": null}
  ]
}


We will use the latest schema to create a Hive table to load data with different versions of schema.

Please note
1. The "name" fields in two schemas need to be the same. Otherwise, although the data can be loaded in Hive table, but cannot be retrieved successfully.

2. A default value is needed for the optional fields in the latest schema. Specifying "null" as default of a union only works if "null" is specified as first type in the union.

Failed with exception java.io.IOException: org.apache.avro.AvroTypeException:
Found Employee, expecting Employee


CREATE TABLE Avro_table
ROW FORMAT SERDE
'org.apache.hadoop.hive.serde2.avro.AvroSerDe'

WITH SERDEPROPERTIES (
    'avro.schema.url'='file:///root/avro_schema/Employee2.avsc')
STORED as INPUTFORMAT
  'org.apache.hadoop.hive.ql.io.avro.AvroContainerInputFormat'
OUTPUTFORMAT

Hive File Not Found Issue: xasecure-audit.xml

When I ran a hive query in HDP2.5, it shows:

FAILED: RuntimeException java.io.FileNotFoundException: /etc/hive/2.5.0.0-817/0/xasecure-audit.xml (No such file or directory)


The issue is related to the fact that Ranger is used for the authorization. 
The xasecure-audit entails older versions of Ranger.

The problem found at this location was that they did not have correct "ranger admin" password in the RANGER CONFIG tab in Ambari. 

Solution:

When the HDFS was restarted, it tries to create a new ranger repository and it fails due to Incorrect "
rangeradmin" password. 
Once the "rangeradmin" password is updated correctly, the HDFS namenode started with Ranger authorizer and was able to audit all activities via Ranger Console.

Or disable Ranger authorization on Hive:

Find Hive config tab in Ambari, select "None" as Authorization in the Security section,
Then, restart Hive.

Reference:


Wednesday, 2 September 2015

Join Multiple Hive Tables without manually typing column names

Sometimes, we have to join multiple hive tables on the same key. Since there might be many column names, we don't want to type them one by one. Select * doesn't work in this case, due to duplicate join key in the column list.

Below python script can help.
It takes table names, join key, and the columns excluded in the results. 

import subprocess as sub
import sys
from monitor import os_error_exit
import ConfigParser

def getColumns(table_name, columns):
    p = sub.Popen(["hive","-e", "describe "+table_name],stdout=sub.PIPE,stderr=sub.PIPE)
    results, errors = p.communicate()
    schema = results.split('\t')
    for i in range(0, len(schema)-1, 2):
        if "# Partition Information" in schema[i] or schema[i].strip()=='':
            break
        columns.add(schema[i].strip())

def getJoinConditions(table_list, column_key):
    alias = ['t'+str(i) for i in range(100)]
    base_table = table_list[0]
    join_condition = base_table+" " +alias[0]
    for i in range(1, len(table_list)):
        current_table = table_list[i]
        join_condition += " left outer join " + current_table +" " +alias[i]+" on "+alias[0]+'.'+column_key+"="+alias[i]+'.'+column_key
    return join_condition


def main():
    
    columns = set()

    for tablename in tables:
        getColumns(tablename, columns)

    selected_columns = columns.difference(discard_columns)
    column_list =  ', '.join(["COALESCE("+ c +",0.0) as "+c for c in selected_columns])
    full_columns = "t0."+column_key+", "+ column_list
    join_condition_str = getJoinConditions(tables, column_key)

    cmd="hive -hiveconf COLUMNS='" + full_columns + "' -hiveconf JOIN_CONDITION='" + join_condition_str +"' -f " + hive_sql_path

if __name__ == "__main__":
    main()


Thursday, 16 July 2015

Issue of Vectorization on Parquet table

When Vectorization is turned on in Hive:
set hive.vectorized.execution.enabled=true;

If the involved table is in parquet rather than orc format, you may see below error.
This error appears in both "tez" and "mr" engine.

Solution: Disable vectorization.


Caused by: java.io.IOException: java.lang.RuntimeException: org.apache.hadoop.hive.ql.metadata.HiveException: Incompatible Bytes vector column and primitive category VOID
at org.apache.hadoop.hive.io.HiveIOExceptionHandlerChain.handleRecordReaderNextException(HiveIOExceptionHandlerChain.java:121)
at org.apache.hadoop.hive.io.HiveIOExceptionHandlerUtil.handleRecordReaderNextException(HiveIOExceptionHandlerUtil.java:77)
at org.apache.hadoop.hive.ql.io.HiveContextAwareRecordReader.doNext(HiveContextAwareRecordReader.java:352)
at org.apache.hadoop.hive.ql.io.HiveRecordReader.doNext(HiveRecordReader.java:79)
at org.apache.hadoop.hive.ql.io.HiveRecordReader.doNext(HiveRecordReader.java:33)
at org.apache.hadoop.hive.ql.io.HiveContextAwareRecordReader.next(HiveContextAwareRecordReader.java:115)
at org.apache.hadoop.mapred.split.TezGroupedSplitsInputFormat$TezGroupedSplitsRecordReader.next(TezGroupedSplitsInputFormat.java:126)
at org.apache.tez.mapreduce.lib.MRReaderMapred.next(MRReaderMapred.java:113)
at org.apache.hadoop.hive.ql.exec.tez.MapRecordSource.pushRecord(MapRecordSource.java:61)
... 15 more
Caused by: java.lang.RuntimeException: org.apache.hadoop.hive.ql.metadata.HiveException: Incompatible Bytes vector column and primitive category VOID
at org.apache.hadoop.hive.ql.io.parquet.VectorizedParquetInputFormat$VectorizedParquetRecordReader.next(VectorizedParquetInputFormat.java:136)
at org.apache.hadoop.hive.ql.io.parquet.VectorizedParquetInputFormat$VectorizedParquetRecordReader.next(VectorizedParquetInputFormat.java:49)
at org.apache.hadoop.hive.ql.io.HiveContextAwareRecordReader.doNext(HiveContextAwareRecordReader.java:347)
... 21 more

Saturday, 11 July 2015

Table Partitioning vs Bucketing


The main difference between Hive partitioning and Bucketing
when we do partitioning, we create a partition for each unique value of the column. But there may be situation where we need to create lot of tiny partitions. But if you use bucketing, you can limit it to a number which you choose and decompose your data into those buckets. In hive a partition is a directory but a bucket is a file.

Drawbacks of Partitions:
A design that creates too many partitions may optimize some queries, but be detrimental for other important queries. 
Having too many partitions is the large number of Hadoop files and directories that are created unnecessarily and overhead to NameNode since it must keep all metadata for the file system in memory.
Advantages of Bucketing:
The number of buckets is fixed so it does not fluctuate with data. If two tables are bucketed by employee_id, Hive can create a logically correct sampling. Bucketing also aids in doing efficient map-side joins etc.

Bucketing Operation
Bucketed tables are fantastic in that they allow much more efficient sampling than do non-bucketed tables, and they may later allow for time saving operations such as mapside joins. However, the bucketing specified at table creation is not enforced when the table is written to, and so it is possible for the table's metadata to advertise properties which are not upheld by the table's actual layout. This should obviously be avoided. Here's how to do it right.
First, 
CREATE TABLE user_info_bucketed(user_id BIGINT, firstname STRING, lastname STRING)
COMMENT 'A bucketed copy of user_info'
PARTITIONED BY(ds STRING)
CLUSTERED BY(user_id) INTO 256 BUCKETS;
Note that we specify a column (user_id) to base the bucketing.
Then we populate the table
set hive.enforce.bucketing = true; 
FROM user_id
INSERT OVERWRITE TABLE user_info_bucketed
PARTITION (ds='2009-02-25')
SELECT userid, firstname, lastname WHERE ds='2009-02-25';
The command set hive.enforce.bucketing = true; allows the correct number of reducers and the cluster by column to be automatically selected based on the table. Otherwise, you would need to set the number of reducers to be the same as the number of buckets a la set mapred.reduce.tasks = 256; and have a CLUSTER BY ... clause in the select.
How does Hive distribute the rows across the buckets? In general, the bucket number is determined by the expressionhash_function(bucketing_column) mod num_buckets. The hash_function depends on the type of the bucketing column. 

Reference:

Friday, 10 July 2015

SQL selects columns where other column is the Aggregation result

SELECT sum(ordertotal), userid,  username FROM orders GROUP BY userid;

Above query generates the infamous “Expression not in GROUP BY key” error, because the username column is not being aggregated but the ordertotal is.

An easy fix is to aggregate the username values using the collect_set function, but output only one of them: You should get the same output as before, but this time the username is included.

SELECT sum(ordertotal), userid, collect_set(username)[0] FROM orders GROUP BY userid;



OVER allows you to get aggregate information without using a GROUP BY. In other words, you can retrieve detail rows, and get aggregate data alongside it.
Break up that resultset into partitions with the use of PARTITION BY.

SELECT userid, itemlist, sum(ordertotal) OVER (PARTITION BY userid)
FROM orders;

SELECT
  [ID],
  [State],
  [Value]
FROM
(
  SELECT 
    [ID],
    [State],
    [Value],
    ROW_NUMBER() OVER (PARTITION BY [ID] ORDER BY [Value] DESC) AS [Rank]
  FROM [t1]
) AS [sub]
WHERE [sub].[Rank] = 1
ORDER BY
  [ID] ASC,
  [State] ASC

Thursday, 9 July 2015

Ngrams in Hive

ngrams() and context_ngrams(): N-gram frequency estimation

N-grams are subsequences of length N drawn from a longer sequence. The purpose of the ngrams() UDAF is to find the k most frequent n-grams from one or more sequences. It can be used in conjunction with the sentences() UDF to analyze unstructured natural language text, or the collect() function to analyze more general string data.
Contextual n-grams are similar to n-grams, but allow you to specify a 'context' string around which n-grams are to be estimated. For example, you can specify that you're only interested in finding the most common two-word phrases in text that follow the context "I love". You could achieve the same result by manually stripping sentences of non-contextual content and then passing them to ngrams(), but context_ngrams() makes it much easier.

SELECT explode(context_ngrams(sentences(lower(tweet)), 2, 100 [, 1000])) FROM twitter;
The command above will return the top-100 bigrams (2-grams) from a hypothetical table called twitter. The tweetcolumn is assumed to contain a string with arbitrary, possibly meaningless, text. The lower() UDF first converts the text to lowercase for standardization, and then sentences() splits up the text into arrays of words. 
The optional fourth argument is the precision factor that control the tradeoff between memory usage and accuracy in frequency estimation. Higher values will be more accurate, but could potentially crash the JVM with an OutOfMemory error. If omitted, sensible defaults are used.
SELECT explode(context_ngrams(sentences(lower(tweet)), array("i","love",null), 100, [, 1000])) FROM twitter;
The command above will return a list of the top 100 words that follow the phrase "i love" in a hypothetical database of Twitter tweets. Each null specifies the position of an n-gram component to estimate; therefore, every query must contain at least one null in the context array.

Reference: