Hadoop外部データラッパーの使用

ApacheHiveまたはApacheSparkを介してHadoopForeignDataWrapperを使用できます。HiveとSparkは両方とも、構成されたメタストアにメタデータを格納します。ここで、データベースとテーブルはHiveQLを使用して作成されます。

Hadoop上のApacheHiveでHDFSFDWを使用する

Apache Hive™ data warehouse software facilitates querying and managing large datasets residing in distributed storage. Hive provides a mechanism to project structure onto this data and query the data using a SQL-like language called HiveQL. At the same time, this language allows traditional map/reduce programmers to plug in their custom mappers and reducers when it is inconvenient or inefficient to express this logic in HiveQL.

Hiveには2つのバージョンがあります- HiveServer1 と HiveServer2 は`ApacheHiveウェブサイト<https://hive.apache.org/downloads.html>`_からダウンロードできます。

注釈

HadoopForeignDataWrapperは HiveServer2 のみをサポートします。

Hadoopの上でApacheHiveでHDFSFDWを使用するには:

ステップ1:`weblogs_parse<http://wiki.pentaho.com/download/attachments/23531451/weblogs_parse.zip?version=1&modificationDate=1327096242000/>`_をダウンロードし、`WikiPentahoウェブサイト<https://wiki.pentaho.com/display/BAD/Transforming+Data+within+Hive/>`_。

ステップ2:これらのコマンドを使用して weblog_parse.txt ファイルをアップロードします:

hadoop fs -mkdir /weblogs
hadoop fs -mkdir /weblogs/parse
hadoop fs -put weblogs_parse.txt /weblogs/parse/part-00000

ステップ3:次のコマンドを使用して、 HiveServer を開始します(まだ実行されていない場合)。

$HIVE_HOME/bin/hiveserver2

または

$HIVE_HOME/bin/hive --service hiveserver2

ステップ4:Hiveビーラインクライアントを使用して HiveServer2 に接続します。例:

$ beeline
Beeline version 1.0.1 by Apache Hive
beeline> !connect jdbc:hive2://localhost:10000/default;auth=noSasl

ステップ5:Hiveでテーブルを作成します。

CREATE TABLE weblogs (
    client_ip           STRING,
    full_request_date   STRING,
    day                 STRING,
    month               STRING,
    month_num           INT,
    year                STRING,
    hour                STRING,
    minute              STRING,
    second              STRING,
    timezone            STRING,
    http_verb           STRING,
    uri                 STRING,
    http_status_code    STRING,
    bytes_returned      STRING,
    referrer            STRING,
    user_agent          STRING)
row format delimited
fields terminated by '\t';

ステップ6:ウェブログテーブルにデータをロードします。

hadoop fs -cp /weblogs/parse/part-00000 /user/hive/warehouse/weblogs/

ステップ7:PostgreSQLからデータにアクセスします。これで、PostgreSQLのウェブログテーブルを使用できます。psqlを使用して接続したら、次の手順に従います。

-- set the GUC variables appropriately, e.g. :
hdfs_fdw.jvmpath='/home/edb/Projects/hadoop_fdw/jdk1.8.0_111/jre/lib/amd64/server/'
hdfs_fdw.classpath='/usr/local/edbas/lib/postgresql/HiveJdbcClient-1.0.jar:
                    /home/edb/Projects/hadoop_fdw/hadoop/share/hadoop/common/hadoop-common-2.6.4.jar:
                    /home/edb/Projects/hadoop_fdw/apache-hive-1.0.1-bin/lib/hive-jdbc-1.0.1-standalone.jar'

-- load extension first time after install
CREATE EXTENSION hdfs_fdw;

-- create server object
CREATE SERVER hdfs_server
         FOREIGN DATA WRAPPER hdfs_fdw
         OPTIONS (host '127.0.0.1');

-- create user mapping
CREATE USER MAPPING FOR postgres
    SERVER hdfs_server OPTIONS (username 'hive_username', password 'hive_password');

