Download icon
Write icon
RunPipeline icon
Log icon
SlackIncomingWebhook icon

Apache Beam Word Count with the DIRECT Runner

Orchestrate the Apache Beam word-count pipeline on Kestra with the DIRECT runner, real ingested text, retries, timeouts, and Slack failure alerts.

Categories
Data

Run the canonical Apache Beam word-count pipeline on Kestra using the Beam DIRECT runner over a real downloaded text file instead of a toy literal. This blueprint is the hello-world of Apache Beam data processing, showing how Kestra internal storage feeds a Beam source, how a Python pipeline tokenizes and aggregates words, and how output shards flow back into Kestra. It solves the problem of learning and testing Beam pipelines without standing up a Flink, Spark, or Dataflow cluster, because the DIRECT runner executes locally inside the task container.

How it works

  1. The download_text task (io.kestra.plugin.core.http.Download) fetches the source_url text file into Kestra internal storage with retry and a timeout, so a transient network blip does not kill the run.
  2. The write_pipeline task (io.kestra.plugin.core.storage.Write) renders a Beam YAML pipeline definition and stores it. Keeping it as a Write (rather than an inline definition) lets the source file name source.txt live in one place, shared with inputFiles.
  3. The run_wordcount task (io.kestra.plugin.beam.RunPipeline) runs the pipeline with sdk: PYTHON and beamRunner: DIRECT on the pinned apache/beam_python3.11_sdk:2.71.0 image. The downloaded file is mounted via inputFiles, the pipeline reads each line, tokenizes words, explodes them, maps each word to a count of 1, sums counts per word, and writes CSV via WriteToCsv. Shards are captured by outputFiles.
  4. The log_result task (io.kestra.plugin.core.log.Log) prints the captured output shard URIs.

What you get

  • A working, end-to-end Beam pipeline you can run in minutes with no external cluster.
  • Per-word counts written as CSV shards, captured into Kestra internal storage.
  • Built-in reliability: task-level retry and timeout on download_text and run_wordcount, plus pipelineTimeoutSeconds for the Beam wait.
  • A reproducible run thanks to the pinned Beam SDK image rather than :latest.

Who it's for

  • Data engineers evaluating Apache Beam or prototyping pipelines before deploying to Dataflow or Flink.
  • Platform teams who want a reference pattern for running Beam under an orchestrator.
  • Anyone learning the Beam YAML framework and the DIRECT runner.

Why orchestrate this with Kestra

The Beam DIRECT runner has no scheduler, no retry policy, and no alerting of its own: it just executes a pipeline once. Kestra fills that gap. You declare the whole flow in YAML, attach event or schedule triggers, get task-level retries and timeouts, capture lineage from the ingested file through the generated pipeline to the output shards, and wire a flow-level errors block that posts a Slack alert on failure. This is the production scaffolding Beam itself leaves to you.

Prerequisites

  • Docker available to the Kestra worker, since RunPipeline runs Beam inside the pinned SDK container.
  • Outbound HTTPS so download_text can fetch the sample text.
  • No Flink, Spark, or JobServer is required for the DIRECT runner.

Secrets

  • SLACK_WEBHOOK: incoming webhook used only by the failure-alert errors block (io.kestra.plugin.slack.notifications.SlackIncomingWebhook). The word-count itself uses no credentials. Remove the errors block if you do not want Slack alerting.

Quick start

  1. Add the flow and execute it. Optionally override source_url with another public text file.
  2. Open the run_wordcount task run and review the captured outputFiles.
  3. Download a counts CSV shard from the Outputs tab to see per-word totals.

How to extend

  • Swap beamRunner: DIRECT for a clustered runner (Flink, Spark, or Dataflow) once the pipeline works locally.
  • Replace the Download source with an object-storage or database ingestion task and feed it through inputFiles.
  • Add a Schedule or event trigger to run the word count on every new file or on a cron.
  • Extend the Beam YAML transforms to filter stop words, lowercase, or write to Parquet instead of CSV.

Links

See How

New to Kestra?

Use blueprints to kickstart your first workflows.