Conecte o EMR a um data lake do Object Storage Service (OSS) usando o modo de cache do JindoFS.
Contexto
O EMR conecta-se a um data lake do OSS pelo modo de cache ou pelo modo de armazenamento em blocos do JindoFS.
O modo de cache oferece compatibilidade com o OSS nativo ao armazenar arquivos como objetos. Para aumentar a eficiência de acesso no cluster EMR, o sistema armazena localmente os arquivos acessados com frequência. Essa abordagem preserva o formato original dos arquivos e garante total compatibilidade com outros clientes OSS. Para mais informações, consulte Instruções de uso do modo de cache do JindoFS.
O modo de armazenamento em blocos proporciona máxima eficiência para leitura, gravação de dados e acesso a metadados. Nesse modo, os dados ficam armazenados como blocos no OSS e um cache local acelera as operações. Um serviço de namespace local gerencia os metadados para garantir acesso de alto desempenho. Para mais informações, consulte Instruções de uso do modo de armazenamento em blocos do JindoFS.
Pré-requisitos
-
Crie um cluster EMR. Para mais informações, consulte Criar um cluster.
Ao criar o cluster, observe os seguintes pontos:
O cluster EMR e o bucket do OSS devem pertencer à mesma conta Alibaba Cloud. Para obter melhores resultados, mantenha-os também na mesma região.
Durante a criação do cluster, ative as opções Assign Public Network IP e Log on to Cluster in SSH Mode. Essas configurações conectam o cluster à rede pública e permitem login remoto no servidor via shell.
Os serviços bigboot e smartdata são obrigatórios para as configurações subsequentes. Caso não estejam selecionados por padrão, selecione-os.
Crie uma tarefa de Data Lake Delivery. Para mais informações, consulte Guia de início rápido.
Procedimento
-
Conecte-se ao OSS e ative o cache usando o modo de cache do JindoFS no EMR. Para mais informações, consulte Instruções de uso do modo de cache do JindoFS.
Esse recurso usa discos locais para armazenar em cache os blocos de dados acessados com frequência. Por padrão, ele vem desativado e todas as operações de leitura acessam os dados diretamente no OSS. Quando o cache está ativado, o serviço Jindo gerencia automaticamente o cache local e o limpa com base em uma marca d'água superior. Configure a proporção de cache conforme suas necessidades.
-
Inicie o Spark SQL.
Use uma ferramenta de login remoto, como o PuTTY, para acessar o servidor principal do EMR.
-
Execute o comando abaixo para iniciar o Spark SQL.
spark-sql --master yarn --num-executors 5 --executor-memory 1g --executor-cores 2
-
Use uma instrução SQL para criar uma tabela externa que aponte para o diretório de dados do OSS.
Use a instrução SQL obtida no console do Table Store. A instrução SQL a seguir serve apenas como referência.
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/' ;Na página Deliver Data to OSS da instância, na coluna Actions da tarefa de entrega, clique em View Statement to Create Table para visualizar e copiar a instrução SQL.
-
Execute a seguinte instrução SQL para carregar as partições de dados da fonte de dados do OSS.
No comando,
lineitemé o nome da tabela externa criada.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> -
Visualize os dados.
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>