Todos os produtos
Search
Central de documentação

CloudFlow:Distributed mode

Última atualização: Jun 28, 2026

O modo distribuído executa cada iteração de um estado Map como uma execução independente de sub-workflow, permitindo processamento paralelo de alta concorrência, ideal para big data e computação paralela.

Utilize o modo distribuído quando um único estado Map não conseguir processar sua carga de trabalho sequencialmente. Esse modo distribui o trabalho entre execuções concorrentes de sub-workflows — cada uma com seu próprio ciclo de vida — e processa dados em paralelo a partir de entradas inline ou buckets do Object Storage Service (OSS).

Conceitos principais

Termo

Definição

Modo distribuído

Modo de processamento do estado Map que trata grandes volumes de dados simultaneamente por meio de execuções de sub-workflows. Objetos em buckets do OSS podem servir como source de dados.

Execução de sub-workflow

Uma única iteração de um estado Map distribuído. A definição vem do campo Processor, e cada sub-workflow recebe um item (ou um lote) como entrada. Os sub-workflows podem ser executados no modo padrão ou no modo expresso. Para mais detalhes, consulte Modo padrão e modo expresso.

Política de tolerância a falhas

Quando uma execução de sub-workflow falha, o estado Map termina com falha. Em estados Map com múltiplos itens, configure políticas de tolerância a falhas para permitir que o processamento continue nos itens restantes mesmo se execuções individuais de sub-workflows falharem.

Como funciona

  1. O estado Map recebe a entrada e extrai um array de itens usando ItemsPath (ou lê itens de uma source de dados do OSS usando ItemReader).

  2. Se ItemConstructor estiver configurado, cada item é transformado.

  3. Caso ItemBatcher esteja configurado, os itens são agrupados em lotes.

  4. Cada item ou lote é enviado para uma execução de sub-workflow definida por Processor, executada no modo especificado por ProcessorConfig.

  5. Após a conclusão de todas as execuções de sub-workflows, o estado Map coleta os resultados em um array JSON (ou os grava no OSS usando ResultWriter).

Importante

Quando ItemsPath, ItemConstructor e ItemBatcher estão configurados simultaneamente, a ordem de execução é: ItemsPath, depois ItemConstructor e, por fim, ItemBatcher.

Exemplo básico

O workflow abaixo define um estado Map distribuído que lê itens de $Input.Items e executa cada item em um estado Pass no modo Express:

Type: StateMachine
Name: MyWorkflow
SpecVersion: v1
StartAt: Map
States:
  - Type: Map
    Name: Map
    ProcessorConfig:
      ExecutionMode: Express
    ItemsPath: $Input.Items
    Processor:
      StartAt: Pass
      States:
        - Type: Pass
          Name: Pass
          End: true
    End: true

Entrada da máquina de estados:

{
  "Items": [
    {"key_1":"value_1"},
    {"key_2":"value_2"},
    {"key_3":"value_3"}
  ]
}

Isso gera três execuções de sub-workflows, cada uma executando a seguinte definição:

Type: StateMachine
Name: Map
SpecVersion: v1
StartAt: Pass
States:
  - Type: Pass
    Name: Pass
    End: true

Entradas dos sub-workflows:

// Sub-workflow 1
{"key_1":"value_1"}

// Sub-workflow 2
{"key_2":"value_2"}

// Sub-workflow 3
{"key_3":"value_3"}

Após a conclusão de todos os sub-workflows, a saída do estado Map é um array JSON contendo todas as saídas dos sub-workflows:

{
    "Items": [
        {
            "key_1": "value_1"
        },
        {
            "key_2": "value_2"
        },
        {
            "key_3": "value_3"
        }
    ]
}

Referência de campos

Todos os campos disponíveis em um estado Map no modo distribuído:

Campo

Tipo

Obrigatório

Descrição

Exemplo

Name

string

Sim

Nome do estado.

my-state-name

Description

string

Não

Descrição do estado.

describe it here

Type

string

Sim

Tipo do estado.

Map

InputConstructor

map[string]any

Não

Construtor de entrada.

Consulte Entradas e saídas.

ItemsPath

string

Sim

Expressão para extrair um array da entrada.

Consulte ItemsPath.

ItemBatcher

ItemBatcher

