diff --git a/.github/workflows/codeql.yml b/.github/workflows/codeql.yml index 30539e8ec08a..95d128d53658 100644 --- a/.github/workflows/codeql.yml +++ b/.github/workflows/codeql.yml @@ -74,7 +74,7 @@ jobs: # your codebase is analyzed, see https://docs.github.com/en/code-security/code-scanning/creating-an-advanced-setup-for-code-scanning/codeql-code-scanning-for-compiled-languages steps: - name: Checkout repository - uses: actions/checkout@v4 + uses: actions/checkout@v6 # Add any setup steps before running the `github/codeql-action/init` action. # This includes steps like installing compilers or runtimes (`actions/setup-node` diff --git a/CHANGES.md b/CHANGES.md index ac8215f35494..698d88b01fab 100644 --- a/CHANGES.md +++ b/CHANGES.md @@ -69,7 +69,6 @@ ## New Features / Improvements -* X feature added (Java/Python) ([#X](https://github.com/apache/beam/issues/X)). * (Java) Enabled state tag encoding v2 by default for new Dataflow Streaming Engine jobs. It can be disabled by passing `--experiments=disable_streaming_engine_state_tag_encoding_v2` or `--updateCompatibilityVersion=2.74.0` pipeline option. Note that the tag encoding version cannot change during a job update. Jobs using tag encoding v2 (enabled by default for new jobs on 2.75.0+) cannot be downgraded to Beam versions prior to 2.73.0, as only versions 2.73.0 and later support tag encoding v2. ([#38705](https://github.com/apache/beam/issues/38705)). ## Breaking Changes diff --git a/sdks/go.mod b/sdks/go.mod index e110296eb540..bd415ede1b52 100644 --- a/sdks/go.mod +++ b/sdks/go.mod @@ -33,9 +33,9 @@ require ( cloud.google.com/go/spanner v1.91.0 cloud.google.com/go/storage v1.62.3 github.com/aws/aws-sdk-go-v2 v1.42.0 - github.com/aws/aws-sdk-go-v2/config v1.32.23 + github.com/aws/aws-sdk-go-v2/config v1.32.24 github.com/aws/aws-sdk-go-v2/credentials v1.19.23 - github.com/aws/aws-sdk-go-v2/feature/s3/manager v1.22.25 + github.com/aws/aws-sdk-go-v2/feature/s3/manager v1.22.26 github.com/aws/aws-sdk-go-v2/service/s3 v1.103.3 github.com/aws/smithy-go v1.27.2 github.com/docker/go-connections v0.7.0 // indirect @@ -60,7 +60,7 @@ require ( golang.org/x/sync v0.21.0 golang.org/x/sys v0.46.0 golang.org/x/text v0.38.0 - google.golang.org/api v0.283.0 + google.golang.org/api v0.284.0 google.golang.org/genproto v0.0.0-20260523011958-0a33c5d7ca68 google.golang.org/grpc v1.81.1 google.golang.org/protobuf v1.36.11 @@ -204,5 +204,5 @@ require ( golang.org/x/tools v0.45.0 // indirect golang.org/x/xerrors v0.0.0-20240903120638-7835f813f4da // indirect google.golang.org/genproto/googleapis/api v0.0.0-20260523011958-0a33c5d7ca68 // indirect - google.golang.org/genproto/googleapis/rpc v0.0.0-20260523011958-0a33c5d7ca68 // indirect + google.golang.org/genproto/googleapis/rpc v0.0.0-20260526163538-3dc84a4a5aaa // indirect ) diff --git a/sdks/go.sum b/sdks/go.sum index b8a399fadc1f..a0db92fd7ec8 100644 --- a/sdks/go.sum +++ b/sdks/go.sum @@ -207,8 +207,8 @@ github.com/aws/aws-sdk-go-v2/aws/protocol/eventstream v1.7.13 h1:p1BBrg/Hhp6uK7z github.com/aws/aws-sdk-go-v2/aws/protocol/eventstream v1.7.13/go.mod h1:8cIfkE9MDhkRZGpQ22aV6/lkYeYSozpz16Smrs5x4Ls= github.com/aws/aws-sdk-go-v2/config v1.15.3/go.mod h1:9YL3v07Xc/ohTsxFXzan9ZpFpdTOFl4X65BAKYaz8jg= github.com/aws/aws-sdk-go-v2/config v1.25.3/go.mod h1:tAByZy03nH5jcq0vZmkcVoo6tRzRHEwSFx3QW4NmDw8= -github.com/aws/aws-sdk-go-v2/config v1.32.23 h1:PYDobtcsJXK6bQe9I8RQk6s19Bz3xa3xRU08Hy1Em3Y= -github.com/aws/aws-sdk-go-v2/config v1.32.23/go.mod h1:QID4dqUQVgEOYPKsPWd1sNWCCR2c5g7o3jeEtIXPOZU= +github.com/aws/aws-sdk-go-v2/config v1.32.24 h1:aEDEj533yGdVvEHfkCY0D/1FbDrjnZr4pIulxRjqpHs= +github.com/aws/aws-sdk-go-v2/config v1.32.24/go.mod h1:yZtrGKJGlqfEW+/m2uTsJK+Jz7xF5R0eZfgcIG9m1ss= github.com/aws/aws-sdk-go-v2/credentials v1.11.2/go.mod h1:j8YsY9TXTm31k4eFhspiQicfXPLZ0gYXA50i4gxPE8g= github.com/aws/aws-sdk-go-v2/credentials v1.16.2/go.mod h1:sDdvGhXrSVT5yzBDR7qXz+rhbpiMpUYfF3vJ01QSdrc= github.com/aws/aws-sdk-go-v2/credentials v1.19.23 h1:Zhu3GOpRCkNjtE/gJpuPDsytSnaCCTQk8neAGsgzG5Y= @@ -219,8 +219,8 @@ github.com/aws/aws-sdk-go-v2/feature/ec2/imds v1.18.29 h1:r6qZHbT+wxgWO/e9vYNUEt github.com/aws/aws-sdk-go-v2/feature/ec2/imds v1.18.29/go.mod h1:QRnaRcTVGKPGRy8w78HMQtKUGRYcnMZAANATkeVA6Mo= github.com/aws/aws-sdk-go-v2/feature/s3/manager v1.11.3/go.mod h1:0dHuD2HZZSiwfJSy1FO5bX1hQ1TxVV1QXXjpn3XUE44= github.com/aws/aws-sdk-go-v2/feature/s3/manager v1.14.0/go.mod h1:UcgIwJ9KHquYxs6Q5skC9qXjhYMK+JASDYcXQ4X7JZE= -github.com/aws/aws-sdk-go-v2/feature/s3/manager v1.22.25 h1:xJ3WVH3J0xESIqkavgbNfvQdMzB98bSkGwFXCyM2Tdw= -github.com/aws/aws-sdk-go-v2/feature/s3/manager v1.22.25/go.mod h1:0EOIx2v10rBBcaoTpOsRRNKDxhtRecRbN42YMrJaKZE= +github.com/aws/aws-sdk-go-v2/feature/s3/manager v1.22.26 h1:YNUEPr7Yiako/MzR/h3woMREbdwj0hGiBsZc5ZM90yE= +github.com/aws/aws-sdk-go-v2/feature/s3/manager v1.22.26/go.mod h1:KswdJ4xh+tUgW5CWx7sarhQuD+3iKgg46wojfmCA8q4= github.com/aws/aws-sdk-go-v2/internal/configsources v1.1.9/go.mod h1:AnVH5pvai0pAF4lXRq0bmhbes1u9R8wTE+g+183bZNM= github.com/aws/aws-sdk-go-v2/internal/configsources v1.2.3/go.mod h1:7sGSz1JCKHWWBHq98m6sMtWQikmYPpxjqOydDemiVoM= github.com/aws/aws-sdk-go-v2/internal/configsources v1.4.29 h1:f3vKqSo13fhTYb+JEcXwXefZQE26I1FB5eTSniU67ko= @@ -1324,8 +1324,8 @@ google.golang.org/api v0.69.0/go.mod h1:boanBiw+h5c3s+tBPgEzLDRHfFLWV0qXxRHz3ws7 google.golang.org/api v0.70.0/go.mod h1:Bs4ZM2HGifEvXwd50TtW70ovgJffJYw2oRCOFU/SkfA= google.golang.org/api v0.71.0/go.mod h1:4PyU6e6JogV1f9eA4voyrTY2batOLdgZ5qZ5HOCc4j8= google.golang.org/api v0.74.0/go.mod h1:ZpfMZOVRMywNyvJFeqL9HRWBgAuRfSjJFpe9QtRRyDs= -google.golang.org/api v0.283.0 h1:0lkp8u0MPwJVHqRL+nJlMAoZVVzbmiXmFHXMOTmSPik= -google.golang.org/api v0.283.0/go.mod h1:6Wssta4c5n9qHq5CBhmlai5h/PUa1djdDAIhYEHyvcM= +google.golang.org/api v0.284.0 h1:i+cKTgeQRcRySkP7QTl5PDO7/pAm8EcMFIUMlNbk4Vc= +google.golang.org/api v0.284.0/go.mod h1:AU44fU+XVZOCcd8uLaBIa/ZgzgPf/0qqY3+m7lQaado= google.golang.org/appengine v1.1.0/go.mod h1:EbEs0AVv82hx2wNQdGPgUI5lhzA/G0D9YwlJXL52JkM= google.golang.org/appengine v1.4.0/go.mod h1:xpcJRLb0r/rnEns0DIKYYv+WjYCduHsrkT7/EB5XEv4= google.golang.org/appengine v1.5.0/go.mod h1:xpcJRLb0r/rnEns0DIKYYv+WjYCduHsrkT7/EB5XEv4= @@ -1423,8 +1423,8 @@ google.golang.org/genproto v0.0.0-20260523011958-0a33c5d7ca68 h1:cTHF8xtqtBN5sQ4 google.golang.org/genproto v0.0.0-20260523011958-0a33c5d7ca68/go.mod h1:RRHjglSYABVCWpQ7USCpdfhcd9t4PkajvVwyynZizTc= google.golang.org/genproto/googleapis/api v0.0.0-20260523011958-0a33c5d7ca68 h1:WVVw1Nl19li0fMX++FJ3ye1z9+S1N35QODDy5qpnaXw= google.golang.org/genproto/googleapis/api v0.0.0-20260523011958-0a33c5d7ca68/go.mod h1:1dCETSCY2YKZNXQE3h4fun3TYwF5p8jejRKZgfWAgAY= -google.golang.org/genproto/googleapis/rpc v0.0.0-20260523011958-0a33c5d7ca68 h1:PvEgGJf9C/1u5CHkInMg7UFYYUoiaQmW2LbtH0pjB78= -google.golang.org/genproto/googleapis/rpc v0.0.0-20260523011958-0a33c5d7ca68/go.mod h1:4Hqkh8ycfw05ld/3BWL7rJOSfebL2Q+DVDeRgYgxUU8= +google.golang.org/genproto/googleapis/rpc v0.0.0-20260526163538-3dc84a4a5aaa h1:mZHHdPZl0dbGHCflZgAq/Q468DWVFcU2whhB2KAo8fk= +google.golang.org/genproto/googleapis/rpc v0.0.0-20260526163538-3dc84a4a5aaa/go.mod h1:4Hqkh8ycfw05ld/3BWL7rJOSfebL2Q+DVDeRgYgxUU8= google.golang.org/grpc v1.19.0/go.mod h1:mqu4LbDTu4XGKhr4mRzUsmM4RtVoemTSY81AxZiDr8c= google.golang.org/grpc v1.20.1/go.mod h1:10oTOabMzJvdu6/UiuZezV6QK5dSlG84ov/aaiqXj38= google.golang.org/grpc v1.21.1/go.mod h1:oYelfM1adQP15Ek0mdvEgi9Df8B9CZIaU1084ijfRaM= diff --git a/sdks/java/io/datadog/src/main/java/org/apache/beam/sdk/io/datadog/DatadogEvent.java b/sdks/java/io/datadog/src/main/java/org/apache/beam/sdk/io/datadog/DatadogEvent.java index 80334b5e4664..e9a2546d9d9e 100644 --- a/sdks/java/io/datadog/src/main/java/org/apache/beam/sdk/io/datadog/DatadogEvent.java +++ b/sdks/java/io/datadog/src/main/java/org/apache/beam/sdk/io/datadog/DatadogEvent.java @@ -26,6 +26,12 @@ @AutoValue public abstract class DatadogEvent { + public static final String SOURCE = "ddsource"; + public static final String TAGS = "ddtags"; + public static final String HOSTNAME = "hostname"; + public static final String SERVICE = "service"; + public static final String MESSAGE = "message"; + public static Builder newBuilder() { return new AutoValue_DatadogEvent.Builder(); } diff --git a/sdks/java/io/datadog/src/main/java/org/apache/beam/sdk/io/datadog/DatadogWriteSchemaTransformConfiguration.java b/sdks/java/io/datadog/src/main/java/org/apache/beam/sdk/io/datadog/DatadogWriteSchemaTransformConfiguration.java new file mode 100644 index 000000000000..059b3b33ed4f --- /dev/null +++ b/sdks/java/io/datadog/src/main/java/org/apache/beam/sdk/io/datadog/DatadogWriteSchemaTransformConfiguration.java @@ -0,0 +1,114 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.apache.beam.sdk.io.datadog; + +import static org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.base.Preconditions.checkArgument; + +import com.google.auto.value.AutoValue; +import javax.annotation.Nullable; +import org.apache.beam.sdk.schemas.AutoValueSchema; +import org.apache.beam.sdk.schemas.annotations.DefaultSchema; +import org.apache.beam.sdk.schemas.annotations.SchemaFieldDescription; +import org.apache.beam.sdk.schemas.transforms.providers.ErrorHandling; +import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.base.Strings; + +/** + * Configuration for writing to Datadog. + * + *

