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.
日本語の概要は準備中です。原文の説明を表示しています。
Explains core Apache Beam programming model concepts including PCollections, PTransforms, Pipelines, and Runners. Use when learning Beam fundamentals or explaining pipeline concepts.
インストール方法を見るインストールする前に、エージェントに与えられる指示の中身を確認できます。
Evolved from Google's MapReduce, FlumeJava, and Millwheel projects. Originally called the "Dataflow Model."
The /model directory defines the official, language-agnostic Protocol Buffer (.proto) and gRPC service specifications that establish the Beam Model and the Beam Portability Framework.
/model Exists (Portability & Decoupling)Without a standardized model representation, supporting $N$ SDK languages across $M$ execution runners would require $N \times M$ separate translation layers. By defining all core pipeline concepts, data encodings, metrics, and worker RPC protocols as Protobuf messages and gRPC services, /model acts as the universal lingua franca:
SDKs compile user pipelines into standardized Runner API protobuf graphs.
Runners inspect, optimize, and distribute these graphs without needing SDK-specific language runtimes.
Workers (SDK Harnesses) execute user code (DoFns) and communicate with runners over standardized Fn API gRPC channels.
/model/pipeline (Runner API & Core Model): Defines the SDK- and runner-independent representation of pipelines (Pipeline, Components, PTransform, PCollection, Coder), timestamps/constants, Beam Schemas (Row, Field), and execution metrics (MonitoringInfo)./model/fn-execution (Fn API & Provisioning): Defines bidirectional gRPC services between runners and worker SDK harnesses for bundle execution (Control), element streaming (Data), state/timer access (State), log forwarding (Logging), and container initialization (Provisioning)./model/job-management (Job, Expansion, & Artifact APIs): Defines gRPC interfaces for submitting and monitoring jobs on remote servers (JobService), resolving cross-language transforms in remote SDKs (ExpansionService), and staging dependency artifacts or container images (ArtifactService)./model/interactive (Interactive API): Defines metadata and stream headers for recording and replaying data in Interactive Beam notebooks.beam:transform:pardo:v1, beam:coder:bytes:v1). When inspecting or creating transforms across languages, always verify URN mappings and registry handlers in both the SDK and Runner runtimes..proto files, as they are used across distributed RPC boundaries and persisted checkpoints./model requires re-generating language bindings (e.g., ./gradlew :model:pipeline:generateProto).class in Java or output in Python, as noted in beam_fn_api.proto comments).A Pipeline encapsulates the entire data processing task, including reading, transforming, and writing data.
// Java
Pipeline p = Pipeline.create(options);
p.apply(...)
.apply(...)
.apply(...);
p.run().waitUntilFinish();
# Python
with beam.Pipeline(options=options) as p:
(p | 'Read' >> beam.io.ReadFromText('input.txt')
| 'Transform' >> beam.Map(process)
| 'Write' >> beam.io.WriteToText('output'))
A distributed dataset that can be bounded (batch) or unbounded (streaming).
A data processing operation that transforms PCollections.
// Java
PCollection<String> output = input.apply(MyTransform.create());
# Python
output = input | 'Name' >> beam.ParDo(MyDoFn())
General-purpose parallel processing.
// Java
input.apply(ParDo.of(new DoFn<String, Integer>() {
@ProcessElement
public void processElement(@Element String element, OutputReceiver<Integer> out) {
out.output(element.length());
}
}));
# Python
class LengthFn(beam.DoFn):
def process(self, element):
yield len(element)
input | beam.ParDo(LengthFn())
# Or simpler:
input | beam.Map(len)
Groups elements by key.
PCollection<KV<String, Integer>> input = ...;
PCollection<KV<String, Iterable<Integer>>> grouped = input.apply(GroupByKey.create());
Joins multiple PCollections by key.
Combines elements (sum, mean, etc.).
// Global combine
input.apply(Combine.globally(Sum.ofIntegers()));
// Per-key combine
input.apply(Combine.perKey(Sum.ofIntegers()));
Merges multiple PCollections.
PCollectionList<String> collections = PCollectionList.of(pc1).and(pc2).and(pc3);
PCollection<String> merged = collections.apply(Flatten.pCollections());
Splits a PCollection into multiple PCollections.
input.apply(Window.into(FixedWindows.of(Duration.standardMinutes(5))));
input | beam.WindowInto(beam.window.FixedWindows(300))
Control when results are emitted.
input.apply(Window.<T>into(FixedWindows.of(Duration.standardMinutes(5)))
.triggering(AfterWatermark.pastEndOfWindow()
.withEarlyFirings(AfterProcessingTime.pastFirstElementInPane()
.plusDelayOf(Duration.standardMinutes(1))))
.withAllowedLateness(Duration.standardHours(1))
.accumulatingFiredPanes());
Additional inputs to ParDo.
PCollectionView<Map<String, String>> sideInput =
lookupTable.apply(View.asMap());
mainInput.apply(ParDo.of(new DoFn<String, String>() {
@ProcessElement
public void processElement(ProcessContext c) {
Map<String, String> lookup = c.sideInput(sideInput);
// Use lookup...
}
}).withSideInputs(sideInput));
Configure pipeline execution.
public interface MyOptions extends PipelineOptions {
@Description("Input file")
@Required
String getInput();
void setInput(String value);
}
MyOptions options = PipelineOptionsFactory.fromArgs(args).as(MyOptions.class);
Strongly-typed access to structured data.
@DefaultSchema(AutoValueSchema.class)
@AutoValue
public abstract class User {
public abstract String getName();
public abstract int getAge();
}
PCollection<User> users = ...;
PCollection<Row> rows = users.apply(Convert.toRows());
TupleTag<String> successTag = new TupleTag<>() {};
TupleTag<String> failureTag = new TupleTag<>() {};
PCollectionTuple results = input.apply(ParDo.of(new DoFn<String, String>() {
@ProcessElement
public void processElement(ProcessContext c) {
try {
c.output(process(c.element()));
} catch (Exception e) {
c.output(failureTag, c.element());
}
}
}).withOutputTags(successTag, TupleTagList.of(failureTag)));
results.get(successTag).apply(WriteToSuccess());
results.get(failureTag).apply(WriteToDeadLetter());
Use transforms from other SDKs.
# Use Java Kafka connector from Python
from apache_beam.io.kafka import ReadFromKafka
result = pipeline | ReadFromKafka(
consumer_config={'bootstrap.servers': 'localhost:9092'},
topics=['my-topic']
)
まだレビューはありません。使ってみた感想をお寄せください。
概要と使いどころ
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.
日本語の概要は準備中です。原文の説明を表示しています。
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.
日本語の概要は準備中です。原文の説明を表示しています。
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.
日本語の概要は準備中です。原文の説明を表示しています。
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.
日本語の概要は準備中です。原文の説明を表示しています。
Guides understanding and using the Gradle build system in Apache Beam. Use when building projects, understanding dependencies, or troubleshooting build issues.
日本語の概要は準備中です。原文の説明を表示しています。
Guides development and usage of I/O connectors in Apache Beam. Use when working with I/O connectors, creating new connectors, or debugging data source/sink issues.
日本語の概要は準備中です。原文の説明を表示しています。