Não

Combina vários itens em um lote como entrada do sub-workflow.

Consulte ItemBatcher.

ItemReader

ItemReader

Não

Lê dados de buckets do OSS.

Consulte ItemReader.

ItemConstructor

ItemConstructor

Não

Transforma itens usando $Item para referenciar valores originais da entrada.

Consulte ItemConstructor.

ResultWriter

ResultWriter

Não

Grava resultados dos sub-workflows em um bucket do OSS especificado.

Consulte ResultWriter.

MaxConcuccency

int

Não

Número máximo de execuções concorrentes de sub-workflows.

40

MaxItems

MaxItems

Não

Número máximo de itens que o estado Map processa.

Consulte MaxItems.

ToleratedFailurePercentage

ToleratedFailurePercentage

Não

Percentual de falhas em execuções de sub-workflows a ser tolerado.

Consulte ToleratedFailurePercentage.

ToleratedFailureCount

ToleratedFailureCount

Não

Quantidade de falhas em execuções de sub-workflows a serem toleradas.

Consulte ToleratedFailureCount.

Processor

Processor

Sim

Processador do Map que define o sub-workflow.

Consulte Processor.

ProcessorConfig

ProcessorConfig

Sim

Configuração do processador.

Consulte ProcessorConfig.

OutputConstructor

map[string]any

Não

Construtor de saída.

Consulte OutputConstructor.

Next

string

Não

Próximo estado a ser executado após a conclusão deste estado. Deixe em branco se End for true.

my-next-state

End

bool

Não

Indica se este é o estado terminal do escopo atual.

true

Retry

Retry

Não

Política de nova tentativa em caso de erro.

Consulte Tratamento de erros.

Catch

Catch

Não

Política de captura de erros.

Consulte Tratamento de erros.

Detalhes dos campos

ItemsPath

Extrai um array da entrada do estado Map. Cada elemento do array se torna a entrada para uma execução de sub-workflow. Referencie dados com as variáveis $Context ou $Input:

$Input.FieldA

Processor

Define o sub-workflow que processa cada item ou lote.

Campo

Tipo

Obrigatório

Descrição

Exemplo

States

array

Sim

Estados presentes no sub-workflow.

Veja o exemplo abaixo.

StartAt

string

Sim

Primeiro estado a ser executado.

my start task

Processor:
   StartAt: Pass1
   States:
     - Type: Pass
       Name: Pass1
       End: true

ProcessorConfig

Especifica o modo de execução para os sub-workflows.

Campo

Tipo

Obrigatório

Descrição

Exemplo

ExecutionMode

string

Sim

Modo de execução dos sub-workflows.

Express

ItemReader

Lê dados de entrada de buckets do OSS, oferecendo suporte a volumes maiores de dados.

Campo

Tipo

Obrigatório

Descrição

Exemplo

SourceType

string

Sim

Tipo da source de dados. Valores válidos: OSS_CSV, OSS_JSON_LIST, OSS_OBJECTS, OSS_INVENTORY_FILES.

OSS_CSV

SourceParameters

SourceParameters

Não

Parâmetros de localização da source de dados.

Consulte SourceParameters.

ReaderConfig

ReaderConfig

Não

Configuração do leitor.

Consulte ReaderConfig.

SourceParameters

Campo

Tipo

Obrigatório

Descrição

Exemplo

Bucket

string

Não

Nome do bucket do OSS.

example-bucket

ObjectName

string

Não

Nome do objeto.

object_name_1

Prefix

string

Não

Filtro de prefixo do nome do objeto. Se deixado em branco, todos os objetos correspondentes serão retornados.

example-prefix

ReaderConfig

Campo

Tipo

Obrigatório

Descrição

Exemplo

CSVHeaders

[]string

Não

Cabeçalhos das colunas na primeira linha do arquivo CSV.

ColA,ColB,ColC

Leitura de arquivos CSV do OSS

Leia dados de um arquivo CSV armazenado em um bucket do OSS. Neste exemplo, example-object.csv está armazenado em example-bucket:

