Todos os produtos
Search
Central de documentação

Platform For AI:Script PyAlink

Última atualização: Jun 27, 2026

O componente PyAlink Script permite chamar qualquer algoritmo Alink por meio de código. Use este componente para executar tarefas como classificação, regressão e recomendação. O PyAlink Script também se integra perfeitamente a outros componentes do Machine Learning Designer, possibilitando a criação e validação de pipelines de negócios de ponta a ponta. Este tópico descreve como usar o componente PyAlink Script.

Contexto

Use o componente PyAlink Script de duas formas: isoladamente ou combinado com outros componentes do Machine Learning Designer. Ele fornece acesso a centenas de componentes Alink e suporta leitura e gravação de tipos de dados via código. Também é possível implantar um PipelineModel gerado pelo componente como um serviço EAS. Para mais informações, consulte Exemplo: Implantar um modelo gerado por um PyAlink Script como serviço EAS.

Conceitos principais

Antes de usar o componente PyAlink Script, familiarize-se com os seguintes conceitos fundamentais.

Conceito

Descrição

operator

No Alink, um operador representa uma função algorítmica. Os operadores dividem-se em batch ou stream. A regressão logística, por exemplo, inclui os seguintes operadores:

  • LogisticRegressionTrainBatchOp: Executa o treinamento de regressão logística.

  • LogisticRegressionPredictBatchOp: Realiza previsão em lote para regressão logística.

  • LogisticRegressionPredictStreamOp: Realiza previsão em fluxo para regressão logística.

Conecte os operadores usando os métodos link ou linkFrom. O código a seguir apresenta um exemplo.

# Define data.
    data = CsvSourceBatchOp()
    # Logistic regression training.
    lrTrain = LogisticRegressionTrainBatchOp()
    # Logistic regression prediction.
    LrPredict = LogisticRegressionPredictBatchOp()
    # Train.
    data.link(lrTrain)
    # Predict.
    LrPredict.linkFrom(lrTrain, data)

Cada operador possui parâmetros. A regressão logística, por exemplo, inclui os seguintes parâmetros.

  • labelCol: Nome da coluna de rótulo na tabela de entrada. Parâmetro obrigatório do tipo String.

  • featureCols: Array de nomes de colunas de recursos. Parâmetro do tipo String[]. O valor padrão é NULL, indicando que todas as colunas estão selecionadas.

Para configurar um parâmetro, use o prefixo set seguido pelo nome do parâmetro em CamelCase. O código abaixo serve como exemplo.

lr = LogisticRegressionTrainBatchOp()\
                .setFeatureCols(colnames)\
                .setLabelCol("label")

Fontes e destinos de dados são tipos especiais de operadores. Após defini-los, conecte-os aos componentes de algoritmo usando os métodos link ou linkFrom.

image

O Alink inclui fontes de dados comuns para stream e batch. O código a seguir mostra um exemplo.

df_data = pd.DataFrame([
    [2, 1, 1],
    [3, 2, 1],
    [4, 3, 2],
    [2, 4, 1],
    [2, 2, 1],
    [4, 3, 2],
    [1, 2, 1],
    [5, 3, 2]
])
input = BatchOperator.fromDataframe(df_data, schemaStr='f0 int, f1 int, label int')
# load data
dataTest = input
colnames = ["f0","f1"]
lr = LogisticRegressionTrainBatchOp().setFeatureCols(colnames).setLabelCol("label")
model = input.link(lr)
predictor = LogisticRegressionPredictBatchOp().setPredictionCol("pred")
predictor.linkFrom(model, dataTest).print()

pipeline

Também é possível usar algoritmos Alink em um pipeline, combinando processamento de dados, geração de recursos e treinamento de modelo em um único fluxo para treinamento, previsão e serviços online. Veja um exemplo no código a seguir.

quantileDiscretizer = QuantileDiscretizer()\
            .setNumBuckets(2)\
            .setSelectedCols("sepal_length")

