This is an automated email from the ASF dual-hosted git repository.
derrickaw pushed a change to branch 20260729_addDL2IceIT
in repository https://gitbox.apache.org/repos/asf/beam.git
omit 61e34caf763 add test set name to run step
omit c0e8d885698 post commit trigger
omit 67062b09ff1 add new delta lake to iceberg blueprint
omit 767604656c1 update postcommit yaml tests workflow to support new
blueprints workflow
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)
add 2c735b67eb3 Bump cloud.google.com/go/spanner from 1.93.0 to 1.94.0 in
/sdks (#39550)
add c6e0f5b630e Bump github.com/aws/aws-sdk-go-v2/config in /sdks (#39551)
add 2fd6af992b5 Bump github.com/aws/aws-sdk-go-v2/service/s3 in /sdks
(#39554)
add 61fad4f672f Bump github.com/aws/aws-sdk-go-v2/feature/s3/manager in
/sdks (#39553)
add 525780e7af7 Provide a better error when beam plugin was supplied but
wasn't staged. (#39440)
add f65e0e0b0a8 Updates the Delta Lake source to support reading bounded
change data (#39426)
add 4f934835059 Fix module-level side effects and global random seeding in
univariate ML anomaly tests (#39462)
add bc991d99d17 Set envs
add 8a7c972c61b Merge pull request #39558 from apache/fix-auditkeys
add 621a78cd412 Add Beam YAML support for DebeziumIO
add 01c668b1455 Add Debezium YAML integration test
add a2267554a16 Fix record schema
add 25ed2f8b5b7 With max num of records
add aae48ec3ff0 Add setters
add 2541fc77cb3 Refactoring
add 315ed95a97d Fix python formatter
add b85646bf8df Remove primaryKeyColumns options
add 5178343b652 Fix spotless
add 0de9a676681 Merge pull request #39457 from apache/debezium-io-yaml
add 1517cd1777c update postcommit yaml tests workflow to support new
blueprints workflow
add 7b415cbf34b add new delta lake to iceberg blueprint
add 78a54516591 post commit trigger
add ede83944734 add test set name to run step
This update added new revisions after undoing existing revisions.
That is to say, some revisions that were in the old version of the
branch are not in the new version. This situation occurs
when a user --force pushes a change and generates a repository
containing something like this:
* -- * -- B -- O -- O -- O (61e34caf763)
\
N -- N -- N refs/heads/20260729_addDL2IceIT (ede83944734)
You should already have received notification emails for all of the O
revisions, and so the following emails describe only the N revisions
from the common base, B.
Any revisions marked "omit" are not gone; other references still
refer to them. Any revisions marked "discard" are gone forever.
No new revisions were added by this update.
Summary of changes:
.../beam_PostCommit_Python_Xlang_IO_Dataflow.json | 2 +-
.../beam_PostCommit_Python_Xlang_IO_Direct.json | 2 +-
.../beam_Infrastructure_AuditUnmanagedKeys.yml | 5 +
.../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/python_dependency_tests.yml | 1 +
.../workflows/tour_of_beam_backend_integration.yml | 1 +
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 | 34 +-
sdks/go.sum | 68 +-
sdks/java/extensions/sql/iceberg/build.gradle | 4 +-
.../org/apache/beam/io/debezium/DebeziumIO.java | 16 +-
.../DebeziumReadSchemaTransformProvider.java | 96 ++-
.../DebeziumReadSchemaTransformProviderTest.java | 139 ++++
.../beam/sdk/io/delta/CreateCDCReadTasksDoFn.java | 296 ++++++++
.../apache/beam/sdk/io/delta/DeltaCDCReadTask.java | 125 +++
.../beam/sdk/io/delta/DeltaCDCSourceDoFn.java | 359 +++++++++
.../java/org/apache/beam/sdk/io/delta/DeltaIO.java | 150 +++-
.../io/delta/DeltaReadSchemaTransformProvider.java | 6 +-
.../apache/beam/sdk/io/delta/DeltaSourceDoFn.java | 2 +-
.../org/apache/beam/sdk/io/delta/DeltaIOTest.java | 841 +++++++++++++++++++--
sdks/java/io/iceberg/build.gradle | 4 +-
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 +
.../apache_beam/ml/anomaly/univariate/mean_test.py | 8 +-
.../apache_beam/ml/anomaly/univariate/perf_test.py | 51 +-
.../ml/anomaly/univariate/quantile_test.py | 8 +-
.../ml/anomaly/univariate/stdev_test.py | 8 +-
.../runners/interactive/interactive_beam_test.py | 8 +-
.../runners/interactive/recording_manager.py | 197 +++--
.../runners/interactive/recording_manager_test.py | 188 ++++-
.../apache_beam/runners/worker/sdk_worker_main.py | 13 +-
.../databases/debezium.yaml} | 41 +-
sdks/python/apache_beam/yaml/integration_tests.py | 73 +-
sdks/python/apache_beam/yaml/standard_io.yaml | 23 +
sdks/python/build.gradle | 2 +
sdks/python/setup.py | 27 +-
sdks/python/tox.ini | 11 +-
83 files changed, 4509 insertions(+), 785 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
create mode 100644
sdks/java/io/debezium/src/test/java/org/apache/beam/io/debezium/DebeziumReadSchemaTransformProviderTest.java
create mode 100644
sdks/java/io/delta/src/main/java/org/apache/beam/sdk/io/delta/CreateCDCReadTasksDoFn.java
create mode 100644
sdks/java/io/delta/src/main/java/org/apache/beam/sdk/io/delta/DeltaCDCReadTask.java
create mode 100644
sdks/java/io/delta/src/main/java/org/apache/beam/sdk/io/delta/DeltaCDCSourceDoFn.java
copy sdks/python/apache_beam/yaml/{tests/map.yaml =>
extended_tests/databases/debezium.yaml} (55%)