Type: StateMachine
Name: MyWorkflow
SpecVersion: v1
StartAt: Map
States:
  - Type: Map
    Name: Map
    ProcessorConfig:
      ExecutionMode: Express
    Processor:
      StartAt: Pass
      States:
        - Type: Pass
          Name: Pass
          End: true
    ItemReader:
      SourceType: OSS_CSV
      SourceParameters:
        Bucket: example-bucket
        ObjectName: example-object.csv
    End: true

Conteúdo do arquivo CSV:

ColA,ColB,ColC
col_a_1,col_b_1,col_c_1

Cada linha se torna a entrada do sub-workflow:

{
  "ColA": "col_a_1",
  "ColB": "col_b_1",
  "ColC": "col_c_1"
}

Leitura de arrays JSON do OSS

Importante

O arquivo JSON deve conter um array JSON.

Leia um array JSON de example-object.json em example-bucket:

Type: StateMachine
Name: MyWorkflow
SpecVersion: v1
StartAt: Map
States:
  - Type: Map
    Name: Map
    ProcessorConfig:
      ExecutionMode: Express
    Processor:
      StartAt: Pass
      States:
        - Type: Pass
          Name: Pass
          End: true
    ItemReader:
      SourceType: OSS_JSON_LIST
      SourceParameters:
        Bucket: example-bucket
        ObjectName: example-object.json
    End: true

Conteúdo do arquivo JSON:

[
  {
    "key_1": "value_1"
  }
]

Cada elemento se torna a entrada do sub-workflow:

{
  "key_1": "value_1"
}

Leitura de objetos do OSS

Liste objetos com um prefixo específico e passe seus metadados como entrada do sub-workflow. Este exemplo lê objetos com o prefixo example-prefix de example-bucket:

Type: StateMachine
Name: MyWorkflow
SpecVersion: v1
StartAt: Map
States:
  - Type: Map
    Name: Map
    ProcessorConfig:
      ExecutionMode: Express
    Processor:
      StartAt: Pass
      States:
        - Type: Pass
          Name: Pass
          End: true
    ItemReader:
      SourceType: OSS_OBJECTS
      SourceParameters:
        Bucket: example-bucket
        Prefix: example-prefix
    End: true

Dada esta estrutura de bucket:

example-bucket
   └── example-prefix/object_1

Os metadados de cada objeto tornam-se a entrada do sub-workflow:

{
  "XMLName": {
    "Space": "",
    "Local": "Contents"
  },
  "Key": "example-prefix/object_1",
  "Type": "Normal",
  "Size": 268435,
  "ETag": "\"50B06D6680D86F04138HSN612EF5DEC6\"",
  "Owner": {
    "XMLName": {
      "Space": "",
      "Local": ""
    },
    "ID": "",
    "DisplayName": ""
  },
  "LastModified": "2024-01-01T01:01:01Z",
  "StorageClass": "Standard",
  "RestoreInfo": ""
}

Leitura de arquivos de inventário do OSS

Leia itens de um manifesto de inventário do OSS. Este exemplo lê inventory/2024-01-01T01-01Z/manifest.json de example-bucket:

Type: StateMachine
Name: MyWorkflow
SpecVersion: v1
StartAt: Map
States:
  - Type: Map
    Name: Map
    ProcessorConfig:
      ExecutionMode: Express
    Processor:
      StartAt: Pass
      States:
        - Type: Pass
          Name: Pass
          End: true
    ItemReader:
      SourceType: OSS_INVENTORY_FILES
      SourceParameters:
        Bucket: example-bucket
        ObjectName: inventory/2024-01-01T01-01Z/manifest.json
    ItemConstructor:
      Key.$: $Item.Key
    End: true

Conteúdo do arquivo de inventário:

"example-bucket","object_name_1"
"example-bucket","object_name_2"

Entrada do primeiro sub-workflow:

{
  "Bucket": "example-bucket",
  "Key": "object_name_1"
}

ItemBatcher

Agrupa vários itens em lotes. Cada lote se torna a entrada para uma única execução de sub-workflow.

Campo

Tipo

Obrigatório

Descrição

Exemplo

MaxItemsPerBatch

int

Não

Máximo de itens por lote.

Consulte o exemplo.

MaxInputBytesPerBatch

int

Não

Máximo de bytes por lote. Unidade: byte.

Consulte o exemplo.

BatchInput

map[string]any

Não

Dados adicionais a serem incluídos na entrada de cada lote.

