Friday, January 11, 2013

Bash single quote escape in running a Hive query

If you want to run this hive query using SSH, you probably will be headache how to escape:
select * from my_table where date_column = '2013-01-11';
The correct command looks like this:
ssh myserver 'hive -e "select * from my_table where date_column = '"'2013-01-11'"';"'
What hack? Here is how bash read this:
1. ssh
2. myserver
3. hive -e "select * from my_table where date_column = 
4. '2013-01-01'
5. ;"
3,4,5 will be concatenated to one line
hive -e "select * from my_table where date_column = '2013-01-01';"

Thursday, December 27, 2012

Have to read all from Ruby PTY output

I wrote a Ruby script to call ssh-copy-id to distribute my public key to a list of hosts. It took me a while to make it work. The code is simple: read the password, and call ssh-copy-id for each host to copy the public key. Of course, the code doesn't consider all scenarios like the key is already copied, and wrong password. The tricky parts are "cp_out.readlines" and "Process.wait(pid)". If you don't read all the data from cp_out (comment out cp_out.readlines), the spawned process won't return.
#!/usr/bin/env ruby
require 'rubygems'

require 'pty'
require 'expect'
require 'io/console'

hosts_file = ARGV[0] || "hosts"

print "Password:"
password = $stdin.noecho(&:gets)
password.chomp!
puts

$expect_verbose = true
File.open(hosts_file).each do |host|
  host.chomp!
  print "Copying id to #{host} ... "
  begin
    PTY.spawn("ssh-copy-id #{host}") do |cp_out, cp_in, pid|
      begin
        pattern = /#{host}'s password:/
        cp_out.expect(pattern, 10) do |m|
          cp_in.printf("#{password}\n")
        end
        cp_out.readlines
      rescue Errno::EIO
      ensure
        Process.wait(pid)
      end
    end
  rescue PTY::ChildExited => e
    puts "Exited: #{e.status}"
  end
  status = $?
    if status == 0 
      puts "Done!"
  else
    puts "Failed with exit code #{status}!"
  end
end

Tuesday, December 11, 2012

Puppet require vs. include vs. class

According to puppet reference include and require:

  • Both include and require are functions;
  • Both include and require will "Evaluate one or more classes";
  • Both include and require cannot handle parametered classes
  • require is a superset of include, because it "adds the required class as a dependency";
  • require could cause "nasty dependency cycle";
  • require is "largely unnecessary"; See puppet language guide.
    Puppet also has a require function, which can be used inside class definitions and which does implicitly declare a class, in the same way that the include function does. This function doesn’t play well with parameterized classes. The require function is largely unnecessary, as class-level dependencies can be managed in other ways.
  • We can include a class multiple times, but cannot declare a class multiple times.
    class inner {
      notice("I'm inner")
    
      file {"/tmp/abc":
        ensure => directory
      }
    }
    
    class outer_a {
      # include inner
      class { "inner": }
    
      notice("I'm outer_a")
    }
    
    class outer_b {
      # include inner
      class { "inner": }
    
      notice("I'm outer_b")
    }
    
    include outer_a
    include outer_b
    
    Duplicate declaration: Class[Inner] is already declared in file /home/bewang/temp/puppet/require.pp at line 11; cannot redeclare at /home/bewang/temp/puppet/require.pp:18 on node pmaster.puppet-test.com
    
  • You can safely include a class, first two examples pass, but you cannot declare class inner after outer_a or outer_b like the third one:
    class inner {
    }
    
    class outer_a {
      include inner
    }
    
    class outer_b {
      include inner
    }
    
    class { "inner": }
    include outer_a
    include outer_b
    
    class { "inner": }
    class { "outer_a": }
    class { "outer_b": }
    
    include outer_a
    include outer_b
    class { "inner": } # Duplicate redeclaration error
    

Thursday, December 6, 2012

Hive Metastore Configuration

