New to Kestra?
Use blueprints to kickstart your first workflows.
Orchestrate the Apache Beam word-count pipeline on Kestra with the DIRECT runner, real ingested text, retries, timeouts, and Slack failure alerts.
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.
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.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.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.log_result task (io.kestra.plugin.core.log.Log) prints the captured output shard URIs.retry and timeout on download_text and run_wordcount, plus pipelineTimeoutSeconds for the Beam wait.:latest.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.
RunPipeline runs Beam inside the pinned SDK container.download_text can fetch the sample text.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.source_url with another public text file.run_wordcount task run and review the captured outputFiles.counts CSV shard from the Outputs tab to see per-word totals.beamRunner: DIRECT for a clustered runner (Flink, Spark, or Dataflow) once the pipeline works locally.Download source with an object-storage or database ingestion task and feed it through inputFiles.Schedule or event trigger to run the word count on every new file or on a cron.