Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
3 changes: 3 additions & 0 deletions CHANGES.md
Original file line number Diff line number Diff line change
Expand Up @@ -101,6 +101,9 @@
* Fixed unbounded checkpoint state growth for splittable DoFns that self-checkpoint on the portable Flink runner (Java) ([#27648](https://github.com/apache/beam/issues/27648)).
* Improved Java pipeline performance by avoiding repeated `DoFn` type descriptor resolution when creating cached invokers ([#39309](https://github.com/apache/beam/issues/39309)).
* (Python) Fixed a memory leak in Python SDK caused by storing exceptions with potentially large stack frames in a cache ([#39406](https://github.com/apache/beam/issues/39406)).
* Portable Java SDK now encodes SchemaCoders in a portable way ([#34672](https://github.com/apache/beam/issues/34672)).
- The new coder translation is protected by the "use_known_schema_coder" PipelineOptions experiment.
- Fixes ([#36496](https://github.com/apache/beam/issues/36496)), ([#30276](https://github.com/apache/beam/issues/30276)), ([#29245](https://github.com/apache/beam/issues/29245)).

## Security Fixes

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -221,7 +221,7 @@ public Boolean visit(OrFinallyTrigger trigger) {
private static byte[] serializeWindowingStrategy(
WindowingStrategy<?, ?> windowingStrategy, PipelineOptions options) {
try {
SdkComponents sdkComponents = SdkComponents.create();
SdkComponents sdkComponents = SdkComponents.create(options);

String workerHarnessContainerImageURL =
DataflowRunner.getContainerImageForJob(options.as(DataflowPipelineOptions.class));
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -1355,19 +1355,20 @@ public DataflowPipelineJob run(Pipeline pipeline) {
// with the SDK harness image (which implements Fn API).
//
// The same Environment is used in different and contradictory ways, depending on whether
// it is a v1 or v2 job submission.
// it is a portable or non-portable job submission.
RunnerApi.Environment defaultEnvironmentForDataflow =
Environments.createDockerEnvironment(workerHarnessContainerImageURL);

// The SdkComponents for portable an non-portable job submission must be kept distinct. Both
// The SdkComponents for portable and non-portable job submission must be kept distinct. Both
// need the default environment.
SdkComponents portableComponents = SdkComponents.create();
portableComponents.registerEnvironment(
defaultEnvironmentForDataflow
.toBuilder()
.addAllDependencies(getDefaultArtifacts())
.addAllCapabilities(Environments.getJavaCapabilities())
.build());
SdkComponents portableComponents =
SdkComponents.create(
options,
defaultEnvironmentForDataflow
.toBuilder()
.addAllDependencies(getDefaultArtifacts())
.addAllCapabilities(Environments.getJavaCapabilities())
.build());

RunnerApi.Pipeline portablePipelineProto =
PipelineTranslation.toProto(pipeline, portableComponents, false);
Expand Down Expand Up @@ -1400,27 +1401,28 @@ public DataflowPipelineJob run(Pipeline pipeline) {
"Skipping Dataflow Streaming Java Runner transform replacements since job will run on Dataflow Portable Runner.");
} else {
// Now rewrite things to be as needed for Dataflow Streaming Java Runner (mutates the
// pipeline)
// pipeline).
// This way the job submitted is valid for Dataflow Streaming Java Runner and Dataflow
// Portable Runner, simultaneously
// Portable Runner, simultaneously.
replaceV1Transforms(pipeline);
}
// Capture the SdkComponents for look up during step translations
SdkComponents dataflowV1Components = SdkComponents.create();
dataflowV1Components.registerEnvironment(
defaultEnvironmentForDataflow
.toBuilder()
.addAllDependencies(getDefaultArtifacts())
.addAllCapabilities(Environments.getJavaCapabilities())
.build());
// Capture the SdkComponents for look up during step translations.
SdkComponents dataflowNonPortableComponents =
SdkComponents.create(
options,
defaultEnvironmentForDataflow
.toBuilder()
.addAllDependencies(getDefaultArtifacts())
.addAllCapabilities(Environments.getJavaCapabilities())
.build());
// No need to perform transform upgrading for the Dataflow Streaming Java Runner proto.
RunnerApi.Pipeline dataflowV1PipelineProto =
PipelineTranslation.toProto(pipeline, dataflowV1Components, true, false);
RunnerApi.Pipeline dataflowNonPortablePipelineProto =
PipelineTranslation.toProto(pipeline, dataflowNonPortableComponents, true, false);

if (LOG.isDebugEnabled()) {
LOG.debug(
"Dataflow v1 pipeline proto:\n{}",
TextFormat.printer().printToString(dataflowV1PipelineProto));
"Dataflow non-portable worker pipeline proto:\n{}",
TextFormat.printer().printToString(dataflowNonPortablePipelineProto));
}

// Set a unique client_request_id in the CreateJob request.
Expand All @@ -1440,7 +1442,11 @@ public DataflowPipelineJob run(Pipeline pipeline) {

JobSpecification jobSpecification =
translator.translate(
pipeline, dataflowV1PipelineProto, dataflowV1Components, this, packages);
pipeline,
dataflowNonPortablePipelineProto,
dataflowNonPortableComponents,
this,
packages);

if (!isNullOrEmpty(dataflowOptions.getDataflowWorkerJar()) && !useUnifiedWorker(options)) {
List<String> experiments =
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -47,6 +47,7 @@
import com.google.api.services.dataflow.model.Job;
import com.google.api.services.dataflow.model.Step;
import com.google.api.services.dataflow.model.WorkerPool;
import com.google.auto.value.AutoValue;
import java.io.File;
import java.io.IOException;
import java.io.Serializable;
Expand Down Expand Up @@ -88,10 +89,13 @@
import org.apache.beam.sdk.io.gcp.bigquery.BigQueryIO;
import org.apache.beam.sdk.io.gcp.bigquery.BigQueryIO.Write.CreateDisposition;
import org.apache.beam.sdk.io.range.OffsetRange;
import org.apache.beam.sdk.options.ExperimentalOptions;
import org.apache.beam.sdk.options.PipelineOptions;
import org.apache.beam.sdk.options.PipelineOptionsFactory;
import org.apache.beam.sdk.options.StreamingOptions;
import org.apache.beam.sdk.options.ValueProvider;
import org.apache.beam.sdk.schemas.AutoValueSchema;
import org.apache.beam.sdk.schemas.annotations.DefaultSchema;
import org.apache.beam.sdk.state.StateSpec;
import org.apache.beam.sdk.state.StateSpecs;
import org.apache.beam.sdk.state.ValueState;
Expand Down Expand Up @@ -166,15 +170,11 @@ public class DataflowPipelineTranslatorTest implements Serializable {
@Rule public transient ExpectedException thrown = ExpectedException.none();

private SdkComponents createSdkComponents(PipelineOptions options) {
SdkComponents sdkComponents = SdkComponents.create();

String containerImageURL =
DataflowRunner.getContainerImageForJob(options.as(DataflowPipelineOptions.class));
RunnerApi.Environment defaultEnvironmentForDataflow =
Environments.createDockerEnvironment(containerImageURL);

sdkComponents.registerEnvironment(defaultEnvironmentForDataflow);
return sdkComponents;
return SdkComponents.create(options, defaultEnvironmentForDataflow);
}

// A Custom Mockito matcher for an initial Job that checks that all
Expand Down Expand Up @@ -1294,15 +1294,16 @@ public String apply(byte[] input) {
file1.deleteOnExit();
File file2 = File.createTempFile("file2-", ".txt");
file2.deleteOnExit();
SdkComponents sdkComponents = SdkComponents.create();
sdkComponents.registerEnvironment(
Environments.createDockerEnvironment(DataflowRunner.getContainerImageForJob(options))
.toBuilder()
.addAllDependencies(
Environments.getArtifacts(
ImmutableList.of("file1.txt=" + file1, "file2.txt=" + file2)))
.addAllCapabilities(Environments.getJavaCapabilities())
.build());
SdkComponents sdkComponents =
SdkComponents.create(
options,
Environments.createDockerEnvironment(DataflowRunner.getContainerImageForJob(options))
.toBuilder()
.addAllDependencies(
Environments.getArtifacts(
ImmutableList.of("file1.txt=" + file1, "file2.txt=" + file2)))
.addAllCapabilities(Environments.getJavaCapabilities())
.build());

RunnerApi.Pipeline pipelineProto = PipelineTranslation.toProto(pipeline, sdkComponents, true);

Expand Down Expand Up @@ -1870,4 +1871,54 @@ public OffsetRange getInitialRange(@SuppressWarnings("unused") @Element String e
return null;
}
}

@AutoValue
@DefaultSchema(AutoValueSchema.class)
public abstract static class SimpleAutoValue {
public abstract String getString();

public abstract int getInt32();

public abstract long getInt64();

public static DataflowPipelineTranslatorTest.SimpleAutoValue of(
String string, int int32, long int64) {
return new AutoValue_DataflowPipelineTranslatorTest_SimpleAutoValue(string, int32, int64);
}
}

@Test
public void testSchemaCoderTranslation() throws Exception {
DataflowPipelineOptions options = buildPipelineOptions();
Pipeline pipeline = Pipeline.create(options);
pipeline
.apply(Impulse.create())
.apply(
MapElements.via(
new SimpleFunction<byte[], SimpleAutoValue>() {
@Override
public SimpleAutoValue apply(byte[] input) {
return SimpleAutoValue.of("foo", 5, 10L);
}
}))
.apply(Window.into(FixedWindows.of(Duration.standardMinutes(1))));
{
// Without the experiment, SchemaCoder is treated as unknown.
SdkComponents sdkComponents = createSdkComponents(options);
RunnerApi.Pipeline pipelineProto = PipelineTranslation.toProto(pipeline, sdkComponents, true);
Map<String, RunnerApi.Coder> coders = pipelineProto.getComponents().getCodersMap();
assertTrue(coders.containsKey("SchemaCoder"));
assertEquals("beam:coders:javasdk:0.1", coders.get("SchemaCoder").getSpec().getUrn());
}

{
// Add the experiment to get the known coder urn instead.
ExperimentalOptions.addExperiment(options, "use_known_schema_coder");
SdkComponents sdkComponents = createSdkComponents(options);
RunnerApi.Pipeline pipelineProto = PipelineTranslation.toProto(pipeline, sdkComponents, true);
Map<String, RunnerApi.Coder> coders = pipelineProto.getComponents().getCodersMap();
assertTrue(coders.containsKey("SchemaCoder"));
assertEquals("beam:coder:schema:v1", coders.get("SchemaCoder").getSpec().getUrn());
}
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -78,7 +78,8 @@ public void testLengthPrefixingOfKeyCoderInStatefulExecutableStage() throws Exce
// Add another stateful stage with a non-standard key coder
Pipeline p = Pipeline.create();
Coder<Void> keycoder = VoidCoder.of();
assertThat(ModelCoderRegistrar.isKnownCoder(keycoder), is(false));
ModelCoderRegistrar coderRegistrar = new ModelCoderRegistrar();
assertThat(coderRegistrar.isKnownCoder(keycoder, p.getOptions()), is(false));
p.apply("impulse", Impulse.create())
.apply(
"create",
Expand Down Expand Up @@ -165,7 +166,8 @@ public void onTimer() {}
public void testLengthPrefixingOfInputCoderExecutableStage() throws Exception {
Pipeline p = Pipeline.create();
Coder<Void> voidCoder = VoidCoder.of();
assertThat(ModelCoderRegistrar.isKnownCoder(voidCoder), is(false));
ModelCoderRegistrar coderRegistrar = new ModelCoderRegistrar();
assertThat(coderRegistrar.isKnownCoder(voidCoder, p.getOptions()), is(false));
p.apply("impulse", Impulse.create())
.apply(
ParDo.of(
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -25,11 +25,13 @@
import org.apache.beam.model.pipeline.v1.RunnerApi;
import org.apache.beam.model.pipeline.v1.RunnerApi.FunctionSpec;
import org.apache.beam.sdk.coders.Coder;
import org.apache.beam.sdk.options.PipelineOptions;
import org.apache.beam.sdk.util.SerializableUtils;
import org.apache.beam.vendor.grpc.v1p69p0.com.google.protobuf.ByteString;
import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.annotations.VisibleForTesting;
import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.BiMap;
import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.ImmutableBiMap;
import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.ImmutableList;
import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.ImmutableMap;
import org.checkerframework.checker.nullness.qual.MonotonicNonNull;
import org.checkerframework.dataflow.qual.Deterministic;
Expand Down Expand Up @@ -62,6 +64,8 @@ private static class DefaultTranslationContext implements TranslationContext {}

private static @MonotonicNonNull BiMap<Class<? extends Coder>, String> knownCoderUrns;

private static @MonotonicNonNull List<CoderTranslatorRegistrar> coderTranslatorRegistrars;

private static @MonotonicNonNull Map<Class<? extends Coder>, CoderTranslator<? extends Coder>>
knownTranslators;

Expand All @@ -80,6 +84,53 @@ static BiMap<Class<? extends Coder>, String> getKnownCoderUrns() {
return knownCoderUrns;
}

private static void initializeCoderTranslatorRegistrars() {
ImmutableList.Builder<CoderTranslatorRegistrar> registrars = ImmutableList.builder();
for (CoderTranslatorRegistrar coderTranslatorRegistrar :
ServiceLoader.load(CoderTranslatorRegistrar.class)) {
registrars.add(coderTranslatorRegistrar);
}
coderTranslatorRegistrars = registrars.build();
}

static boolean isKnownCoder(Coder<?> coder, PipelineOptions options) {
if (coderTranslatorRegistrars == null) {
initializeCoderTranslatorRegistrars();
}
for (CoderTranslatorRegistrar registrar : coderTranslatorRegistrars) {
if (registrar.isKnownCoder(coder, options)) {
return true;
}
}
return false;
}

static CoderTranslator<? extends Coder> getCoderTranslator(Class<? extends Coder> coderClass) {
if (coderTranslatorRegistrars == null) {
initializeCoderTranslatorRegistrars();
}
for (CoderTranslatorRegistrar registrar : coderTranslatorRegistrars) {
CoderTranslator translator = registrar.getCoderTranslator(coderClass);
if (translator != null) {
return translator;
}
}
return null;
}

static Class<? extends Coder> getCoderForUrn(String coderUrn) {
if (coderTranslatorRegistrars == null) {
initializeCoderTranslatorRegistrars();
}
for (CoderTranslatorRegistrar registrar : coderTranslatorRegistrars) {
Class<? extends Coder> coder = registrar.getCoderForUrn(coderUrn);
if (coder != null) {
return coder;
}
}
return null;
}

@VisibleForTesting
@Deterministic
static Map<Class<? extends Coder>, CoderTranslator<? extends Coder>> getKnownTranslators() {
Expand Down Expand Up @@ -107,7 +158,7 @@ public static RunnerApi.MessageWithComponents toProto(Coder<?> coder) throws IOE

public static RunnerApi.Coder toProto(Coder<?> coder, SdkComponents components)
throws IOException {
if (getKnownCoderUrns().containsKey(coder.getClass())) {
if (isKnownCoder(coder, components.getPipelineOptions())) {
return toKnownCoder(coder, components);
}

Expand All @@ -129,7 +180,10 @@ private static RunnerApi.Coder toUnknownCoderWrapper(UnknownCoderWrapper coder)

private static RunnerApi.Coder toKnownCoder(Coder<?> coder, SdkComponents components)
throws IOException {
CoderTranslator translator = getKnownTranslators().get(coder.getClass());
CoderTranslator translator = getCoderTranslator(coder.getClass());
if (translator == null) {
throw new IOException("Unable to find CoderTranslator for known Coder");
}
List<String> componentIds = registerComponents(coder, translator, components);
return RunnerApi.Coder.newBuilder()
.addAllComponentCoderIds(componentIds)
Expand Down Expand Up @@ -186,8 +240,8 @@ private static Coder<?> fromKnownCoder(
components.getComponents().getCodersOrThrow(componentId), components, context);
coderComponents.add(innerCoder);
}
Class<? extends Coder> coderType = getKnownCoderUrns().inverse().get(coderUrn);
CoderTranslator<?> translator = getKnownTranslators().get(coderType);
Class<? extends Coder> coderType = getCoderForUrn(coderUrn);
CoderTranslator<?> translator = getCoderTranslator(coderType);
if (translator != null) {
return translator.fromComponents(
coderComponents, coder.getSpec().getPayload().toByteArray(), context);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -19,6 +19,8 @@

import java.util.Map;
import org.apache.beam.sdk.coders.Coder;
import org.apache.beam.sdk.options.PipelineOptions;
import org.checkerframework.checker.nullness.qual.Nullable;

/** A registrar of {@link Coder} URNs to the associated {@link CoderTranslator}. */
@SuppressWarnings({
Expand All @@ -34,4 +36,18 @@ public interface CoderTranslatorRegistrar {

/** Returns a mapping of URN to {@link CoderTranslator}. */
Map<Class<? extends Coder>, CoderTranslator<? extends Coder>> getCoderTranslators();

/**
* Returns whether the given Coder is known to this CoderTranslatorRegistrar. If the Coder is
* known, then getCoderTranslator() will return a non-null CoderTranslator.
*/
boolean isKnownCoder(Coder<?> coder, PipelineOptions options);

/** Returns the CoderTranslator to use for this Coder, or null if the Coder is not known. */
@Nullable
CoderTranslator<? extends Coder> getCoderTranslator(Class<? extends Coder> coderClass);

/** Returns the Coder to use for the given Urn, or null if the Urn is for an unknown Coder. */
@Nullable
Class<? extends Coder> getCoderForUrn(String coderUrn);
}
Loading
Loading