JindoFS キャッシュモードを使用して、E-MapReduce (EMR) を Object Storage Service (OSS) データレイクに接続します。
背景情報
EMR は、JindoFS キャッシュモードまたは JindoFS ブロックストレージモードのいずれかを使用して、OSS データレイクに接続できます。
-
キャッシュモードは、ファイルをオブジェクトとして保存することで、ネイティブ OSS との互換性を提供します。EMR クラスター内でのアクセス効率を向上させるために、頻繁にアクセスされるファイルはローカルにキャッシュされます。このアプローチでは、元のファイル形式が維持されるため、他の OSS クライアントとの完全な互換性が確保されます。詳細については、「JindoFS キャッシュモードの使用方法」をご参照ください。
-
ブロックストレージモードは、データの読み取り、書き込み、およびメタデータへのアクセスにおいて最大の効率を提供します。このモードでは、データは OSS にブロックとして保存され、ローカルキャッシュが操作を高速化します。ローカルの Namespace サービスがメタデータを管理し、高パフォーマンスのアクセスを保証します。詳細については、「JindoFS ブロックストレージモードの使用方法」をご参照ください。
前提条件
-
EMR クラスターが作成されました。詳細については、「クラスターを作成する」をご参照ください。
クラスターを作成する際は、次の点にご注意ください:
-
EMR クラスターと OSS バケットは、同じ Alibaba Cloud アカウントに属している必要があります。最良の結果を得るためには、同じリージョンに配置することを推奨します。
-
クラスター作成時に、Assign Public Network IP と [SSH モードでクラスターにログイン] を有効にします。これらのオプションにより、クラスターがパブリックネットワークに接続され、シェルを使用してサーバーにリモートでログインできるようになります。
-
後続の設定には、bigboot サービスと smartdata サービスが必要です。デフォルトで選択されていない場合は、必ず選択してください。
-
-
Data Lake Delivery タスクを作成しました。詳細については、「クイックスタート」をご参照ください。
操作手順
-
EMR で JindoFS キャッシュモードを使用して OSS に接続し、キャッシュを有効にします。詳細については、「JindoFS キャッシュモードの使用方法」をご参照ください。
この機能は、ローカルディスクを使用して、頻繁にアクセスされるデータブロックをキャッシュします。デフォルトでは、この機能は無効になっており、すべての読み取り操作は OSS から直接データにアクセスします。キャッシュが有効になっている場合、Jindo サービスはローカルキャッシュを自動的に管理し、高水位標に基づいてクリアします。要件に応じてキャッシュ率を設定してください。
-
Spark SQL を起動します。
-
PuTTY などのリモートログインツールを使用して、EMR ヘッダーサーバーにログインします。
-
次のコマンドを実行して Spark SQL を起動します。
spark-sql --master yarn --num-executors 5 --executor-memory 1g --executor-cores 2
-
-
SQL ステートメントを使用して、OSS データディレクトリを指す外部テーブルを作成します。
Table Store コンソールから取得した SQL ステートメントを使用します。次の SQL ステートメントは参考用です。
CREATE EXTERNAL TABLE lineitem (l_orderkey bigint,l_linenumber bigint,l_receiptdate string,l_returnflag string,l_tax double,l_shipmode string,l_suppkey bigint,l_shipdate string,l_commitdate string,l_partkey bigint,l_quantity double,l_comment string,l_linestatus string,l_extendedprice double,l_discount double,l_shipinstruct string) PARTITIONED BY (`year` int, `month` int) STORED AS PARQUET LOCATION 'jfs://test/' ;インスタンスの OSS へのデータ配信 ページで、配信タスクの [アクション] 列にある ステートメントの表示 をクリックして、SQL ステートメントを表示およびコピーします。
-
次の SQL ステートメントを実行して、OSS データソースからデータパーティションをロードします。
コマンド内の
lineitemは、作成した外部テーブルの名前です。msck repair table lineitem;20/09/22 15:17:04 INFO [main] SparkSQLQueryListener: execution is called 20/09/22 15:17:04 INFO [main] SparkSQLQueryListener: Spark user root executed on 1600759024916 with spark sql successfully. Time taken: 1.377 seconds 20/09/22 15:17:04 INFO [main] SparkSQLCLIDriver: Time taken: 1.377 seconds spark-sql> msck repair table lineitem; 20/09/22 15:17:20 INFO [main] AlterTableRecoverPartitionsCommand: Recover all the partitions in jfs://test/ 20/09/22 15:17:20 INFO [main] AbstractJindoFileSystem: Jboot log name is /var/log/bigboot/jboot-INFO-1600759040539- 20/09/22 15:17:20 INFO [main] OssStore: Filesystem support for magic committers is enabled, write buffer size 1048576 20/09/22 15:17:21 INFO [main] FsStats: cmd=listStatus, src=jfs://test/, dst=null, size=1, parameter=, time-in-ms=444, version=2.7.301 20/09/22 15:17:21 INFO [main] FsStats: cmd=listStatus, src=jfs://test/year=2020, dst=null, size=2, parameter=, time-in-ms=151, version=2.7.301 20/09/22 15:17:21 INFO [main] AlterTableRecoverPartitionsCommand: Found 2 partitions in jfs://test/ 20/09/22 15:17:21 INFO [main] FsStats: cmd=listStatus, src=jfs://test/year=2020/month=8, dst=null, size=21, parameter=, time-in-ms=163, version=2.7 20/09/22 15:17:21 INFO [main] FsStats: cmd=listStatus, src=jfs://test/year=2020/month=9, dst=null, size=21, parameter=, time-in-ms=86, version=2.7 20/09/22 15:17:21 INFO [main] AlterTableRecoverPartitionsCommand: Finished to gather the fast stats for all 2 partitions. 20/09/22 15:17:22 INFO [main] AlterTableRecoverPartitionsCommand: Recovered all partitions (2). 20/09/22 15:17:22 INFO [main] SparkSQLQueryListener: command is called 20/09/22 15:17:22 INFO [main] SparkSQLQueryListener: Spark user root executed on 1600759042070 with spark sql successfully. 20/09/22 15:17:22 INFO [main] SparkSQLQueryListener: execution is called 20/09/22 15:17:22 INFO [main] SparkSQLQueryListener: Spark user root executed on 1600759042100 with spark sql successfully. Time taken: 1.693 seconds 20/09/22 15:17:22 INFO [main] SparkSQLCLIDriver: Time taken: 1.693 seconds spark-sql> -
データをクエリします。
select * from lineitem limit 1;20/09/22 15:18:51 INFO [main] SparkSQLQueryListener: execution is called 20/09/22 15:18:51 INFO [main] SparkSQLQueryListener: Spark user root executed on 1600759131254 with spark sql successfully. 20/09/22 15:18:51 INFO [main] FsStats: cmd=getFileStatus, src=jfs://test/_index, dst=null, size=-1, parameter=null, time-in-ms=22, version=2.7.301 20/09/22 15:18:51 INFO [main] PrunedInMemoryFileIndex: It took 1 ms to list leaf files for 2 paths. 20/09/22 15:18:51 INFO [main] SparkSQLQueryListenerHelper: Partitioned table:default.lineitem;cols:l_orderkey,l_linenumber,l_receiptdate,l_returnflag,l_tax,l_shipmode,l_suppkey,l_shipdate,_commitdate,l_partkey,l_quantity,l_comment,l_linestatus,l_extendedprice,l_discount,l_shipinstruct;parts:year=2020/month=8,year=2020/month=9;paths:jfs://test/year=2020/month=8,jfs://test/year=2020/month=9. 20/09/22 15:18:51 INFO [main] NativeClient: JindoTable put 2 records. 44095908 1 1996-09-19 N 0.03 SHIP 5928453 1996-08-28 1996-06-19 145353442 10.0 lly ironic theo O 14881.8 0.08 TAKE BACK RETURN 020 8 Time taken: 6.22 seconds, Fetched 1 row(s) 20/09/22 15:18:51 INFO [main] SparkSQLCLIDriver: Time taken: 6.22 seconds, Fetched 1 row(s) spark-sql>