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
O estado Map recebe a entrada e extrai um array de itens usando
ItemsPath(ou lê itens de uma source de dados do OSS usandoItemReader).Se
ItemConstructorestiver configurado, cada item é transformado.Caso
ItemBatcheresteja configurado, os itens são agrupados em lotes.Cada item ou lote é enviado para uma execução de sub-workflow definida por
Processor, executada no modo especificado porProcessorConfig.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).
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 |
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 |
|
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
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.
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/ |
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": "*"
}
]
}