O Spark on MaxCompute permite consultar dados em ambientes Hadoop (Hive/E-MapReduce) e Data Lake Formation (DLF) existentes com fontes de dados do Object Storage Service (OSS), sem migrar seus dados ou reescrever seus jobs Spark.
Quando usar esta abordagem
Use o Spark on MaxCompute para acessar fontes de dados externas nas seguintes situações:
Seus jobs Spark já executam em um metastore Hive no E-MapReduce (EMR) ou em um catálogo DLF com OSS, e você deseja adotar o MaxCompute como mecanismo de computação sem mover os dados.
Você precisa consultar tabelas particionadas ou não particionadas existentes durante um período de transição de migração.
Caso planeje migrar totalmente seus dados para o MaxCompute, use o MaxCompute SQL ou as ferramentas padrão de importação de dados do MaxCompute.
Pré-requisitos
Antes de começar, verifique se você tem:
Um projeto externo no MaxCompute mapeado para seu banco de dados Hive (para fontes Hadoop) ou para seu banco de dados DLF (para fontes DLF + OSS).
O runtime do Spark on MaxCompute configurado e acessível.
A versão Spark
spark-2.4.5-odps0.34.0disponível em seu ambiente.
Conceitos principais
Projeto externo: Projeto do MaxCompute mapeado para um catálogo de metadados externo (Hive ou DLF), em vez de armazenar dados no próprio MaxCompute. O Spark on MaxCompute consulta os dados por meio desse mapeamento em tempo de execução.
Acesso a tabela externa: Acesso a tabelas definidas em um projeto externo. Tanto o acesso a tabelas externas quanto o acesso a projetos externos ficam desativados por padrão e devem ser habilitados explicitamente.
Parâmetros de configuração
Defina estes parâmetros na configuração do Spark antes de executar seu job.
Parâmetros obrigatórios (todos os cenários)
|
Parâmetro |
Padrão |
Descrição |
|
|
|
Habilita o acesso a tabelas externas. Defina como |
|
|
|
Habilita o acesso a projetos externos. Defina como |
|
|
— |
Especifica a versão do Spark. Defina como |
Parâmetros adicionais para cenários DLF + OSS
|
Parâmetro |
Padrão |
Descrição |
|
|
— |
Defina como |
|
|
— |
Endpoint do OSS para a região onde seus arquivos de dados estão armazenados (por exemplo, |
|
|
— |
ID da região (por exemplo, |
|
|
— |
Região padrão do OSS (por exemplo, |
Parâmetros do SparkSession
|
Parâmetro |
Valor de exemplo |
Descrição |
|
|
|
Tempo limite de broadcast join em segundos. O padrão é 300 segundos (5 minutos). O exemplo define 1.200 segundos. |
|
|
|
Modo de partição. No modo |
|
|
|
Endpoint do OSS definido diretamente na configuração do SparkSession. |
Acesse um projeto externo baseado em uma fonte de dados Hadoop
O projeto externo hadoop_external_project mapeia para um banco de dados Hive no E-MapReduce. Os exemplos a seguir leem dados de uma tabela não particionada (testtbl) e de uma tabela particionada (testtbl_par).
Consulta com MaxCompute SQL
-- Read from a non-partitioned table
SELECT * FROM hadoop_external_project.testtbl;
-- Read from a partitioned table
SELECT * FROM hadoop_external_project.testtbl_par WHERE b='20220914';
Consulta com Spark on MaxCompute
Adicione os seguintes parâmetros à sua configuração do Spark:
spark.sql.odps.enableExternalTable=true
spark.sql.odps.enableExternalProject=true
spark.hadoop.odps.spark.version=spark-2.4.5-odps0.34.0
Em seguida, execute o seguinte código Scala:
import org.apache.spark.sql.SparkSession
object external_Project_ReadTableHadoop {
def main(args: Array[String]): Unit = {
val spark = SparkSession
.builder()
.appName("external_TableL-on-MaxCompute")
// Broadcast join timeout. Default is 300 seconds.
.config("spark.sql.broadcastTimeout", 20 * 60)
// nonstrict: all partitions can be dynamic. strict: at least one static partition required.
.config("odps.exec.dynamic.partition.mode", "nonstrict")
.config("oss.endpoint", "oss-cn-shanghai-internal.aliyuncs.com")
.getOrCreate()
// List tables in the external project
print("=====show tables in hadoop_external_project6=====")
spark.sql("show tables in hadoop_external_project6").show()
// Read from a non-partitioned table
print("===============hadoop_external_project6.testtbl;================")
spark.sql("desc extended hadoop_external_project6.testtbl").show()
print("===============hadoop_external_project6.testtbl;================")
spark.sql("SELECT * from hadoop_external_project6.testtbl").show()
// Read from a partitioned table
print("===============hadoop_external_project6.testtbl_par;================")
spark.sql("desc extended hadoop_external_project6.testtbl_par").show()
print("===============hadoop_external_project6.testtbl;================")
spark.sql("SELECT * from hadoop_external_project6.testtbl_par where b='20220914'").show()
}
}
Acesse um projeto externo baseado em DLF e OSS
O projeto externo ext_dlf_0713 mapeia para um banco de dados do Data Lake Formation (DLF) com suporte de OSS. O exemplo a seguir lê dados de uma tabela não particionada (tbl_oss1).
Consulta com MaxCompute SQL
-- Read from a non-partitioned table in the DLF-backed external project
SELECT * FROM ext_dlf_0713.tbl_oss1;
Consulta com Spark on MaxCompute
Adicione os seguintes parâmetros à sua configuração do Spark. Os parâmetros de endpoint e região do OSS são obrigatórios para cenários DLF + OSS.
spark.sql.odps.enableExternalTable=true
spark.sql.odps.enableExternalProject=true
spark.hadoop.odps.spark.version=spark-2.4.5-odps0.34.0
# Add this parameter only if the OSS directory was created by EMR
spark.hadoop.odps.oss.location.uri.style=emr
spark.hadoop.odps.oss.endpoint=oss-cn-shanghai-internal.aliyuncs.com
spark.hadoop.odps.region.id=cn-shanghai
spark.hadoop.odps.oss.region.default=cn-shanghai
Em seguida, execute o seguinte código Scala:
import org.apache.spark.sql.{SaveMode, SparkSession}
object external_Project_ReadTable {
def main(args: Array[String]): Unit = {
val spark = SparkSession
.builder()
.appName("external_TableL-on-MaxCompute")
// Broadcast join timeout. Default is 300 seconds.
.config("spark.sql.broadcastTimeout", 20 * 60)
// nonstrict: all partitions can be dynamic. strict: at least one static partition required.
.config("odps.exec.dynamic.partition.mode", "nonstrict")
.config("oss.endpoint", "oss-cn-shanghai-internal.aliyuncs.com")
.getOrCreate()
// List tables in the DLF-backed external project
print("=====show tables in ext_dlf_0713=====")
spark.sql("show tables in ext_dlf_0713").show()
// Read from a non-partitioned table
print("===============ext_dlf_0713.tbl_oss1;================")
spark.sql("desc extended ext_dlf_0713.tbl_oss1").show()
print("===============ext_dlf_0713.tbl_oss1;================")
spark.sql("SELECT * from ext_dlf_0713.tbl_oss1").show()
}
}