Exemplos de código para acessar dados do MaxCompute com o SDK Python.
O MaxCompute expõe interfaces da Storage API por meio do SDK Python. Para mais informações, consulte aliyun-odps-python-sdk.
Pré-requisitos
Ao executar o código em ambiente local, verifique se o PyODPS está instalado. Para mais informações, consulte Instale o PyODPS.
O PyODPS também está disponível nos seguintes ambientes:
DataWorks: O PyODPS vem pré-instalado nos nós PyODPS. Desenvolva e execute tarefas periódicas do PyODPS diretamente nesses nós. Para mais informações, consulte Use PyODPS in DataWorks.
PAI: O PyODPS vem pré-instalado em todas as imagens integradas do PAI e funciona imediatamente em componentes como o componente Python personalizado no PAI-Designer. Em PAI Notebooks, siga o fluxo de trabalho padrão do PyODPS. Para mais informações, consulte Overview of basic operations e DataFrame (Not recommended).
O PyODPS é o SDK Python para MaxCompute. Para mais informações sobre o PyODPS, consulte PyODPS.
Exemplos
Para exemplos completos de código, consulte Exemplos do SDK Python.
-
Configure o ambiente para conexão com o service MaxCompute
import os from odps import ODPS from odps.apis.storage_api import * # Ensure that the ALIBABA_CLOUD_ACCESS_KEY_ID environment variable is set to your Access Key ID, # and the ALIBABA_CLOUD_ACCESS_KEY_SECRET environment variable is set to your Access Key Secret. # For security, avoid hardcoding the Access Key ID and Access Key Secret. # The endpoint of the MaxCompute service. Only connections from VPC networks are supported. o = ODPS( os.getenv('ALIBABA_CLOUD_ACCESS_KEY_ID'), os.getenv('ALIBABA_CLOUD_ACCESS_KEY_SECRET'), project='your-default-project', endpoint='your-end-point' ) # The name of the MaxCompute table to access. table = "<table to access>" # The name of the quota to use for accessing MaxCompute. quota_name = "<quota name>" # Connects to the MaxCompute service and creates an Arrow-format Storage API client. def get_arrow_client(): odps_table = o.get_table(table) client = StorageApiArrowClient(odps=o, table=odps_table, quota_name=quota_name) return clientNotaPara obter o nome da cota de um grupo de recursos exclusivo da Storage API (assinatura):
Grupo de recursos exclusivo da Storage API: Faça login no console do MaxCompute. No canto superior esquerdo, selecione a região desejada. No painel de navegação à esquerda, escolha Workspace > Quotas para visualizar as cotas disponíveis. Para mais informações, consulte Gerencie quotas for computing resources.
Storage API: Faça login no console do MaxCompute. No painel de navegação à esquerda, escolha Tenants > Tenant Property para ativar a Storage API.
-
Leitura de dados da tabela
-
Crie uma sessão de leitura para acessar dados do MaxCompute
import logging import sys from odps.apis.storage_api import * from util import * logger = logging.getLogger(__name__) # Creates a read session. The mode parameter specifies the split strategy: 'size' to split by data size, or 'row' to split by row offset. def create_read_session(mode): client = get_arrow_client() req = TableBatchScanRequest(required_partitions=['pt=test_write_1']) if mode == "size": req.split_options = SplitOptions.get_default_options(SplitOptions.SplitMode.SIZE) elif mode == "row": req.split_options = SplitOptions.get_default_options(SplitOptions.SplitMode.ROW_OFFSET) resp = client.create_read_session(req) if resp.status != Status.OK: logger.info("Create read session failed") return logger.info("Read session id: " + resp.session_id) if __name__ == '__main__': logging.basicConfig(format='%(asctime)s - %(pathname)s[line:%(lineno)d] - %(levelname)s: %(message)s', level=logging.INFO) if len(sys.argv) != 2: raise ValueError("Please provide split mode: size|row") mode = sys.argv[1] if mode != "row" and mode != "size": raise ValueError("Please provide split mode: size|row") create_read_session(mode) -
Verifique o status da sessão de leitura
import logging import sys import time from odps.apis.storage_api import * from util import * logger = logging.getLogger(__name__) # Before reading data, ensure the read session is created and ready. def check_session_status(session_id): client = get_arrow_client() req = SessionRequest(session_id=session_id) resp = client.get_read_session(req) if resp.status != Status.OK: logger.info("Get read session failed") return # Session creation can be time-consuming. You must wait for the session status to become NORMAL before reading data. if resp.session_status == SessionStatus.NORMAL: logger.info("Read session id: " + resp.session_id) else: logger.info("Session status is not expected") if __name__ == '__main__': logging.basicConfig(format='%(asctime)s - %(pathname)s[line:%(lineno)d] - %(levelname)s: %(message)s', level=logging.INFO) if len(sys.argv) != 2: raise ValueError("Please provide session id") session_id = sys.argv[1] check_session_status(session_id) -
Leia dados do MaxCompute
# Reads data rows from MaxCompute for a specified session_id and counts the total number of rows. import logging import sys from odps.apis.storage_api import * from util import * logger = logging.getLogger(__name__) def read_rows(session_id): client = get_arrow_client() req = SessionRequest(session_id=session_id) resp = client.get_read_session(req) if resp.status != Status.OK and resp.status != Status.WAIT: logger.info("Get read session failed") return req = ReadRowsRequest(session_id=session_id) if resp.split_count == -1: req.row_index = 0 req.row_count = resp.record_count else: req.split_index = 0 reader = client.read_rows_arrow(req) total_line = 0 while True: record_batch = reader.read() if record_batch is None: break total_line += record_batch.num_rows if reader.get_status() != Status.OK: logger.info("Read rows failed") return logger.info("Total line is:" + str(total_line)) if __name__ == '__main__': logging.basicConfig(format='%(asctime)s - %(pathname)s[line:%(lineno)d] - %(levelname)s: %(message)s', level=logging.INFO) if len(sys.argv) != 2: raise ValueError("Please provide session id") session_id = sys.argv[1] read_rows(session_id)
-
Documentos relacionados
Para mais informações sobre a Storage API, consulte Storage API overview.