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 3926b886590 fix golangci-lint issue - tour of beam (#39490)
add 4fd1744935c [Gemini] Fix pyrefly check unexpected-keyword (#39528)
add 509af44c580 Bump docker/login-action from 4.5.1 to 4.5.2 (#39540)
add e9d9086ff01 Bump github.com/aws/aws-sdk-go-v2 from 1.43.0 to 1.43.1 in
/sdks (#39542)
add c4a6f802b74 Bump github.com/aws/aws-sdk-go-v2/service/s3 in /sdks
(#39541)
add 79504e2e6a2 Bump google.golang.org/api from 0.290.0 to 0.291.0 in
/sdks (#39544)
add 4738b16cf40 Bump github.com/aws/aws-sdk-go-v2/credentials in /sdks
(#39543)
add b5c6009c057 Bump github.com/aws/aws-sdk-go-v2/feature/s3/manager in
/sdks (#39545)
add f6c7a06d90a Fix Update Python Dependencies (#39534)
add 3c9f9ce84df [Dataflow Streaming] [Multi Key] MultiKey failure handling
+ Integration (#38919)
add 9019efd790c add workflow_dispatch to other workflow files (#39491)
add 08dc50a3b91 Fix flaky BigQuery persistent retry test (#39539)
add 8a1c19bc9d0 [Python] Convert typing and native generic hints in Watch
coder inference (#39547)
add 6e2044b0083 Bump lower and upper bounds for pyarrow + related
dependencies, remove unnecessary CVE hotfix (#39530)
add eda08d8a3c6 Changes SplittableDoFn to call TruncateRestriction on
drain (#39535)
add ec7004f6f57 [Interactive Beam] Fix caching deadlock, wait race
conditions, and stale graph in notebooks (#39161)
add 58bac320ebd [IcebergIO] Raise Java 17 floor for IcebergIO's Java 11
dependents (#39064)
No new revisions were added by this update.
Summary of changes:
.../actions/setup-environment-action/action.yml | 24 +-
.../beam_PostCommit_Python_Xlang_IO_Dataflow.json | 2 +-
.../beam_PostCommit_Python_Xlang_IO_Direct.json | 2 +-
.../beam_PerformanceTests_xlang_KafkaIO_Python.yml | 2 +-
.../beam_PostCommit_Python_Xlang_IO_Dataflow.yml | 2 +-
.../beam_PostCommit_Python_Xlang_IO_Direct.yml | 2 +-
.github/workflows/build_release_candidate.yml | 2 +-
.github/workflows/finalize_release.yml | 2 +-
.github/workflows/python_dependency_tests.yml | 1 +
.../workflows/tour_of_beam_backend_integration.yml | 1 +
.github/workflows/update_python_dependencies.yml | 2 +
CHANGES.md | 4 +
.../org/apache/beam/gradle/BeamModulePlugin.groovy | 12 +-
examples/java/iceberg/build.gradle | 4 +-
it/iceberg/build.gradle | 4 +-
.../core/SplittableParDoViaKeyedWorkItems.java | 70 +++-
.../runners/core/SplittableParDoProcessFnTest.java | 283 +++++++++++++++--
.../dataflow/worker/DataflowExecutionContext.java | 4 +
.../dataflow/worker/MultiKeyBundleOptions.java | 138 ++++++++
.../dataflow/worker/StreamingDataflowWorker.java | 36 ++-
.../worker/StreamingModeExecutionContext.java | 74 ++++-
...Exception.java => WorkCancellingException.java} | 30 +-
.../worker/WorkItemCancelledException.java | 29 +-
.../dataflow/worker/streaming/ActiveWorkState.java | 5 +
.../streaming/BoundedQueueExecutorWorkHandle.java | 7 +-
.../worker/streaming/ComputationState.java | 7 +
.../worker/streaming/ComputationWorkExecutor.java | 5 +-
.../runners/dataflow/worker/streaming/Work.java | 60 +++-
.../dataflow/worker/util/BoundedQueueExecutor.java | 36 ++-
.../windmill/client/commits/CompleteCommit.java | 15 +-
.../commits/StreamingApplianceWorkCommitter.java | 3 +-
.../commits/StreamingEngineWorkCommitter.java | 29 +-
.../client/getdata/StreamGetDataClient.java | 5 +-
.../worker/windmill/state/WindmillStateReader.java | 9 +-
.../processing/ComputationWorkExecutorFactory.java | 9 +-
.../work/processing/StreamingWorkScheduler.java | 181 +++++++----
.../processing/failures/WorkFailureProcessor.java | 73 ++---
.../worker/KeyTokenInvalidExceptionTest.java | 39 ---
.../worker/StreamingDataflowWorkerTest.java | 352 ++++++++++++++++++---
.../worker/StreamingModeExecutionContextTest.java | 286 ++++++++++++++++-
.../worker/WindmillReaderIteratorBaseTest.java | 4 +-
.../worker/WindowingWindmillReaderTest.java | 3 +-
.../dataflow/worker/WorkerCustomSourcesTest.java | 10 +-
.../worker/streaming/ActiveWorkStateTest.java | 33 +-
.../streaming/ComputationStateCacheTest.java | 4 +-
.../worker/streaming/ComputationStateTest.java | 114 +++++++
.../dataflow/worker/streaming/WorkTest.java | 4 +-
.../worker/util/BoundedQueueExecutorTest.java | 82 +++--
.../worker/util/KeyGroupWorkQueueTest.java | 10 +-
.../StreamingApplianceWorkCommitterTest.java | 4 +-
.../commits/StreamingEngineWorkCommitterTest.java | 64 +++-
.../client/grpc/GrpcCommitWorkStreamTest.java | 188 +++++++++++
.../windmill/state/WindmillStateReaderTest.java | 10 +-
.../failures/WorkFailureProcessorTest.java | 101 +++---
.../work/refresh/ActiveWorkRefresherTest.java | 4 +-
sdks/go.mod | 44 +--
sdks/go.sum | 88 +++---
sdks/java/extensions/sql/iceberg/build.gradle | 4 +-
sdks/java/io/iceberg/build.gradle | 4 +-
.../transforms/elementwise/enrichment_test.py | 2 +-
sdks/python/apache_beam/io/gcp/bigquery_test.py | 4 +-
sdks/python/apache_beam/io/watch.py | 12 +-
sdks/python/apache_beam/io/watch_test.py | 19 ++
sdks/python/apache_beam/ml/inference/base.py | 2 +-
.../runners/interactive/interactive_beam_test.py | 8 +-
.../runners/interactive/recording_manager.py | 197 +++++++-----
.../runners/interactive/recording_manager_test.py | 188 ++++++++++-
sdks/python/pyproject.toml | 1 -
sdks/python/setup.py | 27 +-
sdks/python/tox.ini | 11 +-
70 files changed, 2468 insertions(+), 629 deletions(-)
create mode 100644
runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/MultiKeyBundleOptions.java
rename
runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/{KeyTokenInvalidException.java
=> WorkCancellingException.java} (50%)
delete mode 100644
runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/KeyTokenInvalidExceptionTest.java
create mode 100644
runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/streaming/ComputationStateTest.java