Recently I wrote a post for Bad performance of Hive meta store for tables with large number of partitions. I did tests in our environment. Here is what I found:

  • Don't configure a hive client to access remote MySQL database directly as follows. The performance is really bad, especially when you query a table with a large number of partitions.
  •    
        javax.jdo.option.ConnectionURL   
        jdbc:mysql://mysql_server/hive_meta
       
       
        javax.jdo.option.ConnectionDriverName   
        com.mysql.jdbc.Driver
    
    
        javax.jdo.option.ConnectionUserName
        hive_user
    
    
        javax.jdo.option.ConnectionPassword
        password
    
    
  • Must start Hive metastore service on the same server where Hive MySQL database is running.
    • On database server, use the same configuration as above
    • Start the hive metasore service
    • hive --service metastore
      # If use CDH
      yum install hive-metastore
      /sbin/service hive-metastore start
      
    • On hive client machine, use the following configuration.
    •    
          hive.metastore.uris   
          thrift://mysql_server:9083
         
      
    • Don't worry if you see this error message.
    • ERROR conf.HiveConf: Found both hive.metastore.uris and javax.jdo.option.ConnectionURL Recommended to have exactly one of those config keyin configuration
The reason is: when Hive does partition pruning, it will read a list of partitions. The current metastore implementation uses JDO to query the metastore database:
  1. Get a list of partition names using db.getPartitionNames()
  2. Then call db.getPartitionsByName(List<Strin> partNames). If the list is too large, it will load in multiple times, 300 for each load by default. The JDO calls like this
    • For one MPartition object.
    • Send 1 query to retrieve MPartition basic fields.
    • Send 1 query to retrieve MStorageDescriptors
    • Send 1 query to retrieve data from PART_PARAMS.
    • Send 1 query to retrieve data from PARTITION_KEY_VALS.
    • ...
    • Totally 10 queries for one MPartition. Because MPartition will be converted into Partition before send by, all fields will be populated
  3. If one query takes 40ms in my environment. And you can calculate how long does it take for thousands partitions.
  4. Using remote Hive metastore service, all those queries happens locally, it won't take that long for each query, so you can get performance improved significantly. But there are still a lot of queries.

I also wrote ObjectStore using EclipseLink JPA with @BatchFetch. Here is the test result, it will at least 6 times faster than remote metastore service. It will be even faster.

PartitionsJDO Remote MySQL
Remote Service
EclipseLink
Remote MySQL
106,142353569
10057,0763,914940
200116,2165,2541,211
500287,41621,3853,711
1000574,60639,8466,652
3000
132,64519,518

Friday, November 30, 2012

Bad performance of Hive meta store for tables with large number of partitions

Just found this article Batch fetching - optimizing object graph loading

We have some tables with 15K ~ 20K partitions. If I run a query scanning a lot of partitions, Hive could use more than 10 minutes to commit the mapred job.

The problem is caused by ObjectStore.getPartitionsByNames when Hive semantic analyzer tries to prune partitions. This method sends a lot of queries to our MySQL database to retrieve ALL information about partitions. Because MPartition and MStroageDescriptor are converted into Partition and StorageDescriptor, every field will be accessed during conversion, in other words, even the fields has nothing to do with partition pruning, such as BucketCols. In our case, 10 queries for each partition will be sent to the database and each query may take 40ms.

This is known ORM 1+N problem. But it is really bad user experience.

Actually we assembly Partition objects manually, it would only need about 10 queries for a group of partitions (default size is 300). In our environment, it only needs 40 seconds for 30K partitions: 30K / 300 * 10 * 40.

I tried to this way:

  1. Fetch MPartition with fetch group and fetch_size_greedy, so one query can get MPartition's primary fields and MStorageDescriptor cached.
  2. Get all descriptors into a list "msds", run another query to get MStorageDescriptor with filter like "msds.contains(this)", all cached descriptors will be refreshed in one query instead of n queries.