This class is meant to be used with {@link DatadogWriteSchemaTransformProvider}. + */ +@DefaultSchema(AutoValueSchema.class) +@AutoValue +public abstract class DatadogWriteSchemaTransformConfiguration { + + public void validate() { + String invalidConfigMessage = "Invalid Datadog Write configuration: "; + checkArgument(!getUrl().isEmpty(), invalidConfigMessage + "url must be specified."); + checkArgument(!getApiKey().isEmpty(), invalidConfigMessage + "apiKey must be specified."); + Integer batchCount = getBatchCount(); + if (batchCount != null) { + checkArgument(batchCount > 0, invalidConfigMessage + "batchCount must be greater than 0."); + } + Integer minBatchCount = getMinBatchCount(); + if (minBatchCount != null) { + checkArgument( + minBatchCount > 0, invalidConfigMessage + "minBatchCount must be greater than 0."); + } + Long maxBufferSize = getMaxBufferSize(); + if (maxBufferSize != null) { + checkArgument( + maxBufferSize > 0, invalidConfigMessage + "maxBufferSize must be greater than 0."); + } + Integer parallelism = getParallelism(); + if (parallelism != null) { + checkArgument(parallelism > 0, invalidConfigMessage + "parallelism must be greater than 0."); + } + ErrorHandling errorHandling = getErrorHandling(); + if (errorHandling != null) { + checkArgument( + !Strings.isNullOrEmpty(errorHandling.getOutput()), + invalidConfigMessage + "Output must not be empty if error handling specified."); + } + } + + /** Instantiates a {@link DatadogWriteSchemaTransformConfiguration.Builder} instance. */ + public static DatadogWriteSchemaTransformConfiguration.Builder builder() { + return new AutoValue_DatadogWriteSchemaTransformConfiguration.Builder(); + } + + @SchemaFieldDescription("The Datadog API URL.") + public abstract String getUrl(); + + @SchemaFieldDescription("The Datadog API key.") + public abstract String getApiKey(); + + @SchemaFieldDescription("The minimum number of events to batch together for each write.") + public abstract @Nullable Integer getMinBatchCount(); + + @SchemaFieldDescription("The number of events to batch together for each write.") + public abstract @Nullable Integer getBatchCount(); + + @SchemaFieldDescription("The maximum buffer size in bytes.") + public abstract @Nullable Long getMaxBufferSize(); + + @SchemaFieldDescription("The degree of parallelism for writing.") + public abstract @Nullable Integer getParallelism(); + + @SchemaFieldDescription("Specifies how to handle errors.") + public abstract @Nullable ErrorHandling getErrorHandling(); + + @AutoValue.Builder + public abstract static class Builder { + public abstract Builder setUrl(String url); + + public abstract Builder setApiKey(String apiKey); + + public abstract Builder setMinBatchCount(Integer minBatchCount); + + public abstract Builder setBatchCount(Integer batchCount); + + public abstract Builder setMaxBufferSize(Long maxBufferSize); + + public abstract Builder setParallelism(Integer parallelism); + + public abstract Builder setErrorHandling(@Nullable ErrorHandling errorHandling); + + /** Builds the {@link DatadogWriteSchemaTransformConfiguration} configuration. */ + public abstract DatadogWriteSchemaTransformConfiguration build(); + } +} diff --git a/sdks/java/io/datadog/src/main/java/org/apache/beam/sdk/io/datadog/DatadogWriteSchemaTransformProvider.java b/sdks/java/io/datadog/src/main/java/org/apache/beam/sdk/io/datadog/DatadogWriteSchemaTransformProvider.java new file mode 100644 index 000000000000..d0013ac04225 --- /dev/null +++ b/sdks/java/io/datadog/src/main/java/org/apache/beam/sdk/io/datadog/DatadogWriteSchemaTransformProvider.java @@ -0,0 +1,294 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.apache.beam.sdk.io.datadog; + +import static org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.base.Preconditions.checkNotNull; + +import com.google.auto.service.AutoService; +import java.util.Collections; +import java.util.List; +import org.apache.beam.sdk.coders.RowCoder; +import org.apache.beam.sdk.schemas.Schema; +import org.apache.beam.sdk.schemas.transforms.SchemaTransform; +import org.apache.beam.sdk.schemas.transforms.SchemaTransformProvider; +import org.apache.beam.sdk.schemas.transforms.TypedSchemaTransformProvider; +import org.apache.beam.sdk.schemas.transforms.providers.ErrorHandling; +import org.apache.beam.sdk.transforms.Create; +import org.apache.beam.sdk.transforms.DoFn; +import org.apache.beam.sdk.transforms.ParDo; +import org.apache.beam.sdk.values.PCollection; +import org.apache.beam.sdk.values.PCollectionRowTuple; +import org.apache.beam.sdk.values.PCollectionTuple; +import org.apache.beam.sdk.values.Row; +import org.apache.beam.sdk.values.TupleTag; +import org.apache.beam.sdk.values.TupleTagList; + +@AutoService(SchemaTransformProvider.class) +public class DatadogWriteSchemaTransformProvider + extends TypedSchemaTransformProvider { + private static final String IDENTIFIER = "beam:schematransform:org.apache.beam:datadog_write:v1"; + static final String INPUT = "input"; + static final String OUTPUT = "output"; + static final String ERROR = "errors"; + public static final TupleTag ERROR_TAG = new TupleTag() {}; + public static final TupleTag OUTPUT_TAG = new TupleTag() {}; + public static final TupleTag EVENT_TAG = new TupleTag() {}; + + @Override + protected Class configurationClass() { + return DatadogWriteSchemaTransformConfiguration.class; + } + + /** Returns the expected {@link SchemaTransform} of the configuration. */ + @Override + protected SchemaTransform from(DatadogWriteSchemaTransformConfiguration configuration) { + return new DatadogWriteSchemaTransform(configuration); + } + + /** Implementation of the {@link TypedSchemaTransformProvider} identifier method. */ + @Override + public String identifier() { + return IDENTIFIER; + } + + /** Implementation of the {@link TypedSchemaTransformProvider} input collection names method. */ + @Override + public List inputCollectionNames() { + return Collections.singletonList(INPUT); + } + + /** Implementation of the {@link TypedSchemaTransformProvider} output collection names method. */ + @Override + public List outputCollectionNames() { + return Collections.singletonList(ERROR); + } + + /** + * An implementation of {@link SchemaTransform} for Datadog Write jobs configured using {@link + * DatadogWriteSchemaTransformConfiguration}. + */ + static class DatadogWriteSchemaTransform extends SchemaTransform { + private final DatadogWriteSchemaTransformConfiguration configuration; + + DatadogWriteSchemaTransform(DatadogWriteSchemaTransformConfiguration configuration) { + this.configuration = configuration; + } + + @Override + public PCollectionRowTuple expand(PCollectionRowTuple input) { + // Validate configuration parameters + configuration.validate(); + + // Obtain input rows + PCollection inputRows = input.get(INPUT); + + // Check for errors + boolean handleErrors = ErrorHandling.hasOutput(configuration.getErrorHandling()); + + Schema inputSchema = inputRows.getSchema(); + Schema dynamicErrorSchema = + Schema.builder() + .addNullableRowField("failed_row", inputSchema) + .addNullableField("payload", Schema.FieldType.STRING) + .addNullableField("statusCode", Schema.FieldType.INT32) + .addNullableField("statusMessage", Schema.FieldType.STRING) + .build(); + + PCollectionTuple convertResult = + inputRows.apply( + "Convert to DatadogEvent", + ParDo.of(new RowToEventFn(handleErrors, ERROR_TAG, dynamicErrorSchema)) + .withOutputTags(EVENT_TAG, TupleTagList.of(ERROR_TAG))); + + PCollection datadogEvents = + convertResult.get(EVENT_TAG).setCoder(DatadogEventCoder.of()); + PCollection conversionErrors = + convertResult + .get(ERROR_TAG) + .setCoder(org.apache.beam.sdk.coders.RowCoder.of(dynamicErrorSchema)); + + // Configure DatadogIO.Write + DatadogIO.Write.Builder builder = + DatadogIO.writeBuilder(configuration.getMinBatchCount()) + .withUrl(configuration.getUrl()) + .withApiKey(configuration.getApiKey()); + + Integer batchCount = configuration.getBatchCount(); + if (batchCount != null) { + builder = builder.withBatchCount(batchCount); + } + Long maxBufferSize = configuration.getMaxBufferSize(); + if (maxBufferSize != null) { + builder = builder.withMaxBufferSize(maxBufferSize); + } + Integer parallelism = configuration.getParallelism(); + if (parallelism != null) { + builder = builder.withParallelism(parallelism); + } + + DatadogIO.Write write = builder.build(); + + // Apply DatadogIO.Write + PCollection writeErrors = datadogEvents.apply("Write To Datadog", write); + + // Handle errors + ErrorHandling errorHandling = configuration.getErrorHandling(); + if (errorHandling != null) { + PCollection writeErrorRows = + writeErrors + .apply( + "Convert Write Errors to Rows", + org.apache.beam.sdk.transforms.MapElements.into( + org.apache.beam.sdk.values.TypeDescriptors.rows()) + .via( + error -> + Row.withSchema(dynamicErrorSchema) + .addValue(null) + .addValue(error.payload()) + .addValue(error.statusCode()) + .addValue(error.statusMessage()) + .build())) + .setCoder(org.apache.beam.sdk.coders.RowCoder.of(dynamicErrorSchema)); + + PCollection allErrors = + org.apache.beam.sdk.values.PCollectionList.of(conversionErrors) + .and(writeErrorRows) + .apply("Flatten Errors", org.apache.beam.sdk.transforms.Flatten.pCollections()) + .setCoder(org.apache.beam.sdk.coders.RowCoder.of(dynamicErrorSchema)); + + return PCollectionRowTuple.of(errorHandling.getOutput(), allErrors); + } else { + writeErrors.apply("Fail on Write Error", ParDo.of(new FailOnWriteErrorFn())); + PCollection emptyErrors = + input + .getPipeline() + .apply("Empty Errors Placeholder", Create.empty(RowCoder.of(dynamicErrorSchema))); + return PCollectionRowTuple.of(ERROR, emptyErrors); + } + } + } + + static final Schema WRITE_ERROR_SCHEMA = + Schema.builder() + .addNullableField("payload", Schema.FieldType.STRING) + .addNullableField("statusCode", Schema.FieldType.INT32) + .addNullableField("statusMessage", Schema.FieldType.STRING) + .build(); + + static final Schema DATADOG_EVENT_SCHEMA = + Schema.builder() + .addNullableField(DatadogEvent.SOURCE, Schema.FieldType.STRING) + .addNullableField(DatadogEvent.TAGS, Schema.FieldType.STRING) + .addNullableField(DatadogEvent.HOSTNAME, Schema.FieldType.STRING) + .addNullableField(DatadogEvent.SERVICE, Schema.FieldType.STRING) + .addNullableField(DatadogEvent.MESSAGE, Schema.FieldType.STRING) + .build(); + + static Row eventToRow(DatadogEvent event) { + return Row.withSchema(DATADOG_EVENT_SCHEMA) + .addValue(event.ddsource()) + .addValue(event.ddtags()) + .addValue(event.hostname()) + .addValue(event.service()) + .addValue(event.message()) + .build(); + } + + static DatadogEvent rowToEvent(Row row) { + DatadogEvent.Builder builder = DatadogEvent.newBuilder(); + Schema schema = row.getSchema(); + + String ddsource = + schema.hasField(DatadogEvent.SOURCE) ? row.getString(DatadogEvent.SOURCE) : null; + if (ddsource != null) { + builder.withSource(ddsource); + } + String ddtags = schema.hasField(DatadogEvent.TAGS) ? row.getString(DatadogEvent.TAGS) : null; + if (ddtags != null) { + builder.withTags(ddtags); + } + String hostname = + schema.hasField(DatadogEvent.HOSTNAME) ? row.getString(DatadogEvent.HOSTNAME) : null; + if (hostname != null) { + builder.withHostname(hostname); + } + String service = + schema.hasField(DatadogEvent.SERVICE) ? row.getString(DatadogEvent.SERVICE) : null; + if (service != null) { + builder.withService(service); + } + String message = + schema.hasField(DatadogEvent.MESSAGE) ? row.getString(DatadogEvent.MESSAGE) : null; + builder.withMessage(checkNotNull(message, "Message is required.")); + + return builder.build(); + } + + static class RowToEventFn extends DoFn { + private final boolean handleErrors; + private final TupleTag errorOutputTag; + private final Schema errorSchema; + + RowToEventFn(boolean handleErrors, TupleTag errorOutputTag, Schema errorSchema) { + this.handleErrors = handleErrors; + this.errorOutputTag = errorOutputTag; + this.errorSchema = errorSchema; + } + + @ProcessElement + public void processElement(ProcessContext c) { + try { + c.output(rowToEvent(c.element())); + } catch (Exception e) { + if (handleErrors) { + String rowString = c.element().toString(); + String payload = rowString.length() <= 1024 ? rowString : rowString.substring(0, 1024); + c.output( + errorOutputTag, + Row.withSchema(errorSchema) + .addValue(c.element()) + .addValue(payload) + .addValue(java.net.HttpURLConnection.HTTP_BAD_REQUEST) + .addValue(e.getMessage()) + .build()); + } else { + throw new RuntimeException(e); + } + } + } + } + + /** + * A {@link DoFn} that throws a {@link RuntimeException} when a write error is encountered, + * causing the pipeline to fail. This is the default error handling behavior when no error output + * is configured. + */ + static class FailOnWriteErrorFn extends DoFn { + @ProcessElement + public void processElement(@Element DatadogWriteError error) { + String message = error.statusMessage(); + if (error.statusCode() != null) { + throw new RuntimeException( + String.format( + "Datadog write failed with status code %d: %s", error.statusCode(), message)); + } else { + throw new RuntimeException("Datadog write failed: " + message); + } + } + } +} diff --git a/sdks/java/io/datadog/src/test/java/org/apache/beam/sdk/io/datadog/DatadogWriteSchemaTransformProviderTest.java b/sdks/java/io/datadog/src/test/java/org/apache/beam/sdk/io/datadog/DatadogWriteSchemaTransformProviderTest.java new file mode 100644 index 000000000000..534251fb4c19 --- /dev/null +++ b/sdks/java/io/datadog/src/test/java/org/apache/beam/sdk/io/datadog/DatadogWriteSchemaTransformProviderTest.java @@ -0,0 +1,535 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.apache.beam.sdk.io.datadog; + +import static org.apache.beam.sdk.io.datadog.DatadogWriteSchemaTransformProvider.ERROR; +import static org.apache.beam.sdk.io.datadog.DatadogWriteSchemaTransformProvider.INPUT; +import static org.apache.beam.sdk.io.datadog.DatadogWriteSchemaTransformProvider.OUTPUT; +import static org.junit.Assert.assertEquals; +import static org.junit.Assert.assertTrue; +import static org.junit.jupiter.api.Assertions.assertThrows; + +import java.util.Arrays; +import java.util.List; +import java.util.ServiceLoader; +import java.util.stream.Collectors; +import java.util.stream.StreamSupport; +import org.apache.beam.sdk.Pipeline.PipelineExecutionException; +import org.apache.beam.sdk.schemas.AutoValueSchema; +import org.apache.beam.sdk.schemas.NoSuchSchemaException; +import org.apache.beam.sdk.schemas.Schema; +import org.apache.beam.sdk.schemas.SchemaRegistry; +import org.apache.beam.sdk.schemas.transforms.SchemaTransform; +import org.apache.beam.sdk.schemas.transforms.SchemaTransformProvider; +import org.apache.beam.sdk.schemas.transforms.providers.ErrorHandling; +import org.apache.beam.sdk.testing.PAssert; +import org.apache.beam.sdk.testing.TestPipeline; +import org.apache.beam.sdk.transforms.Create; +import org.apache.beam.sdk.values.PCollection; +import org.apache.beam.sdk.values.PCollectionRowTuple; +import org.apache.beam.sdk.values.Row; +import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.Lists; +import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.Sets; +import org.junit.Rule; +import org.junit.Test; +import org.junit.runner.RunWith; +import org.junit.runners.JUnit4; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +@RunWith(JUnit4.class) +public class DatadogWriteSchemaTransformProviderTest { + + private static final Logger LOG = + LoggerFactory.getLogger(DatadogWriteSchemaTransformProviderTest.class); + + @Rule public TestPipeline p = TestPipeline.create(); + + private static final Schema SCHEMA = + Schema.builder() + .addStringField("ddsource") + .addNullableField("ddtags", Schema.FieldType.STRING) + .addStringField("hostname") + .addNullableField("service", Schema.FieldType.STRING) + .addStringField("message") + .build(); + + private static final List ROWS = + Arrays.asList( + Row.withSchema(SCHEMA) + .withFieldValue("ddsource", "my-source") + .withFieldValue("ddtags", "tag1:value1,tag2") + .withFieldValue("hostname", "my-host") + .withFieldValue("service", "my-service") + .withFieldValue("message", "Hello World 1") + .build(), + Row.withSchema(SCHEMA) + .withFieldValue("ddsource", "my-source-2") + .withFieldValue("ddtags", null) + .withFieldValue("hostname", "my-host-2") + .withFieldValue("service", null) + .withFieldValue("message", "Hello World 2") + .build()); + + private List events; + + @org.junit.Before + public void setUp() { + events = + Arrays.asList( + DatadogEvent.newBuilder() + .withSource("my-source") + .withTags("tag1:value1,tag2") + .withHostname("my-host") + .withService("my-service") + .withMessage("Hello World 1") + .build(), + DatadogEvent.newBuilder() + .withSource("my-source-2") + .withHostname("my-host-2") + .withMessage("Hello World 2") + .build()); + } + + @Test + public void testWriteInvalidConfigurations() { + + // apiKey not set + assertThrows( + IllegalStateException.class, + () -> { + DatadogWriteSchemaTransformConfiguration.builder() + .setUrl("http://localhost:8080") + // .setApiKey("test-api-key") # ApiKey is mandatory + .build() + .validate(); + }); + + // url not set + assertThrows( + IllegalStateException.class, + () -> { + DatadogWriteSchemaTransformConfiguration.builder() + // .setUrl("http://localhost:8080") # Url is mandatory + .setApiKey("test-api-key") + .build() + .validate(); + }); + } + + @Test + public void testWriteBuildTransform() { + DatadogWriteSchemaTransformProvider provider = new DatadogWriteSchemaTransformProvider(); + DatadogWriteSchemaTransformConfiguration configuration = + DatadogWriteSchemaTransformConfiguration.builder() + .setApiKey("test-api-key") + .setUrl("http://localhost:8080") + .build(); + + provider.from(configuration); + } + + @Test + public void testWriteBuildTransformAndRun() { + DatadogWriteSchemaTransformProvider provider = new DatadogWriteSchemaTransformProvider(); + DatadogWriteSchemaTransformConfiguration configuration = + DatadogWriteSchemaTransformConfiguration.builder() + .setApiKey("test-api-key") + .setUrl("http://localhost:8080") + .build(); + + SchemaTransform transform = provider.from(configuration); + + PCollection input = p.apply("Create", Create.of(ROWS).withRowSchema(SCHEMA)); + PCollectionRowTuple inputTuple = PCollectionRowTuple.of(INPUT, input); + PCollectionRowTuple output = transform.expand(inputTuple); + assertEquals(1, output.getAll().size()); + assertTrue(output.has(ERROR)); + + assertThrows(PipelineExecutionException.class, () -> p.run().waitUntilFinish()); + } + + @Test + public void testWriteBuildTransformWithCorrectFields() { + ServiceLoader serviceLoader = + ServiceLoader.load(SchemaTransformProvider.class); + List providers = + StreamSupport.stream(serviceLoader.spliterator(), false) + .filter(provider -> provider.getClass() == DatadogWriteSchemaTransformProvider.class) + .collect(Collectors.toList()); + SchemaTransformProvider datadogProvider = providers.get(0); + assertEquals(datadogProvider.outputCollectionNames(), Lists.newArrayList(ERROR)); + + assertEquals( + Sets.newHashSet( + "url", + "api_key", + "min_batch_count", + "batch_count", + "max_buffer_size", + "parallelism", + "error_handling"), + datadogProvider.configurationSchema().getFields().stream() + .map(field -> field.getName()) + .collect(Collectors.toSet())); + } + + @Test + public void testRowToDatadogEvent() { + for (int i = 0; i < ROWS.size(); i++) { + DatadogEvent actual = DatadogWriteSchemaTransformProvider.rowToEvent(ROWS.get(i)); + assertEquals(events.get(i), actual); + } + } + + @Test + public void testRowToDatadogEventWithMissingOptionalFields() { + Schema missingFieldsSchema = + Schema.builder() + .addStringField("ddsource") + .addStringField("hostname") + .addStringField("message") + .build(); + + Row row = + Row.withSchema(missingFieldsSchema) + .withFieldValue("ddsource", "my-source") + .withFieldValue("hostname", "my-host") + .withFieldValue("message", "Hello World 1") + .build(); + + DatadogEvent expectedEvent = + DatadogEvent.newBuilder() + .withSource("my-source") + .withHostname("my-host") + .withMessage("Hello World 1") + .build(); + + DatadogEvent actual = DatadogWriteSchemaTransformProvider.rowToEvent(row); + assertEquals(expectedEvent, actual); + } + + @Test + public void testRowToDatadogEventWithExtraFields_DiscardsExtraFields() { + Schema extraFieldsSchema = + Schema.builder() + .addStringField("ddsource") + .addNullableField("ddtags", Schema.FieldType.STRING) + .addStringField("hostname") + .addNullableField("service", Schema.FieldType.STRING) + .addStringField("message") + .addStringField("extra_field") + .build(); + + Row row = + Row.withSchema(extraFieldsSchema) + .withFieldValue("ddsource", "my-source") + .withFieldValue("ddtags", "tag1:value1,tag2") + .withFieldValue("hostname", "my-host") + .withFieldValue("service", "my-service") + .withFieldValue("message", "Hello World 1") + .withFieldValue("extra_field", "extra_value") + .build(); + + DatadogEvent expectedEvent = + DatadogEvent.newBuilder() + .withSource("my-source") + .withTags("tag1:value1,tag2") + .withHostname("my-host") + .withService("my-service") + .withMessage("Hello World 1") + .build(); + + DatadogEvent actual = DatadogWriteSchemaTransformProvider.rowToEvent(row); + assertEquals(expectedEvent, actual); + } + + @Test(expected = NullPointerException.class) + public void testRowToDatadogEventWithNullRequiredField() { + Schema nullSchema = + Schema.builder() + .addStringField("ddsource") + .addNullableField("ddtags", Schema.FieldType.STRING) + .addStringField("hostname") + .addNullableField("service", Schema.FieldType.STRING) + .addNullableField("message", Schema.FieldType.STRING) + .build(); + + Row row = + Row.withSchema(nullSchema) + .withFieldValue("ddsource", "my-source") + .withFieldValue("ddtags", "tag1:value1,tag2") + .withFieldValue("hostname", "my-host") + .withFieldValue("service", "my-service") + .withFieldValue("message", null) + .build(); + + DatadogWriteSchemaTransformProvider.rowToEvent(row); + } + + @Test + public void testBuildTransformWithAllParameters() { + DatadogWriteSchemaTransformProvider provider = new DatadogWriteSchemaTransformProvider(); + DatadogWriteSchemaTransformConfiguration configuration = + DatadogWriteSchemaTransformConfiguration.builder() + .setApiKey("test-api-key") + .setUrl("http://localhost:8080") + .setBatchCount(10) + .setMaxBufferSize(100L) + .setParallelism(2) + .build(); + + SchemaTransform transform = provider.from(configuration); + + PCollection input = p.apply("Create", Create.of(ROWS).withRowSchema(SCHEMA)); + PCollectionRowTuple inputTuple = PCollectionRowTuple.of(INPUT, input); + PCollectionRowTuple output = transform.expand(inputTuple); + assertEquals(1, output.getAll().size()); + assertTrue(output.has(ERROR)); + + assertThrows(PipelineExecutionException.class, () -> p.run().waitUntilFinish()); + } + + @Test(expected = IllegalStateException.class) + public void testBuildTransformMissingUrl() { + DatadogWriteSchemaTransformProvider provider = new DatadogWriteSchemaTransformProvider(); + provider.from( + DatadogWriteSchemaTransformConfiguration.builder().setApiKey("test-api-key").build()); + } + + @Test(expected = IllegalStateException.class) + public void testBuildTransformMissingApiKey() { + DatadogWriteSchemaTransformProvider provider = new DatadogWriteSchemaTransformProvider(); + provider.from(DatadogWriteSchemaTransformConfiguration.builder().setUrl("test-url").build()); + } + + @Test + public void testBuildTransformWithInvalidParallelism() { + assertThrows( + IllegalArgumentException.class, + () -> { + DatadogWriteSchemaTransformConfiguration.builder() + .setApiKey("test-api-key") + .setUrl("http://localhost:8080") + .setParallelism(0) + .build() + .validate(); + }); + } + + @Test + public void testBuildTransformWithInvalidBatchCount() { + assertThrows( + IllegalArgumentException.class, + () -> { + DatadogWriteSchemaTransformConfiguration.builder() + .setApiKey("test-api-key") + .setUrl("http://localhost:8080") + .setBatchCount(0) + .setMinBatchCount(1) + .build() + .validate(); + }); + } + + @Test + public void testBuildTransformFromRowConfiguration() throws NoSuchSchemaException { + DatadogWriteSchemaTransformProvider provider = new DatadogWriteSchemaTransformProvider(); + Schema configSchema = provider.configurationSchema(); + Schema errorHandlingSchema = configSchema.getField("error_handling").getType().getRowSchema(); + + Row errorHandlingRow = + Row.withSchema(errorHandlingSchema).withFieldValue(OUTPUT, ERROR).build(); + + Row configRow = + Row.withSchema(configSchema) + .withFieldValue("url", "http://localhost:8080") + .withFieldValue("api_key", "test-api-key") + .withFieldValue("min_batch_count", null) + .withFieldValue("batch_count", 10) + .withFieldValue("max_buffer_size", 100L) + .withFieldValue("parallelism", 2) + .withFieldValue("error_handling", errorHandlingRow) + .build(); + + SchemaTransform transform = provider.from(configRow); + + PCollection input = p.apply("Create", Create.of(ROWS).withRowSchema(SCHEMA)); + PCollectionRowTuple inputTuple = PCollectionRowTuple.of(INPUT, input); + PCollectionRowTuple output = transform.expand(inputTuple); + assertEquals(1, output.getAll().size()); + assertTrue(output.has(ERROR)); + + p.run().waitUntilFinish(); + } + + @Test(expected = ClassCastException.class) + public void testRowToDatadogEventWithWrongType() { + Schema wrongSchema = + Schema.builder() + .addStringField("ddsource") + .addInt64Field("ddtags") + .addStringField("hostname") + .addStringField("message") + .build(); + + Row row = + Row.withSchema(wrongSchema) + .withFieldValue("ddsource", "my-source") + .withFieldValue("ddtags", 123L) + .withFieldValue("hostname", "my-host") + .withFieldValue("message", "Hello World 1") + .build(); + + DatadogWriteSchemaTransformProvider.rowToEvent(row); + } + + @Test + public void testErrorHandling() { + DatadogWriteSchemaTransformProvider provider = new DatadogWriteSchemaTransformProvider(); + ErrorHandling errorHandling = ErrorHandling.builder().setOutput(ERROR).build(); + DatadogWriteSchemaTransformConfiguration configuration = + DatadogWriteSchemaTransformConfiguration.builder() + .setApiKey("test-api-key") + .setUrl("http://localhost:8080") + .setBatchCount(10) + .setMaxBufferSize(1L) + .setParallelism(1) + .setErrorHandling(errorHandling) + .build(); + + Schema nullSchema = + Schema.builder() + .addStringField("ddsource") + .addNullableField("ddtags", Schema.FieldType.STRING) + .addStringField("hostname") + .addNullableField("service", Schema.FieldType.STRING) + .addNullableField("message", Schema.FieldType.STRING) + .build(); + + Row row = + Row.withSchema(nullSchema) + .withFieldValue("ddsource", "my-source") + .withFieldValue("ddtags", "tag1:value1,tag2") + .withFieldValue("hostname", "my-host") + .withFieldValue("service", "my-service") + .withFieldValue("message", null) + .build(); + + PCollection input = p.apply(Create.of(row).withRowSchema(nullSchema)); + PCollectionRowTuple inputTuple = PCollectionRowTuple.of(INPUT, input); + + SchemaTransform transform = provider.from(configuration); + PCollectionRowTuple outputTuple = transform.expand(inputTuple); + + assertTrue(outputTuple.has(ERROR)); + PAssert.that(outputTuple.get(ERROR)) + .satisfies( + (errors) -> { + assertEquals(1, errors.spliterator().getExactSizeIfKnown()); + Row error = errors.iterator().next(); + assertEquals(row.toString(), error.getString("payload")); + assertEquals( + (Integer) java.net.HttpURLConnection.HTTP_BAD_REQUEST, + error.getInt32("statusCode")); + assertTrue( + "Expected status message to contain 'Message is required.'", + error.getString("statusMessage").contains("Message is required.")); + return null; + }); + + p.run().waitUntilFinish(); + } + + @Test + public void testErrorHandlingWithMissingRequiredField() { + DatadogWriteSchemaTransformProvider provider = new DatadogWriteSchemaTransformProvider(); + ErrorHandling errorHandling = ErrorHandling.builder().setOutput(ERROR).build(); + DatadogWriteSchemaTransformConfiguration configuration = + DatadogWriteSchemaTransformConfiguration.builder() + .setApiKey("test-api-key") + .setUrl("http://localhost:8080") + .setErrorHandling(errorHandling) + .build(); + + Schema missingFieldSchema = + Schema.builder().addStringField("ddsource").addStringField("hostname").build(); + + Row row = + Row.withSchema(missingFieldSchema) + .withFieldValue("ddsource", "my-source") + .withFieldValue("hostname", "my-host") + .build(); + + PCollection input = p.apply(Create.of(row).withRowSchema(missingFieldSchema)); + PCollectionRowTuple inputTuple = PCollectionRowTuple.of(INPUT, input); + + SchemaTransform transform = provider.from(configuration); + PCollectionRowTuple outputTuple = transform.expand(inputTuple); + + assertTrue(outputTuple.has(ERROR)); + PAssert.that(outputTuple.get(ERROR)) + .satisfies( + (errors) -> { + assertEquals(1, errors.spliterator().getExactSizeIfKnown()); + Row error = errors.iterator().next(); + assertEquals(row.toString(), error.getString("payload")); + assertEquals( + (Integer) java.net.HttpURLConnection.HTTP_BAD_REQUEST, + error.getInt32("statusCode")); + assertTrue( + "Expected status message to contain 'Message is required.'", + error.getString("statusMessage").contains("Message is required.")); + return null; + }); + + p.run().waitUntilFinish(); + } + + @Test + public void testConfigurationSchema() throws NoSuchSchemaException { + SchemaRegistry registry = SchemaRegistry.createDefault(); + registry.registerSchemaProvider(ErrorHandling.class, new AutoValueSchema()); + Schema schema = registry.getSchema(DatadogWriteSchemaTransformConfiguration.class); + Schema errorHandlingSchema = registry.getSchema(ErrorHandling.class); + + LOG.info("Schema fields: {}", schema.getFieldNames()); + + assertEquals(7, schema.getFieldCount()); + assertTrue(schema.hasField("url")); + assertTrue(schema.hasField("apiKey")); + assertTrue(schema.hasField("minBatchCount")); + assertTrue(schema.hasField("batchCount")); + assertTrue(schema.hasField("maxBufferSize")); + assertTrue(schema.hasField("parallelism")); + assertTrue(schema.hasField("errorHandling")); + + LOG.info("URL field type: {}", schema.getField("url").getType().getTypeName()); + assertEquals(Schema.FieldType.STRING.withNullable(false), schema.getField("url").getType()); + assertEquals(Schema.FieldType.STRING.withNullable(false), schema.getField("apiKey").getType()); + assertEquals( + Schema.FieldType.INT32.withNullable(true), schema.getField("batchCount").getType()); + assertEquals( + Schema.FieldType.INT64.withNullable(true), schema.getField("maxBufferSize").getType()); + assertEquals( + Schema.FieldType.INT32.withNullable(true), schema.getField("parallelism").getType()); + assertEquals( + Schema.FieldType.row(errorHandlingSchema).withNullable(true), + schema.getField("errorHandling").getType()); + } +} diff --git a/sdks/java/io/expansion-service/build.gradle b/sdks/java/io/expansion-service/build.gradle index 32894b978094..60ef89ed223b 100644 --- a/sdks/java/io/expansion-service/build.gradle +++ b/sdks/java/io/expansion-service/build.gradle @@ -96,6 +96,7 @@ dependencies { runtimeOnly project(path: ":sdks:java:io:iceberg:bqms", configuration: "shadow") runtimeOnly library.java.bigdataoss_util_hadoop + runtimeOnly project(":sdks:java:io:datadog") runtimeOnly project(":sdks:java:io:mongodb") runtimeOnly library.java.kafka_clients runtimeOnly library.java.slf4j_jdk14 diff --git a/sdks/python/apache_beam/examples/inference/online_clustering/clustering_pipeline/setup.py b/sdks/python/apache_beam/examples/inference/online_clustering/clustering_pipeline/setup.py index 572763492f41..69f6aecf2d5f 100644 --- a/sdks/python/apache_beam/examples/inference/online_clustering/clustering_pipeline/setup.py +++ b/sdks/python/apache_beam/examples/inference/online_clustering/clustering_pipeline/setup.py @@ -29,7 +29,7 @@ REQUIREMENTS = [ "apache-beam[gcp]==2.40.0", "transformers==4.38.0", - "torch==1.13.1", + "torch==2.12.0", "scikit-learn==1.0.2", ] diff --git a/sdks/python/apache_beam/metrics/metric.py b/sdks/python/apache_beam/metrics/metric.py index a66eed640be6..6e6757be11d9 100644 --- a/sdks/python/apache_beam/metrics/metric.py +++ b/sdks/python/apache_beam/metrics/metric.py @@ -46,13 +46,11 @@ from apache_beam.metrics.metricbase import Histogram from apache_beam.metrics.metricbase import MetricName from apache_beam.metrics.metricbase import StringSet -from apache_beam.options.pipeline_options import DebugOptions if TYPE_CHECKING: from apache_beam.internal.metrics.metric import MetricLogger from apache_beam.metrics.execution import MetricKey from apache_beam.metrics.metricbase import Metric - from apache_beam.options.pipeline_options import PipelineOptions from apache_beam.utils.histogram import BucketType __all__ = ['Metrics', 'MetricsFilter', 'Lineage'] @@ -60,37 +58,6 @@ _LOGGER = logging.getLogger(__name__) -class MetricsFlag(object): - """Process-wide flags controlling which user metric kinds are emitted.""" - counter_disabled = False - string_set_disabled = False - bounded_trie_disabled = False - _initialized = False - - @classmethod - def set_default_pipeline_options(cls, options: 'PipelineOptions') -> None: - if cls._initialized: - return - debug_options = options.view_as(DebugOptions) - if debug_options.lookup_experiment('disableCounterMetrics'): - cls.counter_disabled = True - _LOGGER.info('Counter metrics are disabled.') - if debug_options.lookup_experiment('disableStringSetMetrics'): - cls.string_set_disabled = True - _LOGGER.info('StringSet metrics are disabled.') - if debug_options.lookup_experiment('disableBoundedTrieMetrics'): - cls.bounded_trie_disabled = True - _LOGGER.info('BoundedTrie metrics are disabled.') - cls._initialized = True - - @classmethod - def reset(cls) -> None: - cls.counter_disabled = False - cls.string_set_disabled = False - cls.bounded_trie_disabled = False - cls._initialized = False - - class Metrics(object): """Lets users create/access metric objects during pipeline execution.""" @staticmethod @@ -237,17 +204,12 @@ class DelegatingCounter(Counter): def __init__( self, metric_name: MetricName, process_wide: bool = False) -> None: super().__init__(metric_name) - self._updater = MetricUpdater( + self.inc = MetricUpdater( # type: ignore[method-assign] cells.CounterCell, metric_name, default_value=1, process_wide=process_wide) - def inc(self, n: int = 1) -> None: - if MetricsFlag.counter_disabled: - return - self._updater(n) - class DelegatingDistribution(Distribution): """Metrics Distribution Delegates functionality to MetricsEnvironment.""" def __init__( @@ -269,23 +231,13 @@ class DelegatingStringSet(StringSet): """Metrics StringSet that Delegates functionality to MetricsEnvironment.""" def __init__(self, metric_name: MetricName) -> None: super().__init__(metric_name) - self._updater = MetricUpdater(cells.StringSetCell, metric_name) - - def add(self, value: str) -> None: - if MetricsFlag.string_set_disabled: - return - self._updater(value) + self.add = MetricUpdater(cells.StringSetCell, metric_name) # type: ignore[method-assign] class DelegatingBoundedTrie(BoundedTrie): """Metrics BoundedTrie that Delegates functionality to MetricsEnvironment.""" def __init__(self, metric_name: MetricName) -> None: super().__init__(metric_name) - self._updater = MetricUpdater(cells.BoundedTrieCell, metric_name) - - def add(self, value) -> None: - if MetricsFlag.bounded_trie_disabled: - return - self._updater(value) + self.add = MetricUpdater(cells.BoundedTrieCell, metric_name) # type: ignore[method-assign] class MetricResults(object): diff --git a/sdks/python/apache_beam/metrics/metric_test.py b/sdks/python/apache_beam/metrics/metric_test.py index 6937236a8aad..ae66200737b5 100644 --- a/sdks/python/apache_beam/metrics/metric_test.py +++ b/sdks/python/apache_beam/metrics/metric_test.py @@ -32,9 +32,7 @@ from apache_beam.metrics.metric import MetricResults from apache_beam.metrics.metric import Metrics from apache_beam.metrics.metric import MetricsFilter -from apache_beam.metrics.metric import MetricsFlag from apache_beam.metrics.metricbase import MetricName -from apache_beam.options.pipeline_options import PipelineOptions from apache_beam.runners.direct.direct_runner import BundleBasedDirectRunner from apache_beam.runners.worker import statesampler from apache_beam.testing.metric_result_matchers import DistributionMatcher @@ -123,108 +121,6 @@ def test_get_namespace_error(self): with self.assertRaises(ValueError): Metrics.get_namespace(object()) - def test_metrics_flag(self): - MetricsFlag.reset() - try: - self.assertFalse(MetricsFlag.counter_disabled) - self.assertFalse(MetricsFlag.string_set_disabled) - self.assertFalse(MetricsFlag.bounded_trie_disabled) - - options = PipelineOptions(['--experiments=disableCounterMetrics']) - MetricsFlag.set_default_pipeline_options(options) - self.assertTrue(MetricsFlag.counter_disabled) - self.assertFalse(MetricsFlag.string_set_disabled) - self.assertFalse(MetricsFlag.bounded_trie_disabled) - - MetricsFlag.reset() - options = PipelineOptions(['--experiments=disableStringSetMetrics']) - MetricsFlag.set_default_pipeline_options(options) - self.assertFalse(MetricsFlag.counter_disabled) - self.assertTrue(MetricsFlag.string_set_disabled) - self.assertFalse(MetricsFlag.bounded_trie_disabled) - - MetricsFlag.reset() - options = PipelineOptions(['--experiments=disableBoundedTrieMetrics']) - MetricsFlag.set_default_pipeline_options(options) - self.assertFalse(MetricsFlag.counter_disabled) - self.assertFalse(MetricsFlag.string_set_disabled) - self.assertTrue(MetricsFlag.bounded_trie_disabled) - - MetricsFlag.reset() - options = PipelineOptions([ - '--experiments=disableCounterMetrics', - '--experiments=disableStringSetMetrics', - '--experiments=disableBoundedTrieMetrics', - ]) - MetricsFlag.set_default_pipeline_options(options) - self.assertTrue(MetricsFlag.counter_disabled) - self.assertTrue(MetricsFlag.string_set_disabled) - self.assertTrue(MetricsFlag.bounded_trie_disabled) - finally: - MetricsFlag.reset() - - def test_disabled_counter_is_noop(self): - sampler = statesampler.StateSampler('', counters.CounterFactory()) - statesampler.set_current_tracker(sampler) - state = sampler.scoped_state( - 'mystep', 'myState', metrics_container=MetricsContainer('mystep')) - MetricsFlag.reset() - try: - sampler.start() - with state: - container = MetricsEnvironment.current_container() - Metrics.counter('ns', 'baseline').inc() - self.assertEqual(len(container.metrics), 1) - options = PipelineOptions(['--experiments=disableCounterMetrics']) - MetricsFlag.set_default_pipeline_options(options) - Metrics.counter('ns', 'after_disable').inc() - Metrics.counter('ns', 'after_disable').inc(5) - Metrics.counter('ns', 'after_disable').dec() - self.assertEqual(len(container.metrics), 1) - finally: - sampler.stop() - MetricsFlag.reset() - - def test_disabled_string_set_is_noop(self): - sampler = statesampler.StateSampler('', counters.CounterFactory()) - statesampler.set_current_tracker(sampler) - state = sampler.scoped_state( - 'mystep', 'myState', metrics_container=MetricsContainer('mystep')) - MetricsFlag.reset() - try: - sampler.start() - with state: - container = MetricsEnvironment.current_container() - Metrics.string_set('ns', 'baseline').add('seed') - self.assertEqual(len(container.metrics), 1) - options = PipelineOptions(['--experiments=disableStringSetMetrics']) - MetricsFlag.set_default_pipeline_options(options) - Metrics.string_set('ns', 'after_disable').add('value') - self.assertEqual(len(container.metrics), 1) - finally: - sampler.stop() - MetricsFlag.reset() - - def test_disabled_bounded_trie_is_noop(self): - sampler = statesampler.StateSampler('', counters.CounterFactory()) - statesampler.set_current_tracker(sampler) - state = sampler.scoped_state( - 'mystep', 'myState', metrics_container=MetricsContainer('mystep')) - MetricsFlag.reset() - try: - sampler.start() - with state: - container = MetricsEnvironment.current_container() - Metrics.bounded_trie('ns', 'baseline').add(['a']) - self.assertEqual(len(container.metrics), 1) - options = PipelineOptions(['--experiments=disableBoundedTrieMetrics']) - MetricsFlag.set_default_pipeline_options(options) - Metrics.bounded_trie('ns', 'after_disable').add(['a', 'b']) - self.assertEqual(len(container.metrics), 1) - finally: - sampler.stop() - MetricsFlag.reset() - def test_counter_empty_name(self): with self.assertRaises(ValueError): Metrics.counter("namespace", "") diff --git a/sdks/python/apache_beam/pipeline.py b/sdks/python/apache_beam/pipeline.py index 594660d9bea9..750868f7443a 100644 --- a/sdks/python/apache_beam/pipeline.py +++ b/sdks/python/apache_beam/pipeline.py @@ -73,7 +73,6 @@ from apache_beam.coders import typecoders from apache_beam.internal import pickler from apache_beam.io.filesystems import FileSystems -from apache_beam.metrics.metric import MetricsFlag from apache_beam.options.pipeline_options import CrossLanguageOptions from apache_beam.options.pipeline_options import DebugOptions from apache_beam.options.pipeline_options import PipelineOptions @@ -193,7 +192,6 @@ def __init__( self._options = PipelineOptions([]) FileSystems.set_options(self._options) - MetricsFlag.set_default_pipeline_options(self._options) if runner is None: runner = self._options.view_as(StandardOptions).runner diff --git a/sdks/python/apache_beam/runners/worker/sdk_worker_main.py b/sdks/python/apache_beam/runners/worker/sdk_worker_main.py index 58beda96d63d..754a631eaf33 100644 --- a/sdks/python/apache_beam/runners/worker/sdk_worker_main.py +++ b/sdks/python/apache_beam/runners/worker/sdk_worker_main.py @@ -33,7 +33,6 @@ from apache_beam.internal import pickler from apache_beam.io import filesystems -from apache_beam.metrics import metric from apache_beam.options.pipeline_options import DebugOptions from apache_beam.options.pipeline_options import GoogleCloudOptions from apache_beam.options.pipeline_options import PipelineOptions @@ -124,7 +123,6 @@ def create_harness(environment, dry_run=False): RuntimeValueProvider.set_runtime_options(pipeline_options_dict) sdk_pipeline_options = PipelineOptions.from_dictionary(pipeline_options_dict) filesystems.FileSystems.set_options(sdk_pipeline_options) - metric.MetricsFlag.set_default_pipeline_options(sdk_pipeline_options) pickle_library = sdk_pipeline_options.view_as(SetupOptions).pickle_library pickler.set_library(pickle_library) diff --git a/sdks/python/apache_beam/yaml/integration_tests.py b/sdks/python/apache_beam/yaml/integration_tests.py index 150c0ca86254..e319a3d3a9bd 100644 --- a/sdks/python/apache_beam/yaml/integration_tests.py +++ b/sdks/python/apache_beam/yaml/integration_tests.py @@ -21,18 +21,52 @@ import copy import glob import itertools +import json import logging import os import random import secrets import sqlite3 import string +import struct import unittest import uuid from datetime import datetime from datetime import timezone import mock + +from apache_beam.coders import Coder +from apache_beam.coders.coder_impl import CoderImpl +from apache_beam.yaml.test_utils.datadog_test_utils import temp_fake_datadog_server + + +class BigEndianIntegerCoderImpl(CoderImpl): + """Coder implementation for big-endian integers used in cross-language tests. + + This is needed because Java's BigEndianIntegerCoder falls back to the generic + 'beam:coders:javasdk:0.1' URN when used in cross-language pipelines, and + Python's FnApiRunner needs to know how to decode it. + """ + def encode_to_stream(self, value, stream, nested): + stream.write(struct.pack('>i', value)) + + def decode_from_stream(self, stream, nested): + return struct.unpack('>i', stream.read(4))[0] + + +class BigEndianIntegerCoder(Coder): + def get_impl(self): + return BigEndianIntegerCoderImpl() + + +# Register the coder with the fallback URN used by the Java SDK for this coder. +# This allows the Python FnApiRunner to handle data sharded by Java transforms +# using BigEndianIntegerCoder in integration tests. +Coder.register_urn( + 'beam:coders:javasdk:0.1', + None, lambda payload, components, context: BigEndianIntegerCoder()) + import psycopg2 import pytds import sqlalchemy diff --git a/sdks/python/apache_beam/yaml/standard_io.yaml b/sdks/python/apache_beam/yaml/standard_io.yaml index 4f679c4a77c4..520c466b600b 100644 --- a/sdks/python/apache_beam/yaml/standard_io.yaml +++ b/sdks/python/apache_beam/yaml/standard_io.yaml @@ -418,9 +418,11 @@ catalog_properties: 'catalog_properties' config_properties: 'config_properties' triggering_frequency_seconds: 'triggering_frequency_seconds' + manifest_file_size: 'manifest_file_size' location_prefix: 'location_prefix' partition_fields: 'partition_fields' table_properties: 'table_properties' + sort_fields: 'sort_fields' error_handling: 'error_handling' underlying_provider: type: beamJar @@ -429,6 +431,27 @@ config: gradle_target: 'sdks:java:io:expansion-service:shadowJar' +#Datadog +- type: renaming + transforms: + 'WriteToDatadog': 'WriteToDatadog' + config: + mappings: + 'WriteToDatadog': + 'url': 'url' + 'api_key': 'api_key' + 'min_batch_count': 'min_batch_count' + 'batch_count': 'batch_count' + 'max_buffer_size': 'max_buffer_size' + 'parallelism': 'parallelism' + 'error_handling': 'error_handling' + underlying_provider: + type: beamJar + transforms: + 'WriteToDatadog': 'beam:schematransform:org.apache.beam:datadog_write:v1' + config: + gradle_target: 'sdks:java:io:expansion-service:shadowJar' + #MongoDB - type: renaming transforms: diff --git a/sdks/python/apache_beam/yaml/test_utils/__init__.py b/sdks/python/apache_beam/yaml/test_utils/__init__.py new file mode 100644 index 000000000000..89aea21adcf0 --- /dev/null +++ b/sdks/python/apache_beam/yaml/test_utils/__init__.py @@ -0,0 +1,18 @@ +# +# Licensed to the Apache Software Foundation (ASF) under one or more +# contributor license agreements. See the NOTICE file distributed with +# this work for additional information regarding copyright ownership. +# The ASF licenses this file to You under the Apache License, Version 2.0 +# (the "License"); you may not use this file except in compliance with +# the License. You may obtain a copy of the License at +# +# http://www.apache.org/licenses/LICENSE-2.0 +# +# Unless required by applicable law or agreed to in writing, software +# distributed under the License is distributed on an "AS IS" BASIS, +# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +# See the License for the specific language governing permissions and +# limitations under the License. +# + +"""Helper utilities for YAML integration tests.""" diff --git a/sdks/python/apache_beam/yaml/test_utils/datadog_test_utils.py b/sdks/python/apache_beam/yaml/test_utils/datadog_test_utils.py new file mode 100644 index 000000000000..4a87a3466689 --- /dev/null +++ b/sdks/python/apache_beam/yaml/test_utils/datadog_test_utils.py @@ -0,0 +1,131 @@ +# +# Licensed to the Apache Software Foundation (ASF) under one or more +# contributor license agreements. See the NOTICE file distributed with +# this work for additional information regarding copyright ownership. +# The ASF licenses this file to You under the Apache License, Version 2.0 +# (the "License"); you may not use this file except in compliance with +# the License. You may obtain a copy of the License at +# +# http://www.apache.org/licenses/LICENSE-2.0 +# +# Unless required by applicable law or agreed to in writing, software +# distributed under the License is distributed on an "AS IS" BASIS, +# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +# See the License for the specific language governing permissions and +# limitations under the License. +# + +"""Helper utilities for Datadog integration tests.""" + +import contextlib +import gzip +import http.server +import io +import json +import logging +import threading + +_LOGGER = logging.getLogger(__name__) + + +class DatadogConnection: + def __init__(self, url, api_key): + self.url = url + self.api_key = api_key + + +class MockDatadogHandler(http.server.BaseHTTPRequestHandler): + def do_POST(self): + if self.path == "/api/v2/logs": + is_chunked = self.headers.get('Transfer-Encoding', + '').lower() == 'chunked' + is_gzip = self.headers.get('Content-Encoding', '').lower() == 'gzip' + content_len = int(self.headers.get('Content-Length', 0)) + + try: + raw_data = b'' + if is_chunked: + while True: + line = self.rfile.readline().strip() + if not line: + break + chunk_len = int(line, 16) + if chunk_len == 0: + self.rfile.readline() # Clear trail + break + raw_data += self.rfile.read(chunk_len) + self.rfile.readline() # Clear trail + elif content_len > 0: + raw_data = self.rfile.read(content_len) + + if raw_data and is_gzip: + with gzip.GzipFile(fileobj=io.BytesIO(raw_data)) as f: + raw_data = f.read() + + if raw_data: + data = json.loads(raw_data) + with self.server.record_lock: + if isinstance(data, list): + self.server.received_records.extend(data) + else: + self.server.received_records.append(data) + except Exception as e: + logging.error("CRITICAL: Failure unpacking mock datadog payload: %s", e) + + self.send_response(200) + self.send_header('Content-Type', 'application/json') + self.end_headers() + self.wfile.write(b'{"status": "ok"}') + else: + self.send_response(404) + self.end_headers() + + def log_message(self, format, *args): + pass + + +@contextlib.contextmanager +def temp_datadog_mock_server(received_records): + server = http.server.ThreadingHTTPServer(('localhost', 0), MockDatadogHandler) + server.received_records = received_records + server.record_lock = threading.Lock() + ip, port = server.server_address + thread = threading.Thread(target=server.serve_forever) + thread.daemon = True + thread.start() + try: + yield f"http://{ip}:{port}" + finally: + server.shutdown() + server.server_close() + thread.join() + + +@contextlib.contextmanager +def temp_fake_datadog_server(expected_records=None): + """Context manager to provide a temporary fake Datadog server for testing. + """ + received = [] + with temp_datadog_mock_server(received) as mock_url: + try: + yield DatadogConnection( + url=mock_url, + api_key="dummy_key_for_testing", + ) + except Exception as err: + logging.error( + "Error interacting with temporary fake Datadog server: %s", err) + raise err + finally: + if expected_records is not None: + + canonicalize = lambda rec: json.dumps(rec, sort_keys=True) + + actual_strs = sorted([canonicalize(r) for r in received]) + expected_strs = sorted([canonicalize(e) for e in expected_records]) + + assert actual_strs == expected_strs, ( + f"Mismatch in recorded Datadog events!\n" + f"Expected: {expected_strs}\n" + f"Actual: {actual_strs}" + ) diff --git a/sdks/python/apache_beam/yaml/tests/datadog.yaml b/sdks/python/apache_beam/yaml/tests/datadog.yaml new file mode 100644 index 000000000000..24485bfbaf4d --- /dev/null +++ b/sdks/python/apache_beam/yaml/tests/datadog.yaml @@ -0,0 +1,63 @@ +# +# Licensed to the Apache Software Foundation (ASF) under one or more +# contributor license agreements. See the NOTICE file distributed with +# this work for additional information regarding copyright ownership. +# The ASF licenses this file to You under the Apache License, Version 2.0 +# (the "License"); you may not use this file except in compliance with +# the License. You may obtain a copy of the License at +# +# http://www.apache.org/licenses/LICENSE-2.0 +# +# Unless required by applicable law or agreed to in writing, software +# distributed under the License is distributed on an "AS IS" BASIS, +# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +# See the License for the specific language governing permissions and +# limitations under the License. +# + +fixtures: + - name: FAKE_DATADOG_SERVER + type: "apache_beam.yaml.test_utils.datadog_test_utils.temp_fake_datadog_server" + config: + expected_records: + - { ddsource: "apache-beam1", ddtags: "test-tags1", hostname: "test-host", service: "test-service", message: "Event for label 11a" } + - { ddsource: "apache-beam2", ddtags: "test-tags2", hostname: "test-host", service: "test-service", message: "Event for label 37a" } + - { ddsource: "apache-beam3", ddtags: "test-tags3", hostname: "test-host", service: "test-service", message: "Event for label 389a" } + - { ddsource: "apache-beam4", ddtags: "test-tags4", hostname: "test-host", service: "test-service", message: "Event for label 3821b" } + +pipelines: + - pipeline: + type: composite + transforms: + - type: Create + config: + elements: + - { ddsource: "apache-beam1", ddtags: "test-tags1", hostname: "test-host", service: "test-service", message: "Event for label 11a" } + - { ddsource: "apache-beam2", ddtags: "test-tags2", hostname: "test-host", service: "test-service", message: "Event for label 37a" } + - { ddsource: "apache-beam3", ddtags: "test-tags3", hostname: "test-host", service: "test-service", message: "Event for label 389a" } + - { ddsource: "apache-beam4", ddtags: "test-tags4", hostname: "test-host", service: "test-service", message: "Event for label 3821b" } + - { ddsource: "apache-beam-broken", hostname: "test-host" } # Triggers mandatory field failure + - type: WriteToDatadog + input: Create + config: + url: "{FAKE_DATADOG_SERVER.url}" + api_key: "{FAKE_DATADOG_SERVER.api_key}" + min_batch_count: 1 + batch_count: 2 + max_buffer_size: 1000 + parallelism: 1 + error_handling: + output: error_output + - type: MapToFields + input: WriteToDatadog.error_output + config: + language: python + fields: + failed_source: "failed_row.ddsource" + - type: AssertEqual + input: MapToFields + config: + elements: + - { failed_source: "apache-beam-broken" } + # Asserting good records is taken care of by the fixture + diff --git a/sdks/python/apache_beam/yaml/tests/iceberg_add_files.yaml b/sdks/python/apache_beam/yaml/tests/iceberg_add_files.yaml index 1e089f19ac0e..278bbe390d86 100644 --- a/sdks/python/apache_beam/yaml/tests/iceberg_add_files.yaml +++ b/sdks/python/apache_beam/yaml/tests/iceberg_add_files.yaml @@ -55,6 +55,9 @@ pipelines: catalog_properties: type: "hadoop" warehouse: "{TEMP_DIR}/dir" + manifest_file_size: 50 + sort_fields: + - "rank desc" # Pipeline 3: Read from Iceberg and verify the contents - pipeline: diff --git a/sdks/standard_external_transforms.yaml b/sdks/standard_external_transforms.yaml index 057c4e3f47d1..b50402a64d54 100644 --- a/sdks/standard_external_transforms.yaml +++ b/sdks/standard_external_transforms.yaml @@ -21,6 +21,41 @@ # # Last updated on: 2026-05-06 +- default_service: sdks:java:io:expansion-service:shadowJar + description: '' + destinations: + python: apache_beam/io + fields: + - description: The Datadog API key. + name: api_key + nullable: false + type: str + - description: The number of events to batch together for each write. + name: batch_count + nullable: true + type: int32 + - description: Specifies how to handle errors. + name: error_handling + nullable: true + type: Row(output=) + - description: The maximum buffer size in bytes. + name: max_buffer_size + nullable: true + type: int64 + - description: The minimum number of events to batch together for each write. + name: min_batch_count + nullable: true + type: int32 + - description: The degree of parallelism for writing. + name: parallelism + nullable: true + type: int32 + - description: The Datadog API URL. + name: url + nullable: false + type: str + identifier: beam:schematransform:org.apache.beam:datadog_write:v1 + name: DatadogWrite - default_service: sdks:java:io:expansion-service:shadowJar description: 'Outputs a PCollection of Beam Rows, each containing a single INT64 number called "value". The count is produced from the given "start" value and