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 2448f5ac7d6 Expand Beam Python heap dump with process/native memory
stats (#39244) (#39466)
add 92e3f353a85 [Dataflow Streaming] [Multi Key] Flush streaming sinks at
key boundaries and bundle completion (#39961)
add b84f49acff9 Bump docker/setup-buildx-action from 4.3.0 to 4.4.1
(#40196)
add ed29fe588dd Bump github/codeql-action from 4.38.0 to 4.38.1 (#40195)
add a4f4d65e125 Bump jlumbroso/free-disk-space from 1.3.1 to 2.0.0 (#40194)
add ea41a61c8b0 Bump cloud.google.com/go/bigquery from 1.83.0 to 1.84.0 in
/sdks (#40193)
add b41b0a28986 Bump google.golang.org/grpc from 1.83.2 to 1.84.0 in /sdks
(#40192)
add fd4f1b4d284 Bump docker/setup-qemu-action from 4.3.0 to 4.4.0 (#40170)
add 1be252e250a Bump github.com/aws/aws-sdk-go-v2/config from 1.33.3 to
1.33.5 in /sdks (#40174)
add e8c6e2654bf Bump docker/build-push-action from 7.3.0 to 7.4.0 (#40169)
add d5a7b098c54 Bump github.com/dustin/go-humanize from 1.0.1 to 1.1.0 in
/sdks (#40191)
add 7d0420cd333 Bump github.com/aws/aws-sdk-go-v2/feature/s3/manager in
/sdks (#40168)
add 26f9c76bd07 Bump org.codehaus.plexus:plexus-xml from 3.0.2 to 4.2.0
(#40135)
add db063b8de70 Fix data race on Prism runner's artifact cache map (#40197)
add 7143520f90d Add Python example for slowly updating global window side
inputs (#40183)
add 4d45c70ecc7 Several small updates to YAML ML examples (#40185)
add bd24aab048f [FileWrites] Implementing evictWritersWhenFull option in
WriteFiles (#40200)
add 9dbf3e716a9 Expose Public API for Iceberg Side Input Cache (#40140)
add 3e5cfce0b5a Fix race condition in test. (#40201)
add daf7d4aacc1 [Python] Remove per-open metadata RPC from BeamBlobReader
(#40182)
add db27d09e5ec [Python] Make read and write buffer size configurable in
gcsio (#40184)
add c280892881b [Iceberg CDC sink] Commit files stage (#40134)
No new revisions were added by this update.
Summary of changes:
.github/workflows/beam_PostCommit_Go.yml | 2 +-
.../workflows/beam_PostCommit_Go_Dataflow_ARM.yml | 2 +-
.../beam_PostCommit_Java_Examples_Dataflow_ARM.yml | 2 +-
.github/workflows/beam_PostCommit_Python_Arm.yml | 4 +-
.../beam_PostCommit_XVR_GoUsingJava_Dataflow.yml | 2 +-
.../workflows/beam_PreCommit_CommunityMetrics.yml | 2 +-
.github/workflows/beam_PreCommit_PythonDocker.yml | 2 +-
.github/workflows/beam_PreCommit_Python_ML.yml | 2 +-
.../workflows/beam_Publish_Beam_SDK_Snapshots.yml | 4 +-
.../workflows/beam_Publish_Python_VLLM_Image.yml | 2 +-
...beam_Python_ValidatesContainer_Dataflow_ARM.yml | 4 +-
.github/workflows/build_release_candidate.yml | 6 +-
.github/workflows/build_runner_image.yml | 6 +-
.github/workflows/build_wheels.yml | 2 +-
.github/workflows/codeql.yml | 4 +-
.../republish_released_docker_containers.yml | 4 +-
CHANGES.md | 1 +
buildSrc/build.gradle.kts | 2 +-
.../runners/dataflow/worker/PubsubDynamicSink.java | 45 +-
.../beam/runners/dataflow/worker/PubsubSink.java | 46 +-
.../dataflow/worker/SizeReportingSinkWrapper.java | 6 +
.../worker/StreamingModeExecutionContext.java | 168 +-
.../beam/runners/dataflow/worker/WindmillSink.java | 27 +-
.../dataflow/worker/util/common/worker/Sink.java | 6 +
.../worker/util/common/worker/WriteOperation.java | 9 +-
.../work/processing/ExecuteWorkResult.java | 57 +
.../work/processing/StreamingWorkScheduler.java | 72 +-
.../dataflow/worker/PubsubDynamicSinkTest.java | 210 +-
.../runners/dataflow/worker/PubsubSinkTest.java | 174 +-
.../worker/StreamingDataflowWorkerTest.java | 467 ++++
.../worker/StreamingModeExecutionContextTest.java | 15 +-
.../dataflow/worker/WorkerCustomSourcesTest.java | 38 +-
.../util/common/worker/WriteOperationTest.java | 59 +
sdks/go.mod | 14 +-
sdks/go.sum | 28 +-
.../runners/prism/internal/jobservices/artifact.go | 7 +-
.../runners/prism/internal/jobservices/server.go | 12 +-
.../java/org/apache/beam/sdk/io/WriteFiles.java | 277 ++-
.../org/apache/beam/sdk/io/WriteFilesTest.java | 173 ++
.../beam/sdk/io/iceberg/BoundedAsyncTasks.java | 5 +-
.../org/apache/beam/sdk/io/iceberg/IcebergIO.java | 125 +-
.../IcebergWriteSchemaTransformProvider.java | 57 +
.../beam/sdk/io/iceberg/TableMetadataDriver.java | 2 +-
.../beam/sdk/io/iceberg/cdc/sink/CommitDeltas.java | 195 ++
.../sdk/io/iceberg/cdc/sink/CommitterMetrics.java | 135 ++
.../sdk/io/iceberg/cdc/sink/OrderedCommitFn.java | 740 +++++++
.../beam/sdk/io/iceberg/BoundedAsyncTasksTest.java | 15 +-
.../iceberg/IcebergIOSideInputTableCacheTest.java | 615 ++++++
.../IcebergWriteSchemaTransformProviderTest.java | 185 ++
.../sdk/io/iceberg/TableMetadataDriverTest.java | 4 +-
.../sdk/io/iceberg/cdc/sink/CommitDeltasTest.java | 2291 ++++++++++++++++++++
.../io/iceberg/cdc/sink/CommitterMetricsTest.java | 124 ++
.../apache_beam/examples/snippets/snippets.py | 60 +
.../apache_beam/examples/snippets/snippets_test.py | 38 +
sdks/python/apache_beam/io/gcp/gcsio.py | 61 +-
.../apache_beam/io/gcp/gcsio_integration_test.py | 9 +-
.../python/apache_beam/options/pipeline_options.py | 35 +
.../ml/enrich_spanner_with_bigquery.yaml | 3 +-
.../ml/log_analysis/ml_preprocessing.yaml | 4 +-
.../examples/transforms/ml/log_analysis/train.py | 4 +-
sdks/python/setup.py | 2 +-
.../en/documentation/patterns/side-inputs.md | 6 +-
62 files changed, 6366 insertions(+), 312 deletions(-)
create mode 100644
runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/work/processing/ExecuteWorkResult.java
create mode 100644
sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/cdc/sink/CommitDeltas.java
create mode 100644
sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/cdc/sink/CommitterMetrics.java
create mode 100644
sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/cdc/sink/OrderedCommitFn.java
create mode 100644
sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/IcebergIOSideInputTableCacheTest.java
create mode 100644
sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/cdc/sink/CommitDeltasTest.java
create mode 100644
sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/cdc/sink/CommitterMetricsTest.java