This is an automated email from the ASF dual-hosted git repository.
github-actions[bot] pushed a change to branch nightly-refs/heads/master
in repository https://gitbox.apache.org/repos/asf/beam.git
from c164b38e924 Add a security model for Beam (#40421)
add 2485e5dcfe9 Fix Grafana alerting
add bf63b3b3aad Merge pull request #40445 from apache/update-grafana
add 81e55c1a997 Adding release-2.77.0-postrelease to protected branches in
.asf.yaml
add 540f68e3de9 fix bigtable schema issue (#40025)
add 6e305537dd1 Update Beam website to release 2.77.0
add 8cae0843559 Update release date
add c0dc5668e75 Merge pull request #40234 from apache/website-2-77
add 5aa894f3190 [Gemini] Prototype OpenAI RemoteInferenceModelHandler
Implementation (#40239)
add bde760f1995 Update managed-io.md for release 2.77.0-RC4.
add a2760a201c5 Remove whitespace
add 697a4da9667 Merge pull request #40369 from
apache/updates_managed_io_docs_2.77.0_rc4
add be4211d1d24 [Go SDK] Recreate a failed data/state channel without
holding the lock (#40415)
add 274905845dd Add examples for cloudsql (#40027)
add e26d6bacb93 Parallelize per-file rewrites in GcsUtilV2 copy and move
(#40376)
add 0291da13e0f use binding (#40455)
add 5cb33da4830 [Java] Remove CT_CONSTRUCTOR_THROW from the global
SpotBugs filter (#40008)
add 5394822d017 Remove usage of a (leaking) shared handles, which can also
accidentally release the last reference hodling a model in
_SharedMap._keepalive when all DoFns have been gc'ed. (#40441)
No new revisions were added by this update.
Summary of changes:
.asf.yaml | 1 +
.../metrics/kubernetes/beamgrafana-deploy.yaml | 4 +
CHANGES.md | 3 +-
.../datatokenization/utils/SchemasUtils.java | 7 +
.../subprocess/utils/CallingSubProcessUtils.java | 7 +
.../examples/subprocess/utils/ExecutableFile.java | 7 +
.../it/clickhouse/ClickHouseResourceManager.java | 7 +
.../it/gcp/bigtable/BigtableResourceManager.java | 7 +
.../beam/it/splunk/SplunkResourceManager.java | 7 +
.../beam/runners/core/StatefulDoFnRunner.java | 7 +
.../construction/SerializablePipelineOptions.java | 7 +
.../beam/runners/core/metrics/BoundedTrieData.java | 8 +-
.../apache/beam/runners/direct/ParDoEvaluator.java | 7 +
.../wrappers/streaming/DoFnOperator.java | 2 +-
.../FlinkStreamingPortablePipelineTranslator.java | 7 +
.../wrappers/streaming/DoFnOperator.java | 2 +-
.../streaming/io/UnboundedSourceWrapper.java | 7 +
.../streaming/state/FlinkStateInternals.java | 9 +-
.../wrappers/streaming/DoFnOperator.java | 2 +-
.../FlinkStreamingPortablePipelineTranslator.java | 7 +
.../translation/functions/FlinkAssignContext.java | 2 +-
.../wrappers/streaming/DoFnOperator.java | 2 +-
.../streaming/KeyedPushedBackElementsHandler.java | 2 +-
.../streaming/io/UnboundedSourceWrapper.java | 7 +
.../io/source/FlinkSourceSplitEnumerator.java | 7 +
.../io/source/LazyFlinkSourceSplitEnumerator.java | 7 +
.../streaming/stableinput/BufferingDoFnRunner.java | 7 +
.../stableinput/KeyedBufferingElementsHandler.java | 7 +
.../streaming/state/FlinkStateInternals.java | 9 +-
.../control/DefaultJobBundleFactory.java | 7 +
.../org/apache/beam/runners/jet/JetRunner.java | 7 +
.../runners/jet/processors/AbstractParDoP.java | 2 +-
.../beam/runners/jet/processors/AssignWindowP.java | 7 +
.../beam/runners/jet/processors/WindowGroupP.java | 2 +-
.../io/BoundedDatasetFactory.java | 2 +-
.../beam/runners/spark/io/SourceDStream.java | 2 +-
.../spark/metrics/SparkBeamMetricSource.java | 7 +
.../beam/runners/spark/metrics/sink/CsvSink.java | 7 +
.../runners/spark/metrics/sink/GraphiteSink.java | 7 +
.../io/BoundedDatasetFactory.java | 2 +-
.../metrics/SparkBeamMetricSource.java | 7 +
.../metrics/sink/CodahaleCsvSink.java | 7 +
.../metrics/sink/CodahaleGraphiteSink.java | 7 +
.../SparkBatchPortablePipelineTranslator.java | 7 +
.../SparkDatasetPortablePipelineTranslator.java | 7 +
.../SparkStreamingPortablePipelineTranslator.java | 7 +
sdks/go/pkg/beam/core/runtime/harness/datamgr.go | 13 +-
.../pkg/beam/core/runtime/harness/datamgr_test.go | 188 +++
.../src/main/resources/beam/spotbugs-filter.xml | 3 -
.../org/apache/beam/sdk/jmh/schemas/RowBundle.java | 7 +
.../org/apache/beam/sdk/coders/CoderProviders.java | 2 +-
.../sdk/fn/data/BeamFnDataOutboundAggregator.java | 7 +
.../org/apache/beam/sdk/io/CompressedSource.java | 7 +
.../java/org/apache/beam/sdk/io/Compression.java | 2 +-
.../schemas/transforms/providers/JavaRowUdf.java | 6 +
.../beam/sdk/schemas/utils/ByteBuddyUtils.java | 9 +-
.../beam/sdk/schemas/utils/ConvertHelpers.java | 7 +
.../java/org/apache/beam/sdk/testing/PAssert.java | 3 +-
.../beam/sdk/testing/SerializableMatchers.java | 4 +-
.../apache/beam/sdk/transforms/AsyncWrapper.java | 7 +
.../beam/sdk/transforms/InferableFunction.java | 7 +
.../org/apache/beam/sdk/transforms/Partition.java | 2 +-
.../org/apache/beam/sdk/transforms/Sample.java | 7 +
.../apache/beam/sdk/transforms/SimpleFunction.java | 7 +
.../java/org/apache/beam/sdk/transforms/View.java | 2 +-
.../java/org/apache/beam/sdk/transforms/Watch.java | 2 +-
.../beam/sdk/transforms/join/CoGbkResult.java | 7 +
.../reflect/ByteBuddyDoFnInvokerFactory.java | 2 +-
.../sdk/transforms/windowing/FixedWindows.java | 7 +
.../sdk/transforms/windowing/SlidingWindows.java | 7 +
.../apache/beam/sdk/util/ExplicitShardedFile.java | 7 +
.../apache/beam/sdk/util/MutationDetectors.java | 2 +-
.../util/construction/DefaultArtifactResolver.java | 7 +
.../util/construction/WriteFilesTranslation.java | 2 +-
.../util/construction/graph/QueryablePipeline.java | 7 +
.../apache/beam/sdk/values/PCollectionViews.java | 13 +
.../org/apache/beam/sdk/values/TypeDescriptor.java | 7 +
.../sdk/expansion/service/ExpansionServer.java | 7 +
.../ExpansionServiceSchemaTransformProvider.java | 7 +
.../service/JavaClassLookupTransformProvider.java | 2 +-
.../sdk/extensions/gcp/options/GcsOptions.java | 31 +-
.../beam/sdk/extensions/gcp/util/GcsUtilV1.java | 8 +-
.../beam/sdk/extensions/gcp/util/GcsUtilV2.java | 251 +++-
.../sdk/extensions/gcp/util/GcsUtilV2Test.java | 75 ++
.../beam/sdk/extensions/ml/AnnotateImages.java | 7 +
.../extensions/sorter/HadoopExternalSorter.java | 2 +-
.../sdk/extensions/sorter/NativeFileSorter.java | 2 +-
.../sql/meta/provider/delta/DeltaTable.java | 2 +-
.../sql/meta/provider/iceberg/IcebergTable.java | 2 +-
.../fn/harness/jmh/ProcessBundleBenchmark.java | 19 +
.../jmh/logging/BeamFnLoggingClientBenchmark.java | 7 +
.../beam/fn/harness/BeamFnDataReadRunner.java | 7 +
.../apache/beam/fn/harness/FnApiDoFnRunner.java | 9 +-
.../SplittablePairWithRestrictionDoFnRunner.java | 7 +
...littableSplitAndSizeRestrictionsDoFnRunner.java | 7 +
...ittableTruncateSizedRestrictionsDoFnRunner.java | 7 +
.../fn/harness/control/ExecutionStateSampler.java | 7 +
.../harness/data/PTransformFunctionRegistry.java | 7 +
.../fn/harness/logging/BeamFnLoggingClient.java | 7 +
.../fn/harness/state/FnApiTimerBundleTracker.java | 7 +
.../beam/fn/harness/status/MemoryMonitor.java | 2 +-
.../components/throttling/AdaptiveThrottler.java | 7 +
.../components/throttling/ReactiveThrottler.java | 7 +
.../apache/beam/io/debezium/SourceRecordJson.java | 7 +
.../apache/beam/sdk/io/delta/SerializableRow.java | 7 +
.../io/elasticsearch/ElasticsearchIOTestUtils.java | 2 +-
.../FileWriteSchemaTransformProvider.java | 2 +-
.../BigtableReadSchemaTransformProvider.java | 2 +-
...gtableSimpleWriteSchemaTransformProviderIT.java | 5 +-
.../beam/sdk/io/gcp/bigtable/BigtableWriteIT.java | 5 +-
.../BigtableWriteSchemaTransformProviderIT.java | 5 +-
.../apache/beam/sdk/io/iceberg/BeamRowWrapper.java | 7 +
.../apache/beam/sdk/io/iceberg/BundleLifter.java | 7 +
.../apache/beam/sdk/io/iceberg/RecordWriter.java | 7 +
.../beam/sdk/io/iceberg/RecordWriterManager.java | 2 +-
.../beam/sdk/io/iceberg/cdc/DeleteReader.java | 7 +
.../io/iceberg/cdc/SerializableChangelogTask.java | 6 +-
.../io/iceberg/cdc/sink/RecordDeltaTaskWriter.java | 7 +
.../IcebergCdcReadSchemaTransformProviderTest.java | 71 ++
.../beam/sdk/io/influxdb/ShardInformation.java | 2 +-
.../io/jdbc/JdbcReadSchemaTransformProvider.java | 8 +-
.../io/jdbc/JdbcWriteSchemaTransformProvider.java | 8 +-
.../ReadFromSqlServerSchemaTransformProvider.java | 7 +
.../WriteToSqlServerSchemaTransformProvider.java | 7 +
.../org/apache/beam/sdk/io/kafka/ProducerSpEL.java | 2 +-
.../apache/beam/sdk/io/kudu/KuduServiceImpl.java | 2 +-
.../org/apache/beam/sdk/io/parquet/ParquetIO.java | 2 +-
.../sdk/io/pulsar/NaiveReadFromPulsarDoFn.java | 2 +-
.../apache/beam/sdk/io/rabbitmq/RabbitMqIO.java | 2 +-
.../beam/sdk/io/rabbitmq/RabbitMqMessage.java | 7 +
.../org/apache/beam/io/requestresponse/Cache.java | 2 +-
.../io/snowflake/data/text/SnowflakeBinary.java | 7 +
.../io/snowflake/data/text/SnowflakeVarchar.java | 7 +
.../beam/sdk/io/solace/broker/BrokerResponse.java | 7 +
.../beam/sdk/io/splunk/CustomX509TrustManager.java | 7 +
.../java/org/apache/beam/sdk/io/xml/XmlSource.java | 2 +-
.../managed/ManagedSchemaTransformProvider.java | 2 +-
.../org/apache/beam/sdk/loadtests/LoadTest.java | 7 +
.../launcher/TransformServiceLauncher.java | 7 +
sdks/python/apache_beam/ml/inference/base.py | 7 +-
sdks/python/apache_beam/ml/inference/base_test.py | 47 +
.../apache_beam/ml/inference/openai_inference.py | 392 +++++++
.../ml/inference/openai_inference_it_test.py | 211 ++++
.../ml/inference/openai_inference_test.py | 624 ++++++++++
.../inference/openai_tests_requirements.txt} | 2 +-
.../yaml/examples/testing/examples_test.py | 2 +
..._to_bigquery.yaml => cloudsql_to_bigquery.yaml} | 23 +-
.../transforms/blueprint/postgres_to_bigquery.yaml | 5 +-
.../extended_tests/e2e/cloudsql_to_bigquery.yaml | 82 ++
sdks/python/apache_beam/yaml/integration_tests.py | 35 +
sdks/python/apache_beam/yaml/tests/bigtable.yaml | 72 +-
sdks/python/apache_beam/yaml/yaml_provider.py | 22 +-
.../apache_beam/yaml/yaml_provider_unit_test.py | 44 +
sdks/python/apache_beam/yaml/yaml_testing.py | 10 +-
sdks/python/pytest.ini | 1 +
sdks/python/test-suites/dataflow/common.gradle | 28 +
website/www/site/config.toml | 2 +-
website/www/site/content/en/blog/beam-2.77.0.md | 74 ++
.../site/content/en/documentation/io/managed-io.md | 1240 +++++++++++++-------
.../www/site/content/en/get-started/downloads.md | 16 +-
160 files changed, 3686 insertions(+), 621 deletions(-)
create mode 100644 sdks/python/apache_beam/ml/inference/openai_inference.py
create mode 100644
sdks/python/apache_beam/ml/inference/openai_inference_it_test.py
create mode 100644
sdks/python/apache_beam/ml/inference/openai_inference_test.py
copy
sdks/python/apache_beam/{transforms/enrichment_handlers/feast_tests_requirements.txt
=> ml/inference/openai_tests_requirements.txt} (98%)
copy
sdks/python/apache_beam/yaml/examples/transforms/blueprint/{spanner_to_bigquery.yaml
=> cloudsql_to_bigquery.yaml} (76%)
create mode 100644
sdks/python/apache_beam/yaml/extended_tests/e2e/cloudsql_to_bigquery.yaml
create mode 100644 website/www/site/content/en/blog/beam-2.77.0.md