Consulte o exemplo.

Lote por contagem de itens

Defina MaxItemsPerBatch para controlar quantos itens vão para cada lote:

Type: StateMachine
Name: MyWorkflow
SpecVersion: v1
StartAt: Map
States:
  - Type: Map
    Name: Map
    ProcessorConfig:
      ExecutionMode: Express
    ItemsPath: $Input.Items
    Processor:
      StartAt: Pass
      States:
        - Type: Pass
          Name: Pass
          End: true
    ItemBatcher:
      MaxItemsPerBatch: 2
    End: true

Entrada da máquina de estados:

{
  "Items": [
    {"key_1":"value_1"},
    {"key_2":"value_2"},
    {"key_3":"value_3"},
    {"key_4":"value_4"},
    {"key_5":"value_5"}
  ]
}

Isso produz três execuções de sub-workflows:

// Sub-workflow 1
{
  "Items": [
    {"key_1":"value_1"},
    {"key_2":"value_2"}
  ]
}

// Sub-workflow 2
{
  "Items": [
    {"key_1":"value_3"},
    {"key_2":"value_4"}
  ]
}

// Sub-workflow 3
{
  "Items": [
    {"key_1":"value_5"}
  ]
}

Lote por tamanho em bytes

Defina MaxInputBytesPerBatch para limitar o tamanho de cada lote em bytes.

Importante
  • O ItemBatcher adiciona chaves de metadados a cada lote. O tamanho total do lote inclui essas chaves adicionais.

  • A unidade para MaxInputBytesPerBatch é byte.

Type: StateMachine
Name: MyWorkflow
SpecVersion: v1
StartAt: Map
States:
  - Type: Map
    Name: Map
    ProcessorConfig:
      ExecutionMode: Express
    ItemsPath: $Input.Items
    Processor:
      StartAt: Pass
      States:
        - Type: Pass
          Name: Pass
          End: true
    ItemBatcher:
      MaxInputBytesPerBatch: 50
    End: true

Entrada da máquina de estados:

{
  "Items":[
    {"Key":1},
    {"key":2},
    {"Key":3},
    {"Key":4},
    {"Key":5}
  ]
}

Isso produz três execuções de sub-workflows:

// Sub-workflow 1
{
  "Items":[
    {"Key":1},
    {"key":2}
  ]
}

// Sub-workflow 2
{
  "Items":[
    {"Key":3},
    {"key":4}
  ]
}

// Sub-workflow 3
{
  "Items":[
    {"Key":5}
  ]
}

Adição de dados compartilhados aos lotes

Use BatchInput para incluir dados adicionais da entrada da máquina de estados em cada lote:

Type: StateMachine
Name: MyWorkflow
SpecVersion: v1
StartAt: Map
States:
  - Type: Map
    Name: Map
    ProcessorConfig:
      ExecutionMode: Express
    ItemsPath: $Input.Items
    Processor:
      StartAt: Pass
      States:
        - Type: Pass
          Name: Pass
          End: true
    ItemBatcher:
      MaxInputBytesPerBatch: 70
      BatchInput:
        InputKey.$: $Input.Key
    End: true

Entrada da máquina de estados:

{
  "Key":"value",
  "Items":[
    {"Key":1},
    {"key":2},
    {"Key":3},
    {"Key":4},
    {"Key":5}
  ]
}

Cada sub-workflow recebe tanto os dados de BatchInput quanto os itens do lote:

// Sub-workflow 1
{
  "BatchInput":{
    "InputKey":"value"
  },
  "Items":[
    {"Key":1},
    {"key":2}
  ]
}

// Sub-workflow 2
{
  "BatchInput":{
    "InputKey":"value"
  },
  "Items":[
    {"Key":3},
    {"key":4}
  ]
}

// Sub-workflow 3
{
  "BatchInput":{
    "InputKey":"value"
  },
  "Items":[
    {"Key":5}
  ]
}

ItemConstructor

Transforma cada item antes de passá-lo para uma execução de sub-workflow. Use $Item para referenciar o item original e $Input para referenciar a entrada da máquina de estados:

Type: StateMachine
Name: MyWorkflow
SpecVersion: v1
StartAt: Map
States:
  - Type: Map
    Name: Map
    ProcessorConfig:
      ExecutionMode: Express
    ItemsPath: $Input.Items
    Processor:
      StartAt: Pass
      States:
        - Type: Pass
          Name: Pass
          End: true
    ItemBatcher:
      MaxInputBytesPerBatch: 200
      BatchInput:
        InputKey.$: $Input.Key
    ItemConstructor:
      ConstructedKey.$: $Item.Key
      InputKey.$: $Input.Key
    End: true

Entrada da máquina de estados:

{
  "Key":"value",
  "Items":[
    {"Key":1},
    {"Key":2},
    {"Key":3},
    {"Key":4},
    {"Key":5}
  ]
}

Cada item é transformado pelo ItemConstructor antes do agrupamento em lotes:

// Sub-workflow 1
{
  "BatchInput": {
    "InputKey": "value"
  },
  "Items": [
    {
      "InputKey": "value",
      "ConstructedKey": 1
    },
    {
      "InputKey": "value",
      "ConstructedKey": 2
    },
    {
      "InputKey": "value",
      "ConstructedKey": 3
    }
  ]
}

// Sub-workflow 2
{
  "BatchInput": {
    "InputKey": "value"
  },
  "Items": [
    {
      "InputKey": "value",
      "ConstructedKey": 4
    },
    {
      "InputKey": "value",
      "ConstructedKey": 5
    }
  ]
}

ResultWriter

Grava os resultados dos sub-workflows em um bucket do OSS. Configure o ResultWriter quando a saída combinada de todas as execuções de sub-workflows puder exceder o limite de tamanho de saída do estado.

Campo

Tipo

Obrigatório

Descrição

Exemplo

Parameters

Parameters

Sim

Destino no OSS para os resultados.

Veja abaixo.

Parameters

Campo

Tipo

Obrigatório

Descrição

Exemplo

Bucket

string

Sim

Bucket do OSS de destino.

example-bucket

Prefix

string

Sim

Prefixo do nome do objeto para os arquivos de resultado.

example-prefix/

Nota

As entradas e saídas de estados possuem limites de tamanho. Para estados Map com muitos itens, a saída combinada pode exceder esse limite. Configure o ResultWriter para armazenar os resultados no OSS.

O exemplo a seguir combina ResultWriter com tolerância a falhas. O estado Map tolera até 30% de falhas e grava todos os resultados (sucessos e falhas) no OSS:

Type: StateMachine
Name: MyWorkflow
SpecVersion: v1
StartAt: Map
States:
  - Type: Map
    Name: Map
    ProcessorConfig:
      ExecutionMode: Express
    ItemsPath: $Input.Items
    ItemConstructor:
      Key.$: $Item.Key
      FailedValue.$: $Input.FailedValue
    ToleratedFailurePercentage: 30
    Processor:
      StartAt: Choice
      States:
        - Type: Choice
          Name: Choice
          Branches:
            - Condition: $Input.Key > $Input.FailedValue
              Next: Fail
          Default: Succeed
        - Type: Succeed
          Name: Succeed
          End: true
        - Type: Fail
          Name: Fail
          Code: MockError
          End: true
    ResultWriter:
      Parameters:
        Bucket: example-bucket
        Prefix: example-prefix/
    End: true

Entrada da máquina de estados:

{
  "FailedValue": 4,
  "Items": [
    {"Key": 1},
    {"Key": 2},
    {"Key": 3},
    {"Key": 4},
    {"Key": 5}
  ]
}

Estrutura do arquivo de resultados

O ResultWriter produz três tipos de arquivos no local especificado do OSS:

manifest.json (example-prefix/map-run-name/manifest.json):

{
    "DestinationBucket": "example-bucket",
    "MapRunName": "map-run-name",
    "ResultFiles": {
        "FAILED": [
            {
                "ObjectName": "example-prefix/map-run-name/FAILED_0.json",
                "Size": 262
            }
        ],
        "SUCCEED": [
            {
                "ObjectName": "example-prefix/map-run-name/SUCCEED_0.json",
                "Size": 1057
            }
        ]
    }
}

FAILED_0.json (example-prefix/map-run-name/FAILED_0.json):