-- create foreign table
CREATE FOREIGN TABLE weblogs
(
 client_ip                TEXT,
 full_request_date        TEXT,
 day                      TEXT,
 Month                    TEXT,
 month_num                INTEGER,
 year                     TEXT,
 hour                     TEXT,
 minute                   TEXT,
 second                   TEXT,
 timezone                 TEXT,
 http_verb                TEXT,
 uri                      TEXT,
 http_status_code         TEXT,
 bytes_returned           TEXT,
 referrer                 TEXT,
 user_agent               TEXT
)
SERVER hdfs_server
         OPTIONS (dbname 'default', table_name 'weblogs');


-- select from table
postgres=# SELECT DISTINCT client_ip IP, count(*)
           FROM weblogs GROUP BY IP HAVING count(*) > 5000 ORDER BY 1;
       ip        | count
-----------------+-------
 13.53.52.13     |  5494
 14.323.74.653   | 16194
 322.6.648.325   | 13242
 325.87.75.336   |  6500
 325.87.75.36    |  6498
 361.631.17.30   | 64979
 363.652.18.65   | 10561
 683.615.622.618 | 13505
(8 rows)

-- EXPLAIN output showing WHERE clause being pushed down to remote server.
EXPLAIN (VERBOSE, COSTS OFF) SELECT client_ip, full_request_date, uri FROM weblogs WHERE http_status_code = 200;
                                                   QUERY PLAN
----------------------------------------------------------------------------------------------------------------
 Foreign Scan on public.weblogs
   Output: client_ip, full_request_date, uri
   Remote SQL: SELECT client_ip, full_request_date, uri FROM default.weblogs WHERE ((http_status_code = '200'))
(3 rows)

Hadoop上のApacheSparkでHDFSFDWを使用する

Apache Spark™ は、さまざまなユースケースをサポートする汎用分散コンピューティングフレームワークです。リアルタイムストリームとバッチ処理を、スピード、使いやすさ、高度な分析とともに提供します。Sparkは、Hadoop、HBASE、Cassandra、S3などのサードパーティストレージプロバイダーに依存しているため、ストレージレイヤーを提供しません。SparkはHadoopとシームレスに統合し、既存のデータを処理できます。SparkSQLは、 HiveQL と100%互換性があり、 Spark Thrift Server を使用して Hiveserver2 の代替として使用できます。

Hadoop上でApacheSparkでHDFSFDWを使用するには:

手順1:ローカルモードでApacheSparkをダウンロードしてインストールします。

ステップ2:フォルダー $SPARK_HOME/conf に、次の行を含むファイル spark-defaults.conf を作成します。

spark.sql.warehouse.dir hdfs://localhost:9000/user/hive/warehouse

デフォルトでは、sparkはメタデータとデータ自体の両方にダービーを使用します(sparkではウェアハウスと呼ばれます)。Sparkでウェアハウスとしてhadoopを使用するには、このプロパティを追加する必要があります。

ステップ3:SparkThriftServerを起動します。

./start-thriftserver.sh

ステップ4:ログファイルを使用して、Sparkthriftサーバーが実行されていることを確認します。

ステップ5:以下のデータでローカルファイル names.txt を作成します。

$ cat /tmp/names.txt
1,abcd
2,pqrs
3,wxyz
4,a_b_c
5,p_q_r
,

ステップ6:sparkbeelineクライアントを使用してSparkThriftServer2に接続します。例:

$ beeline
Beeline version 1.2.1.spark2 by Apache Hive
beeline> !connect jdbc:hive2://localhost:10000/default;auth=noSasl org.apache.hive.jdbc.HiveDriver

ステップ7:Sparkでサンプルデータを準備します。beelineコマンドラインツールで次のコマンドを実行します。

./beeline
Beeline version 1.2.1.spark2 by Apache Hive
beeline> !connect jdbc:hive2://localhost:10000/default;auth=noSasl org.apache.hive.jdbc.HiveDriver
Connecting to jdbc:hive2://localhost:10000/default;auth=noSasl
Enter password for jdbc:hive2://localhost:10000/default;auth=noSasl:
Connected to: Spark SQL (version 2.1.1)
Driver: Hive JDBC (version 1.2.1.spark2)
Transaction isolation: TRANSACTION_REPEATABLE_READ
0: jdbc:hive2://localhost:10000> create database my_test_db;
+---------+--+
| Result  |
+---------+--+
+---------+--+
No rows selected (0.379 seconds)
0: jdbc:hive2://localhost:10000> use my_test_db;
+---------+--+
| Result  |
+---------+--+
+---------+--+
No rows selected (0.03 seconds)
0: jdbc:hive2://localhost:10000> create table my_names_tab(a int, name string)
                                 row format delimited fields terminated by ' ';
