Workflowschritte parallel ausführen

Parallele Schritte können die Gesamtausführungszeit für einen Workflow verkürzen, indem mehrere blockierende Aufrufe gleichzeitig ausgeführt werden.

Blockierende Aufrufe wie „sleep“, HTTP-Aufrufe und Callbacks können Zeit in Anspruch nehmen, von Millisekunden bis zu Tagen. Parallele Schritte sollen bei solchen gleichzeitigen, lang andauernden Vorgängen helfen. Wenn ein Workflow mehrere blockierende Aufrufe ausführen muss, die unabhängig voneinander sind, kann die Verwendung paralleler Zweige die Gesamtausführungszeit verkürzen, indem die Aufrufe gleichzeitig gestartet werden und gewartet wird, bis alle abgeschlossen sind.

Wenn Ihr Workflow beispielsweise Kundendaten aus mehreren unabhängigen Systemen abrufen muss, bevor er fortgesetzt wird, ermöglichen parallele Zweige gleichzeitige API-Anfragen. Wenn es fünf Systeme gibt und die Antwort jeweils zwei Sekunden dauert, kann die sequenzielle Ausführung der Schritte in einem Workflow mindestens zehn Sekunden dauern. Bei paralleler Ausführung kann es nur zwei Sekunden dauern.

Parallelen Schritt erstellen

Erstellen Sie einen parallel-Schritt, um einen Teil Ihres Workflows zu definieren, in dem zwei oder mehr Schritte gleichzeitig ausgeführt werden können.

YAML

  - PARALLEL_STEP_NAME:
      parallel:
        exception_policy: POLICY
        shared: [VARIABLE_A, VARIABLE_B, ...]
        concurrency_limit: CONCURRENCY_LIMIT
        BRANCHES_OR_FOR:
          ...

JSON

  [
    {
      "PARALLEL_STEP_NAME": {
        "parallel": {
          "exception_policy": "POLICY",
          "shared": [
            "VARIABLE_A",
            "VARIABLE_B",
            ...
          ],
          "concurrency_limit": "CONCURRENCY_LIMIT",
          "BRANCHES_OR_FOR":
          ...
        }
      }
    }
  ]

Ersetzen Sie Folgendes:

  • PARALLEL_STEP_NAME: der Name des parallelen Schritts.
  • POLICY (optional): bestimmt die Aktion, die andere Zweige ausführen, wenn eine nicht behandelte Ausnahme auftritt. Die Standardrichtlinie continueAll führt zu keiner weiteren Aktion und alle anderen Zweige versuchen, ausgeführt zu werden. continueAll ist die einzige Richtlinie, die derzeit unterstützt wird.
  • VARIABLE_A, VARIABLE_B usw.: eine Liste von beschreibbaren Variablen mit übergeordnetem Bereich, die Zuweisungen innerhalb des parallelen Schritts ermöglichen. Weitere Informationen finden Sie unter Gemeinsame Variablen.
  • CONCURRENCY_LIMIT (optional): die maximale Anzahl von Zweigen und Iterationen, die gleichzeitig innerhalb einer einzelnen Workflow-Ausführung ausgeführt werden können, bevor weitere Zweige und Iterationen in die Warteschlange gestellt werden. Dies gilt nur für einen einzelnen parallel-Schritt und wird nicht kaskadiert. Muss eine positive Ganzzahl sein und kann entweder ein Literalwert oder ein Ausdruck sein. Weitere Informationen finden Sie unter Limits für die Parallelität.
  • BRANCHES_OR_FOR: Verwenden Sie entweder branches oder for, um eine der folgenden Optionen anzugeben:
    • Zweige, die gleichzeitig ausgeführt werden können.
    • Eine Schleife, in der Iterationen gleichzeitig ausgeführt werden können.

Wichtige Hinweise:

  • Parallele Zweige und Iterationen können in beliebiger Reihenfolge ausgeführt werden und die Reihenfolge kann sich bei jeder Ausführung ändern.
  • Parallele Schritte können andere, verschachtelte parallele Schritte bis zum Limit für die Tiefe enthalten. Siehe Kontingente und Limits.
  • Weitere Informationen finden Sie auf der Seite mit der Syntaxreferenz für parallele Schritte.

