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.
日本語の概要は準備中です。原文の説明を表示しています。
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.
インストール方法を見るインストールする前に、エージェントに与えられる指示の中身を確認できます。
This guide outlines the modern best practices and mandatory steps for building new Apache Beam I/O connectors. Modern Beam connectors are expected to be schema-aware, available in cross-language (Python, Go) pipelines via the Expansion Service, and seamlessly integrable with Beam YAML and the Managed I/O API.
New Java I/O connectors should reside under sdks/java/io/<connector-name>.
sdks/java/io/<connector-name>/
├── build.gradle
└── src/
├── main/java/org/apache/beam/sdk/io/<connector-name>/
│ ├── <Connector>IO.java
│ ├── <Connector>ReadSchemaTransformProvider.java
│ └── <Connector>WriteSchemaTransformProvider.java
└── test/java/org/apache/beam/sdk/io/<connector-name>/
├── <Connector>IOTest.java
└── <Connector>ReadSchemaTransformProviderTest.java
build.gradle)Your build.gradle must use standard Beam Java module conventions and explicitly declare necessary dependencies.
plugins { id 'org.apache.beam.module' }
applyJavaNature(
// If <connector-name> contains hyphens, convert them to dots or underscores
automaticModuleName: 'org.apache.beam.sdk.io.<connector_name>',
)
description = "Apache Beam :: SDKs :: Java :: IO :: <Connector Name>"
ext.summary = "Integration with <External System>."
dependencies {
implementation project(path: ":sdks:java:core", configuration: "shadow")
implementation project(path: ":model:pipeline", configuration: "shadow") // For URN definitions
// Add external client libraries here
implementation library.java.<client_dependency>
// Handle strict dependency checking if necessary
permitUnusedDeclared library.java.<client_dependency>
// Standard test dependencies
testImplementation project(path: ":sdks:java:core", configuration: "shadowTest")
testImplementation library.java.junit
}
<Connector>IO.java)Follow Beam's canonical AutoValue builder pattern for user-facing API configuration. While core Java I/O connectors can be strongly typed using specific domain classes or Java generics (<T>) for idiomatic Java SDK usage, modern sources should also emphasize Beam Row and Schema support (e.g., via .readRows()). For excellent real-world implementations of this pattern, refer to IcebergIO and DeltaIO.
Instead of legacy Source classes, implement reading via Beam's Splittable DoFns (SDF) framework for advanced features such as dynamic rebalancing and watermark support.
A primary read transform (such as read() or readRows()) typically extends PTransform<PBegin, PCollection<T>> (or PCollection<Row>). Using PCollection as input is meant for "ReadAll" operations (such as reading a collection of file patterns or queries).
An example SDF-based read transform is given below:
public class MyIO {
public static ReadRows readRows() {
return new AutoValue_MyIO_ReadRows.Builder().build();
}
@AutoValue
public abstract static class ReadRows extends PTransform<PBegin, PCollection<Row>> {
public abstract @Nullable String getConfigurationOption();
public abstract @Nullable Schema getSchema();
public abstract Builder toBuilder();
@AutoValue.Builder
public abstract static class Builder {
public abstract Builder setConfigurationOption(String value);
public abstract Builder setSchema(Schema schema);
public abstract ReadRows build();
}
public ReadRows withConfigurationOption(String value) {
return toBuilder().setConfigurationOption(value).build();
}
public ReadRows withSchema(Schema schema) {
return toBuilder().setSchema(schema).build();
}
@Override
public PCollection<Row> expand(PBegin input) {
return input
// `ReaderDoFn` is an SDF or source implementation that reads records and outputs `Row` objects.
.apply(ParDo.of(new ReaderDoFn(getConfigurationOption())))
.setRowSchema(getSchema());
}
}
public static WriteRows writeRows() {
return new AutoValue_MyIO_WriteRows.Builder().build();
}
@AutoValue
public abstract static class WriteRows extends PTransform<PCollection<Row>, PDone> {
public abstract @Nullable String getConfigurationOption();
public abstract Builder toBuilder();
@AutoValue.Builder
public abstract static class Builder {
public abstract Builder setConfigurationOption(String value);
public abstract WriteRows build();
}
public WriteRows withConfigurationOption(String value) {
return toBuilder().setConfigurationOption(value).build();
}
@Override
public PDone expand(PCollection<Row> input) {
input.apply("WriteRecords", ParDo.of(new WriterDoFn(getConfigurationOption())));
return PDone.in(input.getPipeline());
}
}
}
Please make sure that any code you add adheres to the Beam coding standards. These standards are documented here: https://beam.apache.org/contribute/code-guidelines/
Especially scrutinize any logic that involves splitting the data for parallel processing since it is a common source of errors that can lead to data loss or data duplication related issues.
When implementing data egress (Write transforms), avoid creating single-worker bottlenecks. Depending on your target system's transactional requirements, prefer one of the following canonical Beam sink patterns:
DoFn that manages connections per bundle (@StartBundle, @FinishBundle) or utilizes GroupIntoBatches to perform highly efficient, parallel batched requests.PTransform:
PCollection<CommitMessage>).FileIO core infrastructure (FileIO.write() / FileIO.Sink) rather than implementing custom file rolling and sharding logic.PCollection<Row> (or custom WriteResult / PCollectionRowTuple) containing failed records and explicit error metadata.external_transforms.proto)To standardize your transform identifier across SDKs, define its URN in Beam's protobuf schema.
model/pipeline/src/main/proto/org/apache/beam/model/pipeline/v1/external_transforms.proto.ManagedTransforms.Urns).message ManagedTransforms {
enum Urns {
// ... existing entries
MY_SYSTEM_READ = 15 [(org.apache.beam.model.pipeline.v1.beam_urn) =
"beam:schematransform:org.apache.beam:my_system_read:v1"];
MY_SYSTEM_WRITE = 16 [(org.apache.beam.model.pipeline.v1.beam_urn) =
"beam:schematransform:org.apache.beam:my_system_write:v1"];
}
}
./gradlew :model:pipeline:generateProto :model:pipeline:compileJava
SchemaTransformProviderTo expose your connector to cross-language pipelines and Beam YAML, create a typed SchemaTransformProvider.
@AutoService(SchemaTransformProvider.class)
public class MyReadSchemaTransformProvider extends TypedSchemaTransformProvider<Configuration> {
@Override
public String identifier() {
return getUrn(ExternalTransforms.ManagedTransforms.Urns.MY_SYSTEM_READ);
}
@Override
public String description() {
return "Reads records from My System and outputs a PCollection of Beam Rows.";
}
@Override
protected SchemaTransform from(Configuration configuration) {
return new MyReadSchemaTransform(configuration);
}
@Override
public List<String> outputCollectionNames() {
return Collections.singletonList("output");
}
@DefaultSchema(AutoValueSchema.class)
@AutoValue
public abstract static class Configuration {
@SchemaFieldDescription("Configuration option description.")
public abstract String getConfigurationOption();
public static Builder builder() {
return new AutoValue_MyReadSchemaTransformProvider_Configuration.Builder();
}
@AutoValue.Builder
public abstract static class Builder {
public abstract Builder setConfigurationOption(String value);
public abstract Configuration build();
}
}
static class MyReadSchemaTransform extends SchemaTransform {
private final Configuration configuration;
MyReadSchemaTransform(Configuration configuration) {
this.configuration = Objects.requireNonNull(configuration, "configuration cannot be null");
}
@Override
public PCollectionRowTuple expand(PCollectionRowTuple input) {
PCollection<Row> output = input.getPipeline().apply(
MyIO.readRows().withConfigurationOption(configuration.getConfigurationOption()));
return PCollectionRowTuple.of("output", output);
}
}
}
Managed.java)Beam's Managed I/O transform provides a unified interface for data ingest/egress. To support it:
sdks/java/managed/src/main/java/org/apache/beam/sdk/managed/Managed.java.public static final String MY_SYSTEM = "my_system";
READ_TRANSFORMS or WRITE_TRANSFORMS:public static final Map<String, String> READ_TRANSFORMS =
ImmutableMap.<String, String>builder()
// ... existing transforms
.put(MY_SYSTEM, getUrn(ExternalTransforms.ManagedTransforms.Urns.MY_SYSTEM_READ))
.build();
Managed.java to list your new connector.To enable non-Java SDKs (Python, Go) to discover and expand your new connector, include it in the standard Java Expansion Service.
sdks/java/io/expansion-service/build.gradle.dependencies {
// ... existing dependencies
runtimeOnly project(":sdks:java:io:<connector-name>")
}
Once registered in the expansion service, your SchemaTransform can be utilized in Python and YAML. E.g., for Managed support in Python:
sdks/python/apache_beam/transforms/managed.py.__all__ and map it to its URN in Read._READ_TRANSFORMS or Write._WRITE_TRANSFORMS:MY_SYSTEM = 'my_system'
__all__ = [
# ... existing
"MY_SYSTEM",
]
class Read(PTransform):
_READ_TRANSFORMS = {
# ... existing
MY_SYSTEM: ManagedTransforms.Urns.MY_SYSTEM_READ.urn,
}
sdks/python/apache_beam/transforms/external.py, map the URN to the appropriate Expansion Service jar target in MANAGED_TRANSFORM_URN_TO_JAR_TARGET_MAPPING.Verify your new connector thoroughly across multiple abstraction layers:
Test your core builder and SchemaTransformProvider translation.
./gradlew :sdks:java:io:<connector-name>:test
In ManagedSchemaTransformTranslationTest.java (under sdks/java/managed), you can verify the translation structure of your managed transform. Note that ManagedTest.java generally uses dummy/test providers (TestSchemaTransformProvider) to keep dependencies lightweight.
./gradlew :sdks:java:managed:test
Create integration tests to test end-to-end data processing against real system instances, including Managed.read(Managed.MY_SYSTEM) usage.
ResourceManager classes under the it/ directory.Add any necessary documentation for your connector under the website/www/site/content/en/documentation/io/built-in/ directory.
[!TIP] Canonical Reference Implementations: When developing a new connector, we highly recommend studying Apache Iceberg (IcebergIO.java) and Delta Lake (DeltaIO.java) as state-of-the-art reference implementations.
For more details see the Developing I/O connectors guide.
まだレビューはありません。使ってみた感想をお寄せください。
概要と使いどころ
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.
日本語の概要は準備中です。原文の説明を表示しています。
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.
日本語の概要は準備中です。原文の説明を表示しています。
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.
日本語の概要は準備中です。原文の説明を表示しています。