本文へ移動
cccskills
無料GitHub で公開

runners

Guides understanding and working with Apache Beam runners (Direct, Dataflow, Flink, Spark, etc.). Use when configuring pipelines for different execution environments or debugging runner-specific issues.

インストール方法を見る

含まれるファイル(1)

  • SKILL.md6.0 KB

SKILL.md(原文)

インストールする前に、エージェントに与えられる指示の中身を確認できます。

Apache Beam Runners

Overview

Runners execute Beam pipelines on distributed processing backends. Each runner translates the portable Beam model to its native execution engine.

Available Runners

RunnerLocationDescription
Directrunners/direct-java/Local execution for testing
Prismrunners/prism/Portable local runner
Dataflowrunners/google-cloud-dataflow-java/Google Cloud Dataflow
Flinkrunners/flink/Apache Flink
Sparkrunners/spark/Apache Spark
Jetrunners/jet/Hazelcast Jet
Twister2runners/twister2/Twister2

Direct Runner

For local development and testing.

Java

PipelineOptions options = PipelineOptionsFactory.create();
options.setRunner(DirectRunner.class);
Pipeline p = Pipeline.create(options);

Python

options = PipelineOptions()
options.view_as(StandardOptions).runner = 'DirectRunner'
p = beam.Pipeline(options=options)

Command Line

--runner=DirectRunner

Dataflow Runner

Prerequisites

  • GCP project with Dataflow API enabled
  • Service account with Dataflow Admin role
  • GCS bucket for staging

Java Usage

DataflowPipelineOptions options = PipelineOptionsFactory.as(DataflowPipelineOptions.class);
options.setRunner(DataflowRunner.class);
options.setProject("my-project");
options.setRegion("us-central1");
options.setTempLocation("gs://my-bucket/temp");

Python Usage

options = PipelineOptions([
    '--runner=DataflowRunner',
    '--project=my-project',
    '--region=us-central1',
    '--temp_location=gs://my-bucket/temp'
])

Runner v2

--experiments=use_runner_v2

Custom SDK Container

--sdkContainerImage=gcr.io/project/beam_java11_sdk:custom

Flink Runner

Embedded Mode

FlinkPipelineOptions options = PipelineOptionsFactory.as(FlinkPipelineOptions.class);
options.setRunner(FlinkRunner.class);
options.setFlinkMaster("[local]");

Cluster Mode

options.setFlinkMaster("host:port");

Portable Mode (Python)

options = PipelineOptions([
    '--runner=FlinkRunner',
    '--flink_master=host:port',
    '--environment_type=LOOPBACK'  # or DOCKER, EXTERNAL
])

Spark Runner

Java

SparkPipelineOptions options = PipelineOptionsFactory.as(SparkPipelineOptions.class);
options.setRunner(SparkRunner.class);
options.setSparkMaster("local[*]");  # or spark://host:port

Python (Portable)

options = PipelineOptions([
    '--runner=SparkRunner',
    '--spark_master_url=local[*]'
])

Testing with Runners

ValidatesRunner Tests

Tests that validate runner correctness:

# Direct Runner
./gradlew :runners:direct-java:validatesRunner

# Flink Runner
./gradlew :runners:flink:1.18:validatesRunner

# Spark Runner
./gradlew :runners:spark:3:validatesRunner

# Dataflow Runner
./gradlew :runners:google-cloud-dataflow-java:validatesRunner

TestPipeline with Runners

@Rule public TestPipeline pipeline = TestPipeline.create();

// Set runner via system property
-DbeamTestPipelineOptions='["--runner=TestDataflowRunner"]'

Portable Runners

Concept

  • SDK-independent execution via Fn API
  • SDK runs in container, communicates via gRPC

Environment Types

  • DOCKER - SDK in Docker container
  • LOOPBACK - SDK in same process (testing)
  • EXTERNAL - SDK at specified address
  • PROCESS - SDK in subprocess

Job Server

Start Flink job server:

./gradlew :runners:flink:1.18:job-server:runShadow

Start Spark job server:

./gradlew :runners:spark:3:job-server:runShadow

Runner-Specific Options

Dataflow

OptionDescription
--projectGCP project
--regionGCP region
--tempLocationGCS temp location
--stagingLocationGCS staging
--numWorkersInitial workers
--maxNumWorkersMax workers
--workerMachineTypeVM type

Flink

OptionDescription
--flinkMasterFlink master address
--parallelismDefault parallelism
--checkpointingIntervalCheckpoint interval

Spark

OptionDescription
--sparkMasterSpark master URL
--sparkConfAdditional Spark config

Building Runner Artifacts

Dataflow Worker Jar

./gradlew :runners:google-cloud-dataflow-java:worker:shadowJar

Flink Job Server

./gradlew :runners:flink:1.18:job-server:shadowJar

Spark Job Server

./gradlew :runners:spark:3:job-server:shadowJar

Debugging

Direct Runner

  • Enable logging: -Dorg.slf4j.simpleLogger.defaultLogLevel=debug
  • Use --targetParallelism=1 for deterministic execution

Dataflow

  • Check Dataflow UI: console.cloud.google.com/dataflow
  • Use --experiments=upload_graph for graph debugging
  • Worker logs in Cloud Logging

Portable Runners

  • Enable debug logging on job server
  • Check SDK harness logs in worker containers

レビュー

まだレビューはありません。使ってみた感想をお寄せください。

同じリポジトリのスキル

概要と使いどころ

Guide on how to add and propagate new metadata fields in Apache Beam's WindowedValue, extending protos, windmill persistence, and runner interfaces to avoid metadata loss.

日本語の概要は準備中です。原文の説明を表示しています。

apache/beam8,6802026年10月10日 更新

Explains core Apache Beam programming model concepts including PCollections, PTransforms, Pipelines, and Runners. Use when learning Beam fundamentals or explaining pipeline concepts.

日本語の概要は準備中です。原文の説明を表示しています。

apache/beam8,6802026年10月10日 更新

ci-cd

無料

Guides understanding and working with Apache Beam's CI/CD system using GitHub Actions. Use when debugging CI failures, understanding test workflows, or modifying CI configuration.

日本語の概要は準備中です。原文の説明を表示しています。

apache/beam8,6802026年10月10日 更新

Guides the contribution workflow for Apache Beam, including creating PRs, issue management, code review process, release cycles, and rigorous evaluation rules for high-risk core component changes. Use when contributing code, creating PRs, or modifying core Beam components.

日本語の概要は準備中です。原文の説明を表示しています。

apache/beam8,6802026年10月10日 更新

End-to-end guide on developing new Apache Beam I/O connectors correctly, including core IO transforms, SchemaTransforms, URN proto definitions, Managed API integration, cross-language expansion service, and testing.

日本語の概要は準備中です。原文の説明を表示しています。

apache/beam8,6802026年10月10日 更新

Guides understanding and using the Gradle build system in Apache Beam. Use when building projects, understanding dependencies, or troubleshooting build issues.

日本語の概要は準備中です。原文の説明を表示しています。

apache/beam8,6802026年10月10日 更新

apache のスキルをすべて見る

このスキルの問題を報告する