Experimentelle Funktion durch parallelen Schritt ersetzen

Wenn Sie experimental.executions.map verwenden, um parallele Arbeit zu unterstützen, können Sie Ihren Workflow migrieren, um stattdessen parallele Schritte zu verwenden und normale for-Schleifen parallel auszuführen. Beispiele finden Sie unter Experimentelle Funktion durch parallelen Schritt ersetzen.

Beispiele

Diese Beispiele veranschaulichen die Syntax.

Vorgänge parallel ausführen (mit Zweigen)

Wenn Ihr Workflow mehrere und unterschiedliche Schrittfolgen hat, die gleichzeitig ausgeführt werden können, kann die Platzierung in parallelen Zweigen die Gesamtzeit verkürzen, die für die Ausführung dieser Schritte erforderlich ist.

Im folgenden Beispiel wird eine Nutzer-ID als Argument an den Workflow übergeben und Daten werden parallel von zwei verschiedenen Diensten abgerufen. Gemeinsame Variablen ermöglichen das Schreiben von Werten in die Zweige und das Lesen nach Abschluss der Zweige

YAML

main:
  params: [input]
  steps:
    - init:
        assign:
          - userProfile: {}
          - recentItems: []
    - enrichUserData:
        parallel:
          shared: [userProfile, recentItems]  # userProfile and recentItems are shared to make them writable in the branches
          branches:
            - getUserProfileBranch:
                steps:
                  - getUserProfile:
                      call: http.get
                      args:
                        url: '${"https://example.com/users/" + input.userId}'
                      result: userProfile
            - getRecentItemsBranch:
                steps:
                  - getRecentItems:
                      try:
                        call: http.get
                        args:
                          url: '${"https://example.com/items?userId=" + input.userId}'
                        result: recentItems
                      except:
                        as: e
                        steps:
                          - ignoreError:
                              assign:  # continue with an empty list if this call fails
                                - recentItems: []

JSON

{
  "main": {
    "params": [
      "input"
    ],
    "steps": [
      {
        "init": {
          "assign": [
            {
              "userProfile": {}
            },
            {
              "recentItems": []
            }
          ]
        }
      },
      {
        "enrichUserData": {
          "parallel": {
            "shared": [
              "userProfile",
              "recentItems"
            ],
            "branches": [
              {
                "getUserProfileBranch": {
                  "steps": [
                    {
                      "getUserProfile": {
                        "call": "http.get",
                        "args": {
                          "url": "${\"https://example.com/users/\" + input.userId}"
                        },
                        "result": "userProfile"
                      }
                    }
                  ]
                }
              },
              {
                "getRecentItemsBranch": {
                  "steps": [
                    {
                      "getRecentItems": {
                        "try": {
                          "call": "http.get",
                          "args": {
                            "url": "${\"https://example.com/items?userId=\" + input.userId}"
                          },
                          "result": "recentItems"
                        },
                        "except": {
                          "as": "e",
                          "steps": [
                            {
                              "ignoreError": {
                                "assign": [
                                  {
                                    "recentItems": []
                                  }
                                ]
                              }
                            }
                          ]
                        }
                      }
                    }
                  ]
                }
              }
            ]
          }
        }
      }
    ]
  }
}

Elemente parallel verarbeiten (mit einer parallelen Schleife)

Wenn Sie für jedes Element in einer Liste dieselbe Aktion ausführen müssen, können Sie die Ausführung mit einer parallelen Schleife beschleunigen. Mit einer parallelen Schleife können mehrere Schleifeniterationen parallel ausgeführt werden. Im Gegensatz zu regulären „for“-Schleifen können Iterationen in beliebiger Reihenfolge ausgeführt werden.

Im folgenden Beispiel werden eine Reihe von Nutzermitteilungen in einer parallelen for-Schleife verarbeitet:

YAML

main:
  params: [input]
  steps:
    - sendNotifications:
        parallel:
          for:
            value: notification
            in: ${input.notifications}
            steps:
              - notify:
                  call: http.post
                  args:
                    url: https://example.com/sendNotification
                    body:
                      notification: ${notification}

JSON

{
  "main": {
    "params": [
      "input"
    ],
    "steps": [
      {
        "sendNotifications": {
          "parallel": {
            "for": {
              "value": "notification",
              "in": "${input.notifications}",
              "steps": [
                {
                  "notify": {
                    "call": "http.post",
                    "args": {
                      "url": "https://example.com/sendNotification",
                      "body": {
                        "notification": "${notification}"
                      }
                    }
                  }
                }
              ]
            }
          }
        }
      }
    ]
  }
}

Daten aggregieren (mit einer parallelen Schleife)

Sie können eine Reihe von Elementen verarbeiten und gleichzeitig Daten aus den Vorgängen erfassen, die für jedes Element ausgeführt werden. Beispielsweise können Sie die IDs der erstellten Elemente verfolgen oder eine Liste von Elementen mit Fehlern führen.

Im folgenden Beispiel geben zehn separate Abfragen für ein öffentliches BigQuery-Dataset jeweils die Anzahl der Wörter in einem Dokument oder einer Reihe von Dokumenten zurück. Mit einer gemeinsamen Variablen kann die Anzahl der Wörter akkumuliert und nach Abschluss aller Iterationen gelesen werden. Nachdem die Anzahl der Wörter in allen Dokumenten berechnet wurde, gibt der Workflow die Gesamtzahl zurück.

YAML

# Use a parallel loop to make ten queries to a public BigQuery dataset and
# use a shared variable to accumulate a count of words; after all iterations
# complete, return the total number of words across all documents
main:
  params: [input]
  steps:
    - init:
        assign:
          - numWords: 0
          - corpuses:
              - sonnets
              - various
              - 1kinghenryvi
              - 2kinghenryvi
              - 3kinghenryvi
              - comedyoferrors
              - kingrichardiii
              - titusandronicus
              - tamingoftheshrew
              - loveslabourslost
    - runQueries:
        parallel:  # 'numWords' is shared so it can be written within the parallel loop
          shared: [numWords]
          for:
            value: corpus
            in: ${corpuses}
            steps:
              - runQuery:
                  call: googleapis.bigquery.v2.jobs.query
                  args:
                    projectId: ${sys.get_env("GOOGLE_CLOUD_PROJECT_ID")}
                    body:
                      useLegacySql: false
                      query: ${"SELECT COUNT(DISTINCT word) FROM `bigquery-public-data.samples.shakespeare` " + " WHERE corpus='" + corpus + "' "}
                  result: query
              - add:
                  assign:
                    - numWords: ${numWords + int(query.rows[0].f[0].v)}  # first result is the count
    - done:
        return: ${numWords}

JSON

{
  "main": {
    "params": [
      "input"
    ],
    "steps": [
      {
        "init": {
          "assign": [
            {
              "numWords": 0
            },
            {
              "corpuses": [
                "sonnets",
                "various",
                "1kinghenryvi",
                "2kinghenryvi",
                "3kinghenryvi",
                "comedyoferrors",
                "kingrichardiii",
                "titusandronicus",
                "tamingoftheshrew",
                "loveslabourslost"
              ]
            }
          ]
        }
      },
      {
        "runQueries": {
          "parallel": {
            "shared": [
              "numWords"
            ],
            "for": {
              "value": "corpus",
              "in": "${corpuses}",
              "steps": [
                {
                  "runQuery": {
                    "call": "googleapis.bigquery.v2.jobs.query",
                    "args": {
                      "projectId": "${sys.get_env(\"GOOGLE_CLOUD_PROJECT_ID\")}",
                      "body": {
                        "useLegacySql": false,
                        "query": "${\"SELECT COUNT(DISTINCT word) FROM `bigquery-public-data.samples.shakespeare` \" + \" WHERE corpus='\" + corpus + \"' \"}"
                      }
                    },
                    "result": "query"
                  }
                },
                {
                  "add": {
                    "assign": [
                      {
                        "numWords": "${numWords + int(query.rows[0].f[0].v)}"
                      }
                    ]
                  }
                }
              ]
            }
          }
        }
      },
      {
        "done": {
          "return": "${numWords}"
        }
      }
    ]
  }
}

Nächste Schritte