+---------+--+
| Result  |
+---------+--+
+---------+--+
No rows selected (0.11 seconds)
0: jdbc:hive2://localhost:10000>

0: jdbc:hive2://localhost:10000> load data local inpath '/tmp/names.txt'
                                 into table my_names_tab;
+---------+--+
| Result  |
+---------+--+
+---------+--+
No rows selected (0.33 seconds)
0: jdbc:hive2://localhost:10000> select * from my_names_tab;
+-------+---------+--+
|   a   |  name   |
+-------+---------+--+
| 1     | abcd    |
| 2     | pqrs    |
| 3     | wxyz    |
| 4     | a_b_c   |
| 5     | p_q_r   |
| NULL  | NULL    |
+-------+---------+--+

Hadoopの対応するファイルは次のとおりです。

$ hadoop fs -ls /user/hive/warehouse/
Found 1 items
drwxrwxrwx   - org.apache.hive.jdbc.HiveDriver supergroup 0 2020-06-12 17:03 /user/hive/warehouse/my_test_db.db

$ hadoop fs -ls /user/hive/warehouse/my_test_db.db/
Found 1 items
drwxrwxrwx   - org.apache.hive.jdbc.HiveDriver supergroup 0 2020-06-12 17:03 /user/hive/warehouse/my_test_db.db/my_names_tab

ステップ8:PostgreSQLからデータにアクセスします。psqlを使用してPostgresに接続します。

-- set the GUC variables appropriately, e.g. :
hdfs_fdw.jvmpath='/home/edb/Projects/hadoop_fdw/jdk1.8.0_111/jre/lib/amd64/server/'
hdfs_fdw.classpath='/usr/local/edbas/lib/postgresql/HiveJdbcClient-1.0.jar:
                    /home/edb/Projects/hadoop_fdw/hadoop/share/hadoop/common/hadoop-common-2.6.4.jar:
                    /home/edb/Projects/hadoop_fdw/apache-hive-1.0.1-bin/lib/hive-jdbc-1.0.1-standalone.jar'

-- load extension first time after install
CREATE EXTENSION hdfs_fdw;

-- create server object
CREATE SERVER hdfs_server
  FOREIGN DATA WRAPPER hdfs_fdw
  OPTIONS (host '127.0.0.1', port '10000', client_type 'spark', auth_type 'NOSASL');

-- create user mapping
CREATE USER MAPPING FOR postgres
  SERVER hdfs_server OPTIONS (username 'spark_username', password 'spark_password');

-- create foreign table
CREATE FOREIGN TABLE f_names_tab( a int, name varchar(255)) SERVER hdfs_svr
  OPTIONS (dbname 'testdb', table_name 'my_names_tab');

-- select the data from foreign server
select * from f_names_tab;
 a |  name
---+--------
 1 | abcd
 2 | pqrs
 3 | wxyz
 4 | a_b_c
 5 | p_q_r
 0 |
(6 rows)

-- EXPLAIN output showing WHERE clause being pushed down to remote server.
EXPLAIN (verbose, costs off) SELECT name FROM f_names_tab WHERE a > 3;
                                QUERY PLAN
--------------------------------------------------------------------------
 Foreign Scan on public.f_names_tab
   Output: name
   Remote SQL: SELECT name FROM my_test_db.my_names_tab WHERE ((a > '3'))
(3 rows)

注釈

SparkThriftServerはHiveThriftServerと互換性があるため、外部サーバーの作成中に同じポートが使用されていました。Hiveserver2を使用するアプリケーションは、 ANALYZE コマンドの動作と NOSASL の場合の接続文字列を除き、Sparkで動作します。HiveをSparkに置き換える場合は、 ALTER SERVER を使用して client_type オプションを変更することをお勧めします。