binarizer = Binarizer()\
            .setSelectedCol("petal_width")\
            .setOutputCol("bina")\
            .setReservedCols("sepal_length", "petal_width", "petal_length", "category")\
            .setThreshold(1.);

lda = Lda()\
            .setPredictionCol("lda_pred")\
            .setPredictionDetailCol("lda_pred_detail")\
            .setSelectedCol("category")\
            .setTopicNum(2)\
            .setRandomSeed(0)

pipeline = Pipeline()\
    .add(binarizer)\
    .add(binarizer)\
    .add(lda)

pipeline.fit(data1)
pipeline.transform(data2)

vector

Tipo de dado personalizado no Alink que suporta dois formatos:

  • vetor esparso (SparseVector)

    Exemplo: $4$1:0,1 2:0,2. O número entre os cifrões ($) indica o comprimento do vetor. Os valores após o segundo cifrão são pares índice-valor.

  • vetor denso (DenseVector)

    Exemplo: 0,1 0,2 0,3. Representa uma sequência de valores separados por espaços.

Nota

No Alink, se uma coluna for do tipo vetor, o nome do parâmetro geralmente é vectorColName.

Componentes Alink suportados pelo PyAlink Script

Use centenas de componentes Alink em um PyAlink Script, incluindo componentes para processamento de dados, engenharia de recursos e treinamento de modelos.

Nota

Atualmente, o componente PyAlink Script suporta componentes de pipeline e batch, mas não oferece suporte a componentes de stream.

Método 1: Usar o PyAlink Script isoladamente

Este tópico usa um exemplo de pontuação do conjunto de dados movielens com o modelo ItemCf para descrever como usar a plataforma Machine Learning Designer e recursos da Alibaba Cloud para executar um fluxo de trabalho implementado com um PyAlink Script. Siga o procedimento abaixo.

  1. Acesse a página do Machine Learning Designer e crie um pipeline em branco. Para mais informações, consulte Procedimento.

  2. Na lista de pipelines, selecione o pipeline em branco criado e clique em Open.

  3. Na caixa de pesquisa da lista de componentes à esquerda, procure por PyAlink Script e arraste-o para a tela à direita. Um nó de pipeline chamado PyAlink Script será gerado automaticamente na tela.

  4. Na tela, selecione o nó PyAlink Script-1. No painel à direita, configure os parâmetros nas abas PyAlink Script-1 e Parameter Settings.

    • Na aba Execution tuning, escreva seu código. O código a seguir serve como exemplo.

      from pyalink.alink import *
      def main(sources, sinks, parameter):
          PATH = "http://alink-test.oss-cn-beijing.aliyuncs.com/yuhe/movielens/"
          RATING_FILE = "ratings.csv"
          PREDICT_FILE = "predict.csv"
          RATING_SCHEMA_STRING = "user_id long, item_id long, rating int, ts long"
          ratingsData = CsvSourceBatchOp() \
                  .setFilePath(PATH + RATING_FILE) \
                  .setFieldDelimiter("\t") \
                  .setSchemaStr(RATING_SCHEMA_STRING)
          predictData = CsvSourceBatchOp() \
                  .setFilePath(PATH + PREDICT_FILE) \
                  .setFieldDelimiter("\t") \
                  .setSchemaStr(RATING_SCHEMA_STRING)
          itemCFModel = ItemCfTrainBatchOp() \
                  .setUserCol("user_id").setItemCol("item_id") \
                  .setRateCol("rating").linkFrom(ratingsData);
          itemCF = ItemCfRateRecommender() \
                  .setModelData(itemCFModel) \
                  .setItemCol("item_id") \
                  .setUserCol("user_id") \
                  .setReservedCols(["user_id", "item_id"]) \
                  .setRecommCol("prediction_score")
          result = itemCF.transform(predictData)
          result.link(sinks[0])
          BatchOperator.execute()

      Um PyAlink Script suporta quatro portas de saída. No script, use result.link(sinks[0]) para gravar dados na primeira porta de saída. Componentes downstream podem ler os dados do script conectando-se a esta primeira porta. Para mais informações, consulte Leitura e gravação de diferentes tipos de dados em um PyAlink Script.

    • Na aba Parameter Settings, defina o modo de execução e as especificações do nó.

      Parâmetro

      Descrição

      Execution tuning

      Modos suportados:

      • Select job running mode: Recomendado para tarefas com conjuntos de dados pequenos e para fins de depuração e validação.

      • DLC (Single-machine, multi-concurrency): Indicado para tarefas com grandes volumes de dados ou para ambientes de produção.

      • MaxCompute (Distributed): Executa o job no cluster Fully-managed Flink vinculado ao workspace.

      Fully-managed Flink

      Necessário apenas quando o Number of workers é job running mode ou MaxCompute (Distributed). Especifica o número de nós de execução. Se deixado em branco, o sistema aloca nós automaticamente com base nos dados da tarefa. Por padrão, este parâmetro está vazio.

      Fully-managed Flink (Distributed)

      Configure este parâmetro somente quando o Memory per worker (MB) estiver definido como running mode of the job ou MaxCompute (Distributed). Define o tamanho da memória de um único nó em MB. O valor deve ser um número inteiro positivo, sendo o padrão 8192.

      Fully-managed Flink (Distributed)

      Necessário apenas quando o CPU cores per worker estiver configurado para Job running mode ou MaxCompute (distributed). Define o número de núcleos de CPU para um único nó. O valor deve ser um inteiro positivo e, por padrão, permanece vazio.

      Fully-managed Flink (distributed)

      Tipo de recurso do nó DLC. O padrão é 2 vCPU + 8 GB Mem-ecs.g6.large.

  5. Acima da tela, clique em Select node specification to run script e, em seguida, clique no ícone de execução image para executar o PyAlink Script.

  6. Após a conclusão da tarefa, clique com o botão direito no nó Save na tela e selecione PyAlink Script-1 > View Data para visualizar os resultados.

    Nome da coluna

    Descrição

    user_id

    ID do usuário.

    item_id

    ID do filme.

    prediction_score

    Indica a preferência do usuário pelo filme. Esta pontuação serve como referência para recomendações.