[
    {
        "ExecutionName": "execution-name-5",
        "FlowName": "example",
        "Input": "{\"FailedValue\":4,\"Key\":5}",
        "Output": "{\"ErrorCode\":\"MockError\"}",
        "Status": "Failed",
        "StartedTime": "rfc3339-format-time-string",
        "StoppedTime": "rfc3339-format-time-string"
    }
]

SUCCEED_0.json (example-prefix/map-run-name/SUCCEED_0.json):

[
    {
        "ExecutionName": "execution-name-1",
        "FlowName": "example",
        "Input": "{\"FailedValue\":4,\"Key\":1}",
        "Output": "{\"FailedValue\":4,\"Key\":1}",
        "Status": "Succeeded",
        "StartedTime": "rfc3339-format-time-string",
        "StoppedTime": "rfc3339-format-time-string"
    },
    {
        "ExecutionName": "execution-name-2",
        "FlowName": "example",
        "Input": "{\"FailedValue\":4,\"Key\":2}",
        "Output": "{\"FailedValue\":4,\"Key\":2}",
        "Status": "Succeeded",
        "StartedTime": "rfc3339-format-time-string",
        "StoppedTime": "rfc3339-format-time-string"
    },
    {
        "ExecutionName": "execution-name-3",
        "FlowName": "example",
        "Input": "{\"FailedValue\":4,\"Key\":3}",
        "Output": "{\"FailedValue\":4,\"Key\":3}",
        "Status": "Succeeded",
        "StartedTime": "rfc3339-format-time-string",
        "StoppedTime": "rfc3339-format-time-string"
    },
    {
        "ExecutionName": "execution-name-4",
        "FlowName": "example",
        "Input": "{\"FailedValue\":4,\"Key\":4}",
        "Output": "{\"FailedValue\":4,\"Key\":4}",
        "Status": "Succeeded",
        "StartedTime": "rfc3339-format-time-string",
        "StoppedTime": "rfc3339-format-time-string"
    }
]

MaxItems

Limita o número de itens que o estado Map processa. Se a source de dados contiver mais itens do que o valor definido em MaxItems, apenas os primeiros MaxItems itens serão processados.

Por exemplo, se um bucket do OSS contiver 10.000 objetos e MaxItems estiver definido como 1.000, o estado Map processará apenas 1.000 objetos.

MaxConcurrency

Limita o número de execuções de sub-workflows que rodam simultaneamente. Por exemplo, se o estado Map processar 10.000 itens e MaxConcurrency estiver definido como 100, no máximo 100 sub-workflows serão executados ao mesmo tempo.

ToleratedFailurePercentage

Define a porcentagem de execuções de sub-workflows que podem falhar antes que o próprio estado Map falhe. Por exemplo, com 10.000 itens e ToleratedFailurePercentage definido como 10, o estado Map tolera até 1.000 falhas.

ToleratedFailureCount

Define a quantidade de execuções de sub-workflows que podem falhar antes que o próprio estado Map falhe. Por exemplo, com 10.000 itens e ToleratedFailureCount definido como 10, o estado Map tolera até 10 falhas.

Limites de cota

A tabela a seguir lista os limites de cota padrão para o modo distribuído. Para solicitar uma cota maior, envie um ticket.

Nome da cota

Descrição

Valor padrão

MaxOpenMapRun

Número máximo de estados Map distribuídos executados simultaneamente por conta e por região

10

MaxConcurrency

Número máximo de execuções concorrentes de sub-workflows por MapRun

300

MaxItems

Número máximo de itens por MapRun

10.000

Permissões necessárias

A execução de estados Map no modo distribuído requer as seguintes permissões do OSS. Conceda essas ações em uma política de função. Para mais detalhes, consulte Criar uma política.

{
    "Version": "1",
    "Statement": [
        {
            "Effect": "Allow",
            "Action": [
                "oss:HeadObject",
                "oss:GetObjectMeta",
                "oss:GetObject",
                "oss:PutObject",
                "oss:ListObjectsV2",
                "oss:ListObjects",
                "oss:InitiateMultipartUpload",
                "oss:UploadPart",
                "oss:CompleteMultipartUpload",
                "oss:AbortMultipartUpload",
                "oss:ListMultipartUploads",
                "oss:ListParts"
            ],
            "Resource": "*"
        }
    ]
}