This works well for 1-1 relations, but not on 1-N relation like MPartition.values. I didn't find a way to populate those fields in just one query.

Because JDO mapping doesn't work well in the conversion (MPartition - Partition), I'm wondering if it is worth doing like this:

  1. Query each table in SQL directly PARTITIONS, SDS, etcs.
  2. Assembly Partition objects

This is a hack and the code will be really bad. But I didn't find JDO support "FETCH JOIN" or "Batch fetch".

Tuesday, November 27, 2012

CDH4.1.0 Hive Debug Option Bug

Hive 0.9.0 supports remote debug.Run the following command, and Hive will suspend and listening on 8000.

hive --debug
But there is a bug in Hive CDH4.1.0 which blocks you from using this option. You will get this error message:
[bewang@myserver ~]$ hive --debug
ERROR: Cannot load this JVM TI agent twice, check your java command line for duplicate jdwp options.
Error occurred during initialization of VM
agent library failed to init: jdwp

By setting xtrace, I found "-XX:+UseParallelGC -agentlib:jdwp=transport=dt_socket,server=y,address=8000,suspend=y" is actually add twice. In this commit 54abc308164314a6fae0ef0b2f2241a6d4d9f058, HADOOP_CLIENT_OPTS is appended to HADOOP_OPTS, unfortunately this is done in $HADOOP_HOME/bin/hadoop".

--- a/bin/hive
+++ b/bin/hive
@@ -216,6 +216,7 @@ if [ "$DEBUG" ]; then
   else
     get_debug_params "$DEBUG"
     export HADOOP_CLIENT_OPTS="$HADOOP_CLIENT_OPTS $HIVE_MAIN_CLIENT_DEBUG_OPTS"
+    export HADOOP_OPTS="$HADOOP_OPTS $HADOOP_CLIENT_OPTS"
   fi
 fi

Removing this line will fix the issue.

Wednesday, November 21, 2012

Using Puppet to Manage Java Application Deployment

Deployment of a Java application can be done in a different ways. We can do

  • Build one jar with all exploded dependent jars and your application. You don't have to worry about your classpath.
  • Build a tarball. It is much easier using Maven assembly and Maven AppAssembly plugin to build a tarball/zip to include all dependent jars and the wrapper scripts for running application or service. AppAssembly generate the classpath in the wrapper scripts.
  • Build a war and deploy it to the container.

All above methods have the same issue: you will repeatedly deploy the same jars again and again. One of my project uses a lot of third-party libraries, Spring, Drools, Quartz, EclipseLink and Apache CXF. The whole zip file is about 67MB with all transitive dependent jars, but my application is actually less than 1MB. Every time I have to build and store such large files, and include all dependent jars without any change. Even the hard drive is much cheaper today, doing like this is still not a good practice.

It is easy to deploy my application jar and replace the old jar only each time if I don't update any dependency. Otherwise, it is still inconvenient because you have to change multiple jars.

The ideal way is to deploy the application jar and dependent jars in a directory other than the application directory on the target server, and build the classpath each time when a new version or dependency change is deployed. There are the following advantages:

  • Dependent jars only needs to be deployed once.
  • Different applications can share the same dependencies.
  • Multiple versions exist on the target server. It is much easier to rollback to the old version.
  • You save a lot of space and network traffic

You probably get the answer: Maven. Using maven, maven dependency plugin, and maven local repository, you can simply and easily implement such a system.

I have a working puppet module can do the followings:

  • setup maven from tarball;
  • setup maven settings.xml, local repository, and you maven repository on the intranet;
  • resolve the dependencies of a maven artifact: download jars into the local repository;
  • setup symlinks of dependent jars into a directory so that you don't have different copies for different applications;
  • generate a file having the classpath.

I will publish it to github once I get the time.

You can find similar puppet modules in github, like puppet-nexus and puppet-maven. They just copy the specified jars to a directory.