Método 2: Combinar PyAlink Script com outros componentes

As portas de entrada e saída de um componente PyAlink Script são idênticas às de outros componentes de algoritmo no Machine Learning Designer. Conecte-as para criar um pipeline combinado, conforme ilustrado na figura a seguir.组合使用

Ler e gravar dados no PyAlink Script

  • Leitura de dados.

    • Ler de uma tabela MaxCompute: O script lê dados transmitidos por um componente upstream através de uma porta de entrada. O código a seguir apresenta um exemplo.

      train_data = sources[0]
      test_data = sources[1]

      No código, sources[0] representa a tabela MaxCompute conectada à primeira porta de entrada, e sources[1] representa a tabela conectada à segunda porta. O componente suporta até quatro portas de entrada.

    • Ler de um sistema de arquivos de rede: O script lê dados utilizando componentes de origem do Alink, como CsvSourceBatchOp e AkSourceBatchOp, dentro do código. É possível ler os seguintes tipos de arquivos:

      • Ler um arquivo compartilhado de uma rede via HTTP. O código abaixo serve como exemplo:

        ratingsData = CsvSourceBatchOp() \
                    .setFilePath(PATH + RATING_FILE) \
                    .setFieldDelimiter("\t") \
                    .setSchemaStr(RATING_SCHEMA_STRING)
      • Ler um arquivo OSS. Primeiro, localize o caminho de armazenamento temporário de dados do pipeline na aba pipeline properties. Este caminho, formatado como oss://<bucket-name>/<path>, é usado para operações de leitura/gravação no OSS. O código a seguir mostra um exemplo.

        model_data = AkSourceBatchOp().setFilePath("oss://xxxxxxxx/model_20220323.ak")
  • Gravação de dados.

    • Gravar em uma tabela MaxCompute: O script grava dados em um componente downstream através de uma porta de saída. Veja um exemplo no código a seguir.

      result0.link(sinks[0])
      result1.link(sinks[1])
      BatchOperator.execute()

      A linha result0.link(sinks[0]) grava os dados e os torna acessíveis pela primeira porta de saída. É possível gravar em até quatro tabelas de resultados, correspondentes a quatro portas de saída.

    • Gravar em um arquivo OSS. Use o caminho OSS definido no campo temporary data storage path do pipeline, localizado na aba pipeline properties. O código a seguir apresenta um exemplo.

      result.link(AkSinkBatchOp() \
                  .setFilePath("oss://xxxxxxxx/model_20220323.ak") \
                  .setOverwriteSink(True))
      BatchOperator.execute()

