Os external volumes funcionam como sistemas de arquivos distribuídos no MaxCompute, com suporte do Object Storage Service (OSS). Monte um external volume em um diretório do OSS para que seus jobs do Spark on MaxCompute e MapReduce leiam e gravem arquivos por meio do sistema de permissões do MaxCompute, sem conceder acesso direto ao OSS a cada usuário.
Cada projeto do MaxCompute pode ter vários external volumes.
Casos de uso
Os external volumes são úteis quando você precisa:
Carregar dependências de job na inicialização — baixe automaticamente arquivos JAR, wheels do Python ou arquivos de modelos para o diretório de trabalho do job antes da execução
Ler e gravar arquivos do OSS no código Spark — acessar arquivos armazenados no OSS usando o esquema de caminho
odps://diretamente no código do job SparkAplicar controle de permissão granular — usar o sistema de permissões do MaxCompute para controlar quem pode ler ou gravar em caminhos específicos do volume, em vez de gerenciar políticas de bucket do OSS por usuário
Armazenar saídas de jobs de ML — salve dados de índice ou arquivos de modelo gerados por engines como o Proxima CE de volta no OSS por meio de um volume
Faturamento
Os dados nos external volumes residem no OSS. Não há cobrança de armazenamento no MaxCompute. As taxas de computação se aplicam quando uma engine do MaxCompute lê ou processa dados em um external volume — por exemplo, ao executar um job do Spark on MaxCompute ou MapReduce. As saídas gravadas de volta no OSS (como dados de índice do Proxima CE) estão sujeitas às taxas padrão de armazenamento do OSS.
Pré-requisitos
Antes de começar, verifique se você:
Enviou e recebeu aprovação para uso experimental dos external volumes. Consulte Solicitar uso experimental de novos recursos
Tem o cliente do MaxCompute (odpscmd) V0.43.0 ou posterior instalado. Consulte Cliente do MaxCompute (odpscmd). Se usar o SDK for Java, a versão V0.43.0 ou posterior é obrigatória. Consulte Atualizações de versão
Crie um bucket do OSS. Consulte Criar buckets
Autorizou seu projeto do MaxCompute a acessar o OSS. Consulte Configure um método de acesso ao OSS
Início rápido
Etapa 1: Conceder as permissões necessárias
Para usar external volumes, sua conta precisa das seguintes permissões: CreateInstance, CreateVolume, List, Read e Write. Consulte Permissões do MaxCompute.
-
Verifique se sua conta possui a permissão
CreateVolume:SHOW GRANTS FOR <user_name>; -
Se a permissão
CreateVolumeestiver ausente, conceda-a:GRANT CreateVolume ON project <project_name> TO USER <user_name>;Para revogar a permissão posteriormente:
REVOKE CreateVolume ON project <project_name> FROM USER <user_name>; Execute
SHOW GRANTSnovamente para confirme a concessão da permissão.
Etapa 2: Crie um external volume
Execute o comando a seguir com a conta que possui a permissão CreateVolume:
vfs -create <volume_name>
-storage_provider oss
-url oss://<oss_endpoint>/<bucket>/<path>
-acd <true|false>
-role_arn <arn:aliyun:xxx/aliyunodpsdefaultrole>
Para obter detalhes sobre os parâmetros e outras operações de volume, consulte Operações de external volume.
Após a criação, o caminho do volume será odps://[project_name]/[volume_name]. Use este caminho nos jobs do Spark on MaxCompute e MapReduce.
Etapa 3: Verifique o volume
Liste todos os volumes no projeto atual para confirme a criação do volume:
vfs -ls /;
Usar Spark on MaxCompute com external volumes
O Spark on MaxCompute é compatível com o Spark open source e executa nos recursos de computação integrados, conjuntos de dados e sistema de permissões do MaxCompute.
Há duas maneiras de acessar external volumes a partir de um job Spark:
Referenciar arquivos na inicialização do job — os arquivos do volume são baixados para o diretório de trabalho do job antes do início da execução
Acessar arquivos no código — use o esquema de caminho
odps://diretamente no código Spark para ler e gravar arquivos do volume durante a execução
Referenciar arquivos na inicialização do job
Configure os parâmetros a seguir na seção Parameters do nó ODPS Spark do DataWorks ou no arquivo spark-defaults.conf. Não é possível defina esses parâmetros dentro do código do job.
|
Parâmetro |
Descrição |
|
|
Arquivos para baixe no diretório de trabalho do job antes da inicialização. Separe vários arquivos com vírgulas. Cada valor deve incluir o nome completo do arquivo. |
|
|
Arquivos compactados ( |
Formato do valor:
odps://[project_name]/[volume_name]/[path_to_file]
Exemplo — arquivos:
spark.hadoop.odps.cupid.volume.files=
odps://mc_project/external_volume/data/mllib/kmeans_data.txt,
odps://mc_project/external_volume/target/PythonKMeansExample/KMeansModel/data/part-00000-a2d44ac5-54f6-49fd-b793-f11e6a189f90-c000.snappy.parquet
Após o início do job, o diretório de trabalho conterá kmeans_data.txt e part-00000-a2d44ac5-54f6-49fd-b793-f11e6a189f90-c000.snappy.parquet.
Exemplo — arquivos compactados:
spark.hadoop.odps.cupid.volume.archives=
odps://spark_test_wj2/external_volume/pyspark-3.1.1.zip,
odps://spark_test_wj2/external_volume/python-3.7.9-ucs4.tar.gz
Após o início do job, o diretório de trabalho conterá o conteúdo descompactado de pyspark-3.1.1.zip e python-3.7.9-ucs4.tar.gz.
Acessar arquivos no código
Para ler e gravar arquivos de external volume a partir do código do job Spark, defina os parâmetros a seguir no código:
|
Parâmetro |
Valor |
Descrição |
|
|
|
Habilita o reconhecimento de external volume. Padrão: |
|
|
|
O caminho do volume a ser acessado. Padrão: vazio. |
|
|
|
Classe de implementação para acesso ao OSS. |
|
|
|
Classe de implementação abstrata do sistema de arquivos. |
Exemplo — Clusterização K-means com external volume:
O exemplo a seguir usa o algoritmo K-means. Ele lê dados de treinamento de odps://ms_proj1_dev/volume_yyy1/, treina um modelo e salve a saída de volta no mesmo volume.
Todos os caminhos de arquivo no código usam o esquema odps:// para ler e gravar no external volume.
Defina os quatro parâmetros acima no arquivo spark-defaults.conf ou na seção Parameters do nó ODPS Spark do DataWorks antes de execute este código. O exemplo também requer os parâmetros adicionais a seguir para acesso ao OSS, SDK do JindoFS e runtime do Python:
-- Parameters
spark.hadoop.odps.cupid.volume.paths=odps://ms_proj1_dev/volume_yyy1/
spark.hadoop.odps.volume.common.filesystem=true
spark.hadoop.fs.odps.impl=org.apache.hadoop.fs.aliyun.volume.OdpsVolumeFileSystem
spark.hadoop.fs.AbstractFileSystem.odps.impl=org.apache.hadoop.fs.aliyun.volume.abstractfsimpl.OdpsVolumeFs
spark.hadoop.odps.access.id=xxxxxxxxx
spark.hadoop.odps.access.key=xxxxxxxxx
spark.hadoop.fs.oss.endpoint=oss-cn-beijing-internal.aliyuncs.com
spark.hadoop.odps.cupid.resources=ms_proj1_dev.jindofs-sdk-3.8.0.jar
spark.hadoop.fs.oss.impl=com.aliyun.emr.fs.oss.JindoOssFileSystem
spark.hadoop.odps.cupid.resources=public.python-2.7.13-ucs4.tar.gz
spark.pyspark.python=./public.python-2.7.13-ucs4.tar.gz/python-2.7.13-ucs4/bin/python
spark.hadoop.odps.spark.version=spark-2.4.5-odps0.34.0
-- Code
from numpy import array
from math import sqrt
from pyspark import SparkContext
from pyspark.mllib.clustering import KMeans, KMeansModel
if __name__ == "__main__":
sc = SparkContext(appName="KMeansExample")
# Read training data from the external volume
data = sc.textFile("odps://ms_proj1_dev/volume_yyy1/kmeans_data.txt")
parsedData = data.map(lambda line: array([float(x) for x in line.split(' ')]))
# Train the K-means model
clusters = KMeans.train(parsedData, 2, maxIterations=10, initializationMode="random")
# Evaluate the model
def error(point):
center = clusters.centers[clusters.predict(point)]
return sqrt(sum([x**2 for x in (point - center)]))
WSSSE = parsedData.map(lambda point: error(point)).reduce(lambda x, y: x + y)
print("Within Set Sum of Squared Error = " + str(WSSSE))
# Save the model to the external volume
clusters.save(sc, "odps://ms_proj1_dev/volume_yyy1/target/PythonKMeansExample/KMeansModel")
print(parsedData.map(lambda feature: clusters.predict(feature)).collect())
# Load and use the saved model
sameModel = KMeansModel.load(sc, "odps://ms_proj1_dev/volume_yyy1/target/PythonKMeansExample/KMeansModel")
print(parsedData.map(lambda feature: sameModel.predict(feature)).collect())
sc.stop()
Após a conclusão do job, visualize os arquivos de saída no diretório do OSS mapeado para o volume.
Usar Proxima CE para vetorização no MaxCompute
O Proxima CE realiza indexação vetorial e busca de vizinhos mais próximos em dados armazenados em tabelas do MaxCompute. Os resultados são salvos em um external volume no OSS.
Limitações
O Proxima SDK for Java suporta apenas Linux e macOS. Os arquivos JAR contêm dependências específicas do Linux e não podem ser executados no cliente do MaxCompute em Windows.
O Proxima CE executa dois tipos de tarefas: tarefas locais (sem envolvimento de SQL, MapReduce ou Graph) e tarefas do MaxCompute (executadas via engines SQL, MapReduce ou Graph). Os dois tipos são executados alternadamente. Na inicialização, o Proxima CE tenta carregar o kernel do Proxima na máquina local. Se o kernel for carregado com sucesso, certos módulos serão executados localmente; se o carregamento falhar, erros serão relatados, mas o job continuará usando funções de fallback.
Envie a tarefa usando o cliente do MaxCompute (odpscmd). Os nós MapReduce do DataWorks não são suportados porque a versão subjacente do cliente do MaxCompute está sendo atualizada.
Execute uma tarefa de vetorização do Proxima CE
Etapa 1: Instale o pacote de recursos do Proxima CE.
Etapa 2: Preparar os dados de entrada.
Crie as tabelas de entrada e insira dados de amostra:
-- Create a base table and a query table
CREATE TABLE doc_table_float_smoke(pk STRING, vector STRING) PARTITIONED BY (pt STRING);
CREATE TABLE query_table_float_smoke(pk STRING, vector STRING) PARTITIONED BY (pt STRING);
-- Insert data into the base table
ALTER TABLE doc_table_float_smoke ADD PARTITION(pt='20230116');
INSERT OVERWRITE TABLE doc_table_float_smoke PARTITION (pt='20230116') VALUES
('1.nid','1~1~1~1~1~1~1~1'),
('2.nid','2~2~2~2~2~2~2~2'),
('3.nid','3~3~3~3~3~3~3~3'),
('4.nid','4~4~4~4~4~4~4~4'),
('5.nid','5~5~5~5~5~5~5~5'),
('6.nid','6~6~6~6~6~6~6~6'),
('7.nid','7~7~7~7~7~7~7~7'),
('8.nid','8~8~8~8~8~8~8~8'),
('9.nid','9~9~9~9~9~9~9~9'),
('10.nid','10~10~10~10~10~10~10~10');
-- Insert data into the query table
ALTER TABLE query_table_float_smoke ADD PARTITION(pt='20230116');
INSERT OVERWRITE TABLE query_table_float_smoke PARTITION (pt='20230116') VALUES
('q1.nid','1~1~1~1~2~2~2~2'),
('q2.nid','4~4~4~4~3~3~3~3'),
('q3.nid','9~9~9~9~5~5~5~5');
Etapa 3: Envie a tarefa do Proxima CE.
jar -libjars proxima-ce-aliyun-1.0.0.jar
-classpath proxima-ce-aliyun-1.0.0.jar com.alibaba.proxima2.ce.ProximaCERunner
-doc_table doc_table_float_smoke
-doc_table_partition 20230116
-query_table query_table_float_smoke
-query_table_partition 20230116
-output_table output_table_float_smoke
-output_table_partition 20230116
-data_type float
-dimension 8
-topk 1
-job_mode train:build:seek:recall
-external_volume shanghai_vol_ceshi
-owner_id 1248953xxx
;
Etapa 4: Verifique os resultados.
Consulte a tabela de saída para verifique os resultados de vizinhos mais próximos:
SELECT * FROM output_table_float_smoke WHERE pt='20230116';
Saída esperada:
+------------+------------+------------+------------+
| pk | knn_result | score | pt |
+------------+------------+------------+------------+
| q1.nid | 2.nid | 4.0 | 20230116 |
| q1.nid | 1.nid | 4.0 | 20230116 |
| q1.nid | 3.nid | 20.0 | 20230116 |
| q2.nid | 4.nid | 4.0 | 20230116 |
| q2.nid | 3.nid | 4.0 | 20230116 |
| q2.nid | 2.nid | 20.0 | 20230116 |
| q3.nid | 7.nid | 32.0 | 20230116 |
| q3.nid | 8.nid | 40.0 | 20230116 |
| q3.nid | 6.nid | 40.0 | 20230116 |
+------------+------------+------------+------------+
Próximos passos
Operações de external volume — crie, liste e gerencie external volumes
Acessar o OSS a partir do Spark on MaxCompute — acesso direto ao OSS sem external volumes
Permissões do MaxCompute — gerencie permissões de usuário para volumes e projetos