Monte e use o Alibaba Cloud OSS como armazenamento distribuído no MaxFrame com o decorador with_fs_mount. O FS Mount oferece acesso estável a dados externos no nível do sistema de arquivos, ideal para processamento de dados em grande escala.
Casos de uso
O FS Mount é ideal para análises de big data em jobs do MaxFrame que interagem com armazenamento de objetos persistente, como o OSS. Exemplos de uso:
Carregar, limpar e processar dados brutos diretamente do OSS.
Gravar resultados intermediários no OSS para consumo por tarefas subsequentes.
Compartilhar recursos estáticos, como arquivos de modelos treinados e de configuração.
Métodos tradicionais de leitura e escrita, como pd.read_csv("oss://..."), têm limitações de desempenho do SDK e sobrecarga de rede em ambientes distribuídos. O FS Mount permite acessar arquivos do OSS como se estivessem em um disco local, o que melhora significativamente a eficiência do desenvolvimento.
Procedimento
Ativar serviços e conceder permissões
-
Activate OSS e crie um bucket.
Faça login no OSS console.
No painel de navegação à esquerda, clique em Buckets.
-
Na página Buckets, clique em Create Bucket.
Neste exemplo, o nome do bucket é
xxx-oss-test-sh.
-
Crie uma função do RAM para o MaxCompute e conceda acesso ao ambiente de execução.
Faça login no RAM console.
No painel de navegação à esquerda, selecione .
Na página Roles, clique em Create Role.
-
No canto superior direito da página Create Role, clique em Create Service Linked Role.
Na página Create Role, defina Principal Type como Cloud Service.
Para Principal Type, selecione MaxCompute.
-
Na aba Manage Permissions, clique em Create Authorization. No painel Create Authorization exibido, selecione as políticas a serem concedidas à função e clique em OK.
Selecione as seguintes políticas:
AliyunOSSFullAccess: Concede permissões para gerenciar o OSS.
AliyunMaxComputeFullAccess: Concede permissões para gerenciar o MaxCompute.
Usar with_fs_mount para montar o OSS
-
Recomendado: Autenticar com um ARN de função
from maxframe.udf import with_fs_mount @with_fs_mount( "oss://oss-cn-xxxx-internal.aliyuncs.com/xxx-oss-test-sh/test/", "/mnt/oss_data", storage_options={ "role_arn": "acs:ram::xxx:role/maxframe-oss" }, ) def _process(batch_df): import os if os.path.exists('/mnt/oss_data'): print(f"Mounted files: {os.listdir('/mnt/oss_data')}") else: print("/mnt/oss_data not mounted!") return batch_df * 2 -
Não recomendado: Codificar credenciais diretamente
Este método destina-se apenas a testes e não é recomendado para ambientes de produção.
storage_options={ "access_key_id": "LTAI5t...", "access_key_secret": "Wp9H..." }ImportanteEvite codificar seu AccessKey diretamente. Use
role_arnpara permitir que o sistema solicite automaticamente um token STS temporário, impedindo a exposição do seu par de AccessKey.
Usar with_running_options para controlar a alocação de recursos
Use o decorador with_running_options para alocar recursos de CPU e memória para sua tarefa:
from maxframe.udf import with_running_options
@with_running_options(engine="dpe", cpu=2, memory=16)
@with_fs_mount(...)
def _process(batch_df):
...
|
Parameter |
Recommended value |
Description |
|
|
Fixo |
O FS Mount atualmente suporta apenas o mecanismo DPE. |
|
|
1–4 |
Aumente este valor para tarefas intensivas em I/O ou com alta demanda de descompactação. |
|
|
Comece com 8 GB |
Para carregamento de arquivos grandes, recomenda-se 16 GB ou mais. |
Exemplo
Padrão recomendado: Processe dados em lotes.
Em cenários de processamento de dados em grande escala, use o recurso apply_chunk do MaxFrame para processar os dados de entrada em lotes.
Criar uma sessão do MaxFrame
import os
from odps import ODPS
from maxframe import new_session
from maxframe.udf import with_fs_mount, with_running_options
# Initialize the ODPS client.
# We recommend setting the ALIBABA_CLOUD_ACCESS_KEY_ID and ALIBABA_CLOUD_ACCESS_KEY_SECRET
# environment variables instead of hard-coding the AccessKey ID and AccessKey Secret strings.
o = ODPS(
os.getenv('ALIBABA_CLOUD_ACCESS_KEY_ID'),
os.getenv('ALIBABA_CLOUD_ACCESS_KEY_SECRET'),
project='<your-project>',
endpoint='https://service.cn-<region>.maxcompute.aliyun.com/api',
)
# Set the runtime image.
# The `maxframe_service_dpe_runtime` image includes the required ossfs2 dependency.
# If you use a custom image, you must download the dependency and include it in your image.
# You can find the package at the link below this code block.
options.sql.settings = { "odps.session.image": "maxframe_service_dpe_runtime"}
# Start the session.
session = new_session(o)
print("LogView:", session.get_logview_address())
print("Session ID:", session.session_id)
@with_running_options(engine="dpe", cpu=2, memory=8)
@with_fs_mount(
"oss://oss-cn-<region>-internal.aliyuncs.com/wzy-oss-test-sh/test/",
"/mnt/oss_data",
storage_options={
"role_arn": "acs:ram::<uid>:role/maxframe-oss"
},
)
Pacote de dependência do OSSFS: ossfs2_2.0.3.1_linux_x86_64.deb
Criar uma função definida pelo usuário (UDF)
def _process(batch_df):
import pandas as pd
import os
# Step 1: Check if the mount was successful.
mount_point = "/mnt/oss_data"
if not os.path.exists(mount_point):
raise RuntimeError("OSS mount failed!")
# Step 2: Load data, such as a mapping table or dictionary.
mapping_file = os.path.join(mount_point, "category_map.csv")
if os.path.isfile(mapping_file):
mapping_df = pd.read_csv(mapping_file)
# Step 3: Process the current chunk.
result = batch_df.copy()
result['F'] = result['A'] * 10
return result
Construir um DataFrame e aplicar a UDF
import maxframe.dataframe as md
data = [[1.0, 2.0, 3.0, 4.0, 5.0], ...]
df = md.DataFrame(data, columns=['A', 'B', 'C', 'D', 'E'])
# Apply the UDF by using apply_chunk.
result_df = df.mf.apply_chunk(
_process,
skip_infer=True,
output_type="dataframe",
dtypes=df.dtypes,
index=df.index
)
# Execute the operation and fetch the result.
result = result_df.execute().fetch()
Definir skip_infer=True ignora a inferência de tipos e melhora a velocidade de execução. No entanto, garanta que dtypes e index sejam passados corretamente.
Solução de problemas
Verificar o status da montagem
Adicione logs de depuração na função _process:
import os
print("Mount path exists:", os.path.exists("/mnt/oss_data"))
print("Files in mount:", os.listdir("/mnt/oss_data") if os.path.exists("/mnt/oss_data") else [])
Verifique a saída do LogView para confirmar se foram gerados logs semelhantes aos seguintes:
FS Mount successful! /mnt/oss_data: ['data.csv', 'config.json', 'model.pkl']
Processing batch with shape: (1000, 5)