Exemplo: Implantar um modelo como serviço EAS

  1. Gere o modelo a ser implantado.

    A implantação de um modelo como serviço EAS só é possível se ele for um PipelineModel gerado pelo componente PyAlink Script. Use o código a seguir para gerar um arquivo PipelineModel. Para instruções sobre como executar o script, consulte Método 1: Usar o componente PyAlink Script isoladamente.

    from pyalink.alink import *
    def main(sources, sinks, parameter):
        PATH = "http://alink-test.oss-cn-beijing.aliyuncs.com/yuhe/movielens/"
        RATING_FILE = "ratings.csv"
        PREDICT_FILE = "predict.csv"
        RATING_SCHEMA_STRING = "user_id long, item_id long, rating int, ts long"
        ratingsData = CsvSourceBatchOp() \
                .setFilePath(PATH + RATING_FILE) \
                .setFieldDelimiter("\t") \
                .setSchemaStr(RATING_SCHEMA_STRING)
        predictData = CsvSourceBatchOp() \
                .setFilePath(PATH + PREDICT_FILE) \
                .setFieldDelimiter("\t") \
                .setSchemaStr(RATING_SCHEMA_STRING)
        itemCFModel = ItemCfTrainBatchOp() \
                .setUserCol("user_id").setItemCol("item_id") \
                .setRateCol("rating").linkFrom(ratingsData);
        itemCF = ItemCfRateRecommender() \
                .setModelData(itemCFModel) \
                .setItemCol("item_id") \
                .setUserCol("user_id") \
                .setReservedCols(["user_id", "item_id"]) \
                .setRecommCol("prediction_score")
        model = PipelineModel(itemCF)
        model.save().link(AkSinkBatchOp() \
                .setFilePath("oss://<your_bucket_name>/model.ak") \
                .setOverwriteSink(True))
        BatchOperator.execute()

    Substitua <your_bucket_name> pelo nome do seu bucket OSS.

    Importante

    Certifique-se de ter permissões de leitura para o caminho do conjunto de dados configurado em PATH. Caso contrário, o componente falhará durante a execução.

  2. Gere um arquivo de configuração EAS.

    Execute o script a seguir para gravar a saída em um arquivo config.json.

    # EAS configuration file
    import json
    # Generate EAS model configuration
    model_config = {}
    # Schema of the data that EAS receives
    model_config['inputDataSchema'] = "id long, movieid long" 
    model_config['modelVersion'] = "v0.2"
    eas_config = {
        "name": "recomm_demo",
        "model_path": "http://xxxxxxxx/model.ak",
        "processor": "alink_outer_processor",
        "metadata": {
            "instance": 1,
            "memory": 2048,
            "region":"China (Beijing)"
        },
        "model_config": model_config
    }
    print(json.dumps(eas_config, indent=4))

    Principais parâmetros no arquivo config.json:

    • name: Nome do serviço de modelo implantado.

    • model_path: Caminho OSS onde o arquivo PipelineModel está armazenado. Altere este valor para o caminho OSS real do seu arquivo de modelo.

    Para obter explicações sobre outros parâmetros do arquivo config.json, consulte Referência de comandos.

  3. Implante o modelo como um serviço EAS.

    Implante o modelo usando o cliente eascmd. Para instruções de configuração do cliente, consulte Baixar e configurar o cliente. Por exemplo, em um sistema Windows de 64 bits, execute este comando:

    eascmdwin64.exe create config.json