Skip to main content
Version: 3.0 (next)

Google Cloud Dataflow Google Cloud Dataflow Integration Guide

Connect to Google Cloud Dataflow to run Apache Beam stream and batch pipelines from MaestroHub. This guide covers connection setup, function configuration, and pipeline integration.

Overview​

Dataflow is GCP's managed runner for Apache Beam. The connector covers three things a pipeline needs from it:

  • Launching — start a job from a Flex Template or a classic template, with per-launch parameters and worker overrides
  • Following — read a job's state and timings, list the jobs in a region, read its metric snapshot, and read the log messages that say why it failed
  • Stopping — cancel a job outright, or drain a streaming one so buffered work finishes first
  • Flexible authentication — a service-account key, or Application Default Credentials (GKE Workload Identity, GOOGLE_APPLICATION_CREDENTIALS, the GCE metadata server)
Launched, not waited on

Dataflow jobs run for minutes to hours. Launch Flex Template and Launch Classic Template return as soon as Dataflow accepts the request — they do not block until the pipeline finishes. Pair them with Get Job in a later node to read the state, so a two-hour job never holds a MaestroHub pipeline open.

Dataflow is regional

A job launched in us-central1 is invisible from europe-west1. The region is set on the connection, and every function on that connection addresses jobs in it.

Connection Configuration​

Creating a Dataflow Connection​

Go to Connect, click New Connection, and choose Google Cloud Dataflow.

1. Profile Information​

FieldRequiredDescription
Profile NameYesA descriptive name for this connection
DescriptionNoFree text
LabelsNoKey/value labels for organisation

2. Project & Region​

FieldRequiredDescription
Project IDYesThe Google Cloud project that owns the Dataflow jobs
RegionYesThe Dataflow regional endpoint holding the jobs, e.g. us-central1. Pick Custom Region… for a region not in the list

3. Authentication​

FieldRequiredDescription
Service Account JSON KeyNoThe whole service-account key file. Leave empty to use Application Default Credentials

Leaving the key empty is the right choice on GKE with Workload Identity, or anywhere GOOGLE_APPLICATION_CREDENTIALS is set. A pasted value that is not a service-account key is refused at save time, naming what is wrong — a console URL and an authorized_user key are the two common mistakes.

4. Advanced​

FieldRequiredDescription
Custom EndpointNoA different Dataflow API endpoint, for an emulator or a private service endpoint. Leave empty for Google Cloud
Request TimeoutNoDefault bound on a single Dataflow API call (1s–1h, default 30s)
The timeout bounds the API call, not the job

Request Timeout bounds the HTTP call to Dataflow. The job it launches runs on Google Cloud and is not bounded by it — a launch returns in under a second whether the job then runs for one minute or four hours.

IAM Permissions​

The connection test calls dataflow.jobs.list. The roles/dataflow.developer role covers every function in this connector:

OperationPermission
Launch Flex Templatedataflow.jobs.create
Launch Classic Templatedataflow.jobs.create
Get Jobdataflow.jobs.get
List Jobsdataflow.jobs.list
Cancel or Drain Jobdataflow.jobs.update
Get Job Metricsdataflow.jobs.get
List Job Messagesdataflow.messages.list

A launch additionally needs:

  • iam.serviceAccounts.actAs on the worker service account — the project's Compute Engine default service account unless the function overrides it
  • read access to the Cloud Storage bucket holding the template

Testing the Connection​

Test Connection issues a one-item jobs.list against the configured project and region. It proves three things at once: the project exists, the region is enabled for Dataflow, and the credentials carry the role. A project with no jobs is still a successful test.

Function Builder​

Creating Dataflow Functions​

Open a saved connection, go to the Functions tab, and click New Function. The picker shows all seven operations:

Select Google Cloud Dataflow Function Type dialog

Choosing a Dataflow function type

Launch Flex Template Function​

Launches a job from a Flex Template spec in Cloud Storage and returns the job ID.

FieldRequiredDescription
Job NameYesName for the new job. Must be unique among running jobs in the region, start with a lowercase letter, and contain only lowercase letters, digits and hyphens. Templatable
Template Spec PathYesCloud Storage path of the Flex Template spec file, e.g. gs://my-bucket/templates/aggregate.json. Templatable
Template Parameters (JSON object)NoJSON object of template parameters. Dataflow requires string values, so numbers and booleans must be quoted. Templatable
Machine TypeNoWorker machine type for this launch, e.g. n1-standard-4
Max WorkersNoAutoscaling ceiling for this launch (1–1000)
Initial WorkersNoStarting worker count (1–1000). Cannot exceed Max Workers
Worker Service AccountNoThe identity the workers run as. Empty uses the project's Compute Engine default
Temp LocationNoCloud Storage path for temporary files
SubnetworkNoregions/REGION/subnetworks/SUBNETWORK. Empty uses the default network
TimeoutNoBound on the Launch API call itself

Result:

{
"jobId": "2026-09-24_03_15_22-1234567890123456789",
"jobName": "nightly-aggregation",
"jobType": "JOB_TYPE_BATCH",
"region": "us-central1",
"createTime": "2026-09-24T03:15:22Z"
}

jobName is the name Dataflow accepted, which can differ from the one asked for when Dataflow disambiguates. Dataflow reports no state on a launch, so the job's first state — and every one after it — come from Get Job.

Template parameters are strings

{"maxRows": 5000} is rejected: Dataflow's template parameters are a string→string map, so write {"maxRows": "5000"}. The error names the offending key.

Paths are gs://, not console URLs

A https://console.cloud.google.com/storage/… link is a browser URL, not a Cloud Storage path. Both launch functions refuse it at save time; Dataflow itself would report only a generic "unable to read template".

Launch Classic Template Function​

Launches a job from a classic template staged in Cloud Storage — the shape every Google-provided template under gs://dataflow-templates uses.

FieldRequiredDescription
Job NameYesAs above. Templatable
Template PathYesCloud Storage path of the staged classic template, e.g. gs://dataflow-templates/latest/PubSub_to_BigQuery. Templatable
Template Parameters (JSON object)NoAs above. Templatable
Machine Type, Max Workers, Initial Workers, Worker Service Account, Temp Location, SubnetworkNoAs above. Classic templates usually require Temp Location
TimeoutNoBound on the Launch API call itself

The result shape is identical to the Flex launch.

Get Job Function​

Reads one job by ID. This is the polling half of a launch.

FieldRequiredDescription
Job IDYesThe job ID returned by a launch. Templatable
TimeoutNoBound on this single operation

Result:

{
"jobId": "2026-09-24_03_15_22-1234567890123456789",
"jobName": "nightly-aggregation",
"currentState": "JOB_STATE_DONE",
"currentStateTime": "2026-09-24T03:41:08Z",
"requestedState": "",
"jobType": "JOB_TYPE_BATCH",
"region": "us-central1",
"createTime": "2026-09-24T03:15:22Z",
"startTime": "2026-09-24T03:16:40Z",
"sdkVersion": "2.58.0",
"labels": {"env": "prod"}
}

Every key is always present. startTime is empty while the job is still queued — that is how a polling loop tells a queued job from a running one without parsing the state string. requestedState is empty on a job nobody has asked to stop.

The job's template path is not returned: it lives in the job's Beam pipeline options, which the summary view Dataflow serves here does not include. Fetching the full view to recover one string would ship the whole serialized pipeline graph on every poll. For a Google-provided template the job's labels name it instead.

List Jobs Function​

Returns the jobs in the configured region.

FieldRequiredDescription
FilterNoALL, ACTIVE (running and starting) or TERMINATED (done, failed, cancelled, drained). Applied by Dataflow
Name FilterNoCase-sensitive substring the job name must contain. Templatable
Max ItemsNoMaximum jobs to return (1–1000, default 100)
TimeoutNoBound on this single operation

Result:

{
"jobs": [
{
"jobId": "2026-09-24_03_15_22-1234567890123456789",
"jobName": "nightly-aggregation",
"currentState": "JOB_STATE_DONE",
"currentStateTime": "2026-09-24T03:41:08Z",
"jobType": "JOB_TYPE_BATCH",
"region": "us-central1",
"createTime": "2026-09-24T03:15:22Z"
}
],
"jobCount": 1,
"truncated": false
}

jobs is always a list — empty rather than absent — so a ForEach over a region with no jobs does nothing rather than failing on a missing source.

Filter is server-side; Name Filter is not

Dataflow applies Filter itself. It offers no server-side name filter, so Name Filter narrows what the page returned. The connector reads full pages while a name filter is set, so a match further down the region is still found.

Cancel or Drain Job Function​

Asks Dataflow to stop a running job.

FieldRequiredDescription
Job IDYesThe job to stop. Templatable
Requested StateYesJOB_STATE_CANCELLED stops at once and discards in-flight data. JOB_STATE_DRAINED stops ingestion and lets buffered work finish — streaming jobs only
TimeoutNoBound on the Update API call itself

Result:

{
"jobId": "2026-09-24_03_15_22-1234567890123456789",
"jobType": "JOB_TYPE_BATCH",
"requestedState": "JOB_STATE_CANCELLED"
}
This tells you what was asked, not what happened

Dataflow answers a state change with almost nothing — it reports no state of its own. The result carries the job, its type, and the state that was requested. Use Get Job to watch the job reach JOB_STATE_CANCELLING / JOB_STATE_DRAINING and then its terminal state.

Draining a batch job is refused by Dataflow — only a streaming job can be drained.

Get Job Metrics Function​

Returns the metric values Dataflow recorded for a job.

FieldRequiredDescription
Job IDYesThe job whose metrics to read. Templatable
SinceNoRFC 3339 timestamp; only metrics that changed after this instant are returned. Templatable
Max ItemsNoMaximum metric entries (1–5000, default 500)
TimeoutNoBound on this single operation

Result:

{
"metrics": [
{
"name": "ElementCount",
"origin": "dataflow/v1b3",
"step": "WriteToBigQuery",
"kind": "sum",
"value": 184320,
"cumulative": true,
"updateTime": "2026-09-24T03:41:02Z"
}
],
"metricCount": 1,
"metricTime": "2026-09-24T03:41:05Z",
"truncated": false
}

value is whichever of Dataflow's typed value fields this metric's kind populated, so a downstream expression reads one key rather than probing six. Counters the Beam pipeline publishes itself appear with origin: "user".

List Job Messages Function​

Returns the messages Dataflow logged for a job.

FieldRequiredDescription
Job IDYesThe job whose messages to read. Templatable
Minimum ImportanceNoJOB_MESSAGE_DEBUG, _DETAILED, _BASIC (default), _WARNING or _ERROR
From / UntilNoRFC 3339 window bounds. Templatable
Max ItemsNoMaximum messages (1–1000, default 100)
TimeoutNoBound on this single operation

Result:

{
"messages": [
{
"id": "msg-8a1f3c5e",
"importance": "JOB_MESSAGE_ERROR",
"text": "Workflow failed. Causes: The template parameters were invalid.",
"time": "2026-09-24T03:40:11Z"
}
],
"messageCount": 1,
"truncated": false
}
This is where a failure says why

A failed job records the cause here and nowhere else in the Dataflow API. JOB_STATE_FAILED from Get Job tells you that it failed; JOB_MESSAGE_ERROR from this function tells you why. An alert that carries both is the one an operator can act on.

Using Parameters​

Any field marked templatable accepts ((parameterName)). Parameters are detected as you type and listed in the Function Parameters card:

Auto-detected Dataflow function parameters

A parameter detected from ((day)) in the template parameters

A pipeline supplies their values at execution time, so one Launch function can serve every day of a backfill, and one Get Job function can follow whichever job a launch produced.

Pipeline Integration​

Each function becomes a node in the pipeline editor under Databases. See the Google Cloud Dataflow node reference for node-level configuration.

The shape that comes up most often is launch → wait → read:

  1. Launch Flex Template starts the job and emits jobId
  2. A Delay node waits
  3. Get Job reads the state using $node["Launch Flex Template"].result.jobId
  4. A Condition node branches on currentState — loop back to the delay while it is JOB_STATE_RUNNING, continue on JOB_STATE_DONE, and on JOB_STATE_FAILED run List Job Messages to collect the cause

Common Use Cases​

Nightly aggregation with a per-day parameter​

Launch Flex Template with {"inputDate": "((day))"} runs the same template for whichever day the trigger supplies. The job ID goes into an execution log; Get Job reads the outcome on the next pass.

Blue/green redeploy of a streaming job​

Cancel or Drain Job with JOB_STATE_DRAINED stops the old streaming job without losing buffered work, Get Job waits for it to reach JOB_STATE_DRAINED, and Launch Classic Template starts the replacement.

Failure alerts that carry the cause​

A pipeline polls Get Job; on JOB_STATE_FAILED it calls List Job Messages with Minimum Importance JOB_MESSAGE_ERROR and posts the text to Slack, so the alert names the failure instead of a bare state.

Asserting a load before trusting it​

Get Job Metrics reads the ElementCount for the write step; a Condition node compares it against the expected row count before a downstream node consumes the output.

Troubleshooting​

SymptomCauseFix
Dataflow found nothing at …Wrong project ID, region, or job IDUse List Jobs on the same connection to see what exists in that region
the credentials cannot use Dataflow in project … region …The service account lacks the Dataflow roleGrant roles/dataflow.developer
no Google credentials available for …No key pasted and ADC does not resolve herePaste a service-account key, or run where Workload Identity / GOOGLE_APPLICATION_CREDENTIALS is available
Service account key must be valid JSONA console URL or a partial file was pastedPaste the whole key file
This is a "authorized_user" key, not a service-account keyA gcloud user credential was pastedDownload a key for a service account instead
Job name must start with a lowercase letter…Uppercase or an invalid character in the job nameUse lowercase letters, digits and hyphens
Template spec path must be a Cloud Storage path beginning with gs://A console URL was pastedUse the gs:// path
Value of "x" must be a stringA template parameter is a number or booleanQuote it: {"x": "5000"}
Initial Workers cannot exceed Max WorkersThe worker override pair is inconsistentLower Initial Workers, or raise Max Workers
Draining is only supported for streaming jobsDrain was requested on a batch jobUse JOB_STATE_CANCELLED
requestedState must be JOB_STATE_CANCELLED or JOB_STATE_DRAINEDAnother state reached the APIDataflow treats other states as a silent no-op; the connector refuses them instead
Listing looks shortThe Max Items budget stopped itRaise Max Items; truncated in the result says the cut happened
A launched job has no template parametersThe parameters JSON was empty or its placeholders were unresolvedSupply values for every ((name)) the field uses — an unresolved placeholder resolves to an empty string