This is an automated email from the ASF dual-hosted git repository.

dependabot[bot] pushed a change to branch 
dependabot/gradle/com.diffplug.spotless-spotless-plugin-gradle-8.9.0
in repository https://gitbox.apache.org/repos/asf/beam.git


    omit d922932d02a Bump com.diffplug.spotless:spotless-plugin-gradle from 
5.6.1 to 8.9.0
     add 41609956aa3 OTEL in kafka. (#39151)
     add ec93d37c602 Bump golang.org/x/oauth2 from 0.7.0 to 0.27.0 in 
/playground/backend (#39524)
     add ad1e278da3e sdks/java: re-enable nullness checks in WithKeys (#39506)
     add 9ae08d5a068 [Solace] Close the HTTP response content stream in 
BrokerResponse (#39404)
     add 9c82e053f02 Avoid output inside try-catch in Java IO (#39124)
     add dd463fa934d Fix flaky AsyncWrapper reset_state test on Python 3.14 
(#39521)
     add 2eb3323c555 Bump github.com/moby/moby/client from 0.5.0 to 0.5.1 in 
/sdks (#39517)
     add 905ade4ed42 Bump github.com/aws/smithy-go from 1.27.4 to 1.27.5 in 
/sdks (#39519)
     add cae3e1749f2 Bump actions/stale from 10 to 11 (#39520)
     add f456c02459d Bump scikit-learn (#39525)
     add 39c0dde7080 [Gemini] Fix pyrefly check bad-typed-dict-key (#39415)
     add 24346193cc8 Add registerSqlOperator() to BeamSqlEnv for custom SQL 
operators (#39432)
     add f258e3e8ae4 OTEL in pubsub (#39150)
     add 55fde074af8 Add the directory with staged files to sys.path and 
document the usage (#39434)
     add fc0d9895080 Replace non-PEP 585 types in watch.py (#39527)
     add f20da8a1c4a Update ruff and pyrefly dependencies (#39531)
     add 8473d90e93f [Java IO] Add ArrowFlight IO connector (#37904)
     add 2ee432b2b61 Support JmsIO SchemaTransform and cross-lang (#39437)
     add 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)
     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 b9e4e2f27a8 Bump docker/login-action from 4.5.2 to 4.6.0 (#39552)
     add f7d8d7c8b58 Preserve partitioning on temp FILE_LOADS tables (#38833)
     add 141804ab568 fix AddFilesIT filter for BigLake (#39533)
     add f5feab8c598 Enhance Python Timestamp to be precision-variable up to 
nanos, and map it to Timestamp logical type (#39537)
     add 7b9380b1c25 Buffer BufferedLogger by newline to avoid log splitting 
(#39288)
     add 03db1a07096 Fix DataflowOutputCounter calculation for 
ValueInEmptyWindows (#39487)
     add b8d77b86055 [IcebergIO] Upgrade Iceberg dependency to 1.11.0 (#39559)
     add 980c11432a1 Clean up legacy references to apitools in GCS I/O (#39433)
     add 9c561e2983e [Iceberg] Make timestamptz return new Timestamp.MICROS 
logical type (#39344)
     add 459b7ee5036 Fix flaky unit test to pass post-submit checks (#39562)
     add 42c693f3c1e Bump github/codeql-action from 4 to 4.37.3 (#39564)
     add 5c58b58cff6 Support IBM MQ for Python JmsIO (#39467)
     add 0619156f8b5 use Java 17 harness (#39570)
     add 2621e9e047c Remove remaining artifacts from dataflow apitools client 
(#39439)
     add 51966762894 update containers (#39575)
     add e4779cf79f2 Fix flaky FileIOTest.testMatchWatchForNewFiles test under 
CI filesystems (#38047)
     add 7f96ee4f2bf Bump google.golang.org/grpc from 1.82.1 to 1.83.0 in /sdks 
(#39585)
     add 93af6786bbb Bump github/codeql-action from 4.37.3 to 4.37.4 (#39586)
     add 8b4c6751963 Support core dump analysis with pystack and gdb. (#39484)
     add a72451d8072 Bump github.com/nats-io/nats-server/v2 from 2.14.3 to 
2.14.4 in /sdks (#39584)
     add 3a02af8146c [Docs] Update Flink version references on the Flink runner 
page (#39212)
     add 539b048eb8d Fix dataframe CSV tests on Windows (#39563)
     add 4d3e1f017e5 Support array-valued schema options in Python (#39583)
     add f0734081aca Feat: new cleaning rule to orphaned subscriptions (#39538)
     add db57e4ae351 [Docs] Add CHANGES entries for Python UnboundedSource and 
Watch (#39579)
     add a6b3399e41d Fix internal test failure after #39487 (#39591)
     add 52f6e46bede Add query_output_schema to ReadFromBigQuery for BEAM_ROW + 
query support (#39160)
     add f3e12fee1b2 [DebeziumIO] Upgrade to Debezium 3.5.2.Final (#39569)
     add 82c6ee4e976 [Python] Bound Watch state with a timestamp cursor (#39090)
     add 789e1d1480c Redistribute - trace propagation (#39590)
     add 83821eb2fd1 update containers (#39596)
     add a89b9f7120e remove gsutil usage (#39448)
     add c4b4bda2a5d [Iceberg CDC] Add Changelog readers and update resolver 
(#38837)
     add c8bacb4dabc Update activemq to 5.19.5 (#39593)
     add 44d4089a67e Potential fix for environment variable built from 
user-controlled sources (#38942)
     add 8ded79b7278 Part 1: Log systemName in DataflowWorkUnitClient, Commit, 
and core worker states (#39561)
     add 8b326b96561 Create span in spanner CDC to start new trace when otel is 
enabled. (#39567)
     add f2c622b3c03 Fix OpenTelemetry dependencies in published POMs (#39608)
     add 24eb1e60d37 Bump github.com/aws/aws-sdk-go-v2/config in /sdks (#39606)
     add 945fcfc3f8f Bump github.com/aws/aws-sdk-go-v2/service/s3 in /sdks 
(#39609)
     add f1125f7c6ca Bump com.gradle.common-custom-user-data-gradle-plugin 
(#39602)
     add 3a0985a07da [Go SDK] Add GroupIntoBatches transform (#19868) (#38220)
     add 77a11347746 Bump github.com/aws/aws-sdk-go-v2/feature/s3/manager in 
/sdks (#39605)
     add 14c61d13821 Bump zizmorcore/zizmor-action from 0.6.1 to 0.6.2 (#39604)
     add 86ca05321d5 Adds the Delta Lake CDC read transforms to the Managed I/O 
API (#39599)
     add aa7f74cea25 mention otel in changes (#39618)
     add 15082278fd1 Support sharded coder for Prism runner cross-lang (#39623)
     add c2c910a70bf Fix Dataflow ValueProvider serialization (#39614)
     add a0f3518d076 Deflake JmsIO tests (#39571)
     add 16f471eee5c [Docs] Add a contributor guide for running Python on a 
local Flink cluster (#39580)
     add 0d27f279d69 Bump com.diffplug.spotless:spotless-plugin-gradle from 
5.6.1 to 8.9.0

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   (d922932d02a)
            \
             N -- N -- N   
refs/heads/dependabot/gradle/com.diffplug.spotless-spotless-plugin-gradle-8.9.0 
(0d27f279d69)

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:
 .../actions/setup-environment-action/action.yml    |  24 +-
 .../IO_Iceberg_Integration_Tests.json              |   2 +-
 .../beam_PostCommit_Java_Delta_IO_Dataflow.json    |   2 +-
 .github/trigger_files/beam_PostCommit_Python.json  |   2 +-
 .../beam_PostCommit_Python_Xlang_Gcp_Direct.json   |   2 +-
 .../beam_PostCommit_Python_Xlang_IO_Dataflow.json  |   2 +-
 .../beam_PostCommit_Python_Xlang_IO_Direct.json    |   2 +-
 ...m_PostCommit_Python_Xlang_Messaging_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/beam_PreCommit_GHA.yml           |  18 +-
 .../workflows/beam_Publish_Beam_SDK_Snapshots.yml  |  17 +-
 .github/workflows/build_release_candidate.yml      |   2 +-
 .github/workflows/build_wheels.yml                 |  12 +-
 .github/workflows/codeql.yml                       |   4 +-
 .github/workflows/finalize_release.yml             |   2 +-
 .github/workflows/python_dependency_tests.yml      |   1 +
 .../workflows/run_rc_validation_go_wordcount.yml   |  12 +-
 .../run_rc_validation_python_mobile_gaming.yml     |   6 +-
 .../workflows/run_rc_validation_python_yaml.yml    |   6 +-
 .github/workflows/stale.yml                        |   2 +-
 .github/workflows/tour_of_beam_backend.yml         |   4 +-
 .../workflows/tour_of_beam_backend_integration.yml |   1 +
 .github/workflows/update_python_dependencies.yml   |   2 +
 .test-infra/dataproc/flink_cluster.sh              |   2 +-
 .test-infra/metrics/build.gradle                   |   1 +
 .test-infra/metrics/influxdb/Dockerfile            |   9 +-
 .test-infra/metrics/influxdb/gsutil/.boto          |  24 -
 .test-infra/metrics/influxdb/gsutil/Dockerfile     |  25 -
 .../kubernetes/beam-influxdb-autobackup.yaml       |   7 +-
 .test-infra/tools/stale_cleaner.py                 |  15 +-
 .test-infra/tools/test_stale_cleaner.py            |  61 +-
 CHANGES.md                                         |  34 +-
 build.gradle.kts                                   |   1 +
 .../org/apache/beam/gradle/BeamModulePlugin.groovy |  15 +-
 contributor-docs/README.md                         |   1 +
 contributor-docs/local-flink-python.md             | 204 +++++
 examples/java/iceberg/build.gradle                 |   4 +-
 .../beam/examples/complete/game/UserScore.java     |   2 +-
 examples/multi-language/README.md                  |   6 +-
 .../beam-ml/automatic_model_refresh.ipynb          |   4 +-
 .../get-started/learn_beam_basics_by_doing.ipynb   |   4 +-
 .../learn_beam_transforms_by_doing.ipynb           |   2 +-
 .../learn_beam_windowing_by_doing.ipynb            |   2 +-
 .../notebooks/get-started/try-apache-beam-go.ipynb |   4 +-
 .../get-started/try-apache-beam-java.ipynb         |   4 +-
 .../notebooks/get-started/try-apache-beam-py.ipynb |   4 +-
 it/iceberg/build.gradle                            |   4 +-
 learning/tour-of-beam/backend/function.go          |   6 +-
 .../backend/integration_tests/client.go            |   8 +-
 .../backend/internal/fs_content/yaml.go            |  10 +-
 .../backend/internal/storage/datastore.go          |  11 +-
 .../tour-of-beam/backend/internal/storage/mock.go  |   8 +-
 .../beam/model/fnexecution/v1/standard_coders.yaml |  31 +
 .../model/pipeline/v1/external_transforms.proto    |   2 +
 .../org/apache/beam/model/pipeline/v1/schema.proto |  13 +
 playground/backend/go.mod                          |   7 +-
 playground/backend/go.sum                          |  10 +-
 .../main/groovy/mobilegaming-java-dataflow.groovy  |  12 +-
 .../groovy/mobilegaming-java-dataflowbom.groovy    |  12 +-
 .../main/groovy/quickstart-java-dataflow.groovy    |   8 +-
 .../python_release_automation_utils.sh             |   6 +-
 .../run_release_candidate_python_quickstart.sh     |   6 +-
 .../core/GroupAlsoByWindowViaWindowSetNewDoFn.java |   2 +-
 .../apache/beam/runners/core/KeyedWorkItem.java    |   9 +
 .../apache/beam/runners/core/ReduceFnRunner.java   |  18 +-
 .../core/SplittableParDoViaKeyedWorkItems.java     |  70 +-
 .../runners/core/SplittableParDoProcessFnTest.java | 283 ++++++-
 runners/google-cloud-dataflow-java/build.gradle    |   4 +-
 .../dataflow/worker/DataflowExecutionContext.java  |   4 +
 .../dataflow/worker/DataflowOutputCounter.java     |  68 +-
 .../dataflow/worker/DataflowWorkUnitClient.java    |  12 +-
 .../worker/IntrinsicMapTaskExecutorFactory.java    |  11 +-
 .../dataflow/worker/MultiKeyBundleOptions.java     | 138 ++++
 .../dataflow/worker/SimpleParDoFnHelpers.java      |   5 +-
 .../dataflow/worker/StreamingDataflowWorker.java   |  36 +-
 .../StreamingGroupAlsoByWindowViaWindowSetFn.java  |   2 +-
 .../worker/StreamingModeExecutionContext.java      |  80 +-
 .../dataflow/worker/WindmillKeyedWorkItem.java     |  28 +-
 .../beam/runners/dataflow/worker/WindmillSink.java |  14 +-
 ...Exception.java => WorkCancellingException.java} |  30 +-
 .../worker/WorkItemCancelledException.java         |  29 +-
 .../logging/DataflowWorkerLoggingHandler.java      |   5 +-
 .../worker/logging/DataflowWorkerLoggingMDC.java   |  15 +-
 .../dataflow/worker/streaming/ActiveWorkState.java |   5 +
 .../streaming/BoundedQueueExecutorWorkHandle.java  |   7 +-
 .../worker/streaming/ComputationState.java         |  11 +
 .../worker/streaming/ComputationWorkExecutor.java  |   5 +-
 .../runners/dataflow/worker/streaming/Work.java    |  60 +-
 .../streaming/harness/MetricsDataProvider.java     |   4 +-
 .../harness/StreamingWorkerStatusReporter.java     |   2 +-
 .../dataflow/worker/util/BoundedQueueExecutor.java |  36 +-
 .../worker/windmill/client/commits/Commit.java     |   8 +-
 .../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    | 185 +++--
 .../processing/failures/WorkFailureProcessor.java  |  73 +-
 .../dataflow/worker/DataflowOutputCounterTest.java | 108 +++
 .../worker/DataflowWorkUnitClientTest.java         |   8 +-
 .../IntrinsicMapTaskExecutorFactoryTest.java       |  14 +-
 .../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 +-
 .../logging/DataflowWorkerLoggingHandlerTest.java  |   4 +-
 .../worker/streaming/ActiveWorkStateTest.java      |  33 +-
 .../streaming/ComputationStateCacheTest.java       |   4 +-
 .../worker/streaming/ComputationStateTest.java     | 114 +++
 .../dataflow/worker/streaming/WorkTest.java        |   4 +-
 .../worker/testing/RestoreDataflowLoggingMDC.java  |   8 +-
 .../testing/RestoreDataflowLoggingMDCTest.java     |  10 +-
 .../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                                        |  64 +-
 sdks/go.sum                                        | 128 +--
 sdks/go/README.md                                  |   2 +-
 sdks/go/container/tools/buffered_logging.go        |  64 +-
 sdks/go/container/tools/buffered_logging_test.go   | 168 +++-
 sdks/go/pkg/beam/coder.go                          |  17 +
 sdks/go/pkg/beam/core/graph/coder/coder.go         | 116 +++
 sdks/go/pkg/beam/core/graph/coder/coder_test.go    |  66 ++
 sdks/go/pkg/beam/core/graph/coder/registry.go      |  60 +-
 .../pkg/beam/core/graph/coder/sharded_key_test.go  |  81 ++
 sdks/go/pkg/beam/core/runtime/exec/coder.go        |  63 ++
 sdks/go/pkg/beam/core/runtime/exec/coder_test.go   |  84 ++
 sdks/go/pkg/beam/core/runtime/graphx/coder.go      |  21 +
 sdks/go/pkg/beam/core/runtime/symbols.go           |  20 +
 sdks/go/pkg/beam/core/typex/class.go               |   4 +-
 sdks/go/pkg/beam/core/typex/fulltype.go            |  23 +
 sdks/go/pkg/beam/core/typex/special.go             |  18 +-
 sdks/go/pkg/beam/core/util/reflectx/call.go        |  31 +
 sdks/go/pkg/beam/pcollection.go                    |  16 +
 sdks/go/pkg/beam/runners/prism/internal/coders.go  |  11 +
 .../pkg/beam/runners/prism/internal/coders_test.go |  16 +
 .../prism/internal/engine/elementmanager.go        |   1 +
 sdks/go/pkg/beam/transforms/batch/batch.go         | 677 +++++++++++++++
 .../pkg/beam/transforms/batch/batch_prism_test.go  | 222 +++++
 .../batch/batch_test.go}                           |  49 +-
 sdks/go/pkg/beam/transforms/batch/doc.go           |  58 ++
 sdks/go/pkg/beam/transforms/batch/size.go          |  88 ++
 sdks/go/pkg/beam/transforms/batch/size_test.go     |  91 ++
 .../test/integration/io/xlang/debezium/debezium.go |   2 +-
 .../integration/io/xlang/debezium/debezium_test.go |   2 +-
 .../apache/beam/sdk/schemas/SchemaTranslation.java |   2 +
 .../beam/sdk/schemas/logicaltypes/Timestamp.java   |   9 +-
 .../apache/beam/sdk/transforms/Redistribute.java   |  23 +-
 .../java/org/apache/beam/sdk/transforms/Reify.java |   4 +-
 .../org/apache/beam/sdk/transforms/WithKeys.java   |  27 +-
 .../java/org/apache/beam/sdk/io/FileIOTest.java    |  22 +-
 .../beam/sdk/schemas/SchemaTranslationTest.java    |   7 +
 .../beam/sdk/extensions/arrow/ArrowConversion.java |  61 ++
 .../sdk/extensions/arrow/ArrowConversionTest.java  |  33 +
 sdks/java/extensions/sql/iceberg/build.gradle      |   4 +-
 .../beam/sdk/extensions/sql/impl/BeamSqlEnv.java   |  13 +
 .../extensions/sql/impl/CalciteQueryPlanner.java   |  10 +-
 .../sdk/extensions/sql/impl/JdbcConnection.java    |  22 +
 .../sql/impl/BeamSqlEnvRegisterOperatorTest.java   | 100 +++
 .../io/{synthetic => arrow-flight}/build.gradle    |  28 +-
 .../beam/sdk/io/arrowflight/ArrowFlightIO.java     | 840 +++++++++++++++++++
 .../beam/sdk/io/arrowflight}/package-info.java     |  13 +-
 .../beam/sdk/io/arrowflight/ArrowFlightIOTest.java | 330 ++++++++
 sdks/java/io/debezium/build.gradle                 |  17 +-
 .../io/debezium/expansion-service/build.gradle     |   6 +-
 sdks/java/io/debezium/src/README.md                |   8 +-
 .../org/apache/beam/io/debezium/DebeziumIO.java    |  15 +-
 .../DebeziumReadSchemaTransformProvider.java       |  96 ++-
 .../beam/io/debezium/KafkaSourceConsumerFn.java    |   5 +
 .../io/debezium/DebeziumIOMySqlConnectorIT.java    |   6 +-
 .../debezium/DebeziumIOPostgresSqlConnectorIT.java |   4 +-
 .../apache/beam/io/debezium/DebeziumIOTest.java    |   3 +-
 .../DebeziumReadSchemaTransformProviderTest.java   | 139 ++++
 .../debezium/DebeziumReadSchemaTransformTest.java  |  29 +-
 sdks/java/io/delta/build.gradle                    |   2 +-
 .../beam/sdk/io/delta/CreateCDCReadTasksDoFn.java  | 296 +++++++
 .../apache/beam/sdk/io/delta/DeltaCDCReadTask.java | 125 +++
 .../beam/sdk/io/delta/DeltaCDCSourceDoFn.java      | 363 ++++++++
 ...va => DeltaCdcReadSchemaTransformProvider.java} |  80 +-
 .../java/org/apache/beam/sdk/io/delta/DeltaIO.java | 192 ++++-
 .../io/delta/DeltaReadSchemaTransformProvider.java |   6 +-
 .../apache/beam/sdk/io/delta/DeltaSourceDoFn.java  |   2 +-
 .../org/apache/beam/sdk/io/delta/DeltaIOIT.java    | 168 +++-
 .../org/apache/beam/sdk/io/delta/DeltaIOTest.java  | 920 +++++++++++++++++++--
 .../beam/sdk/io/delta/DeltaWriteTestUtils.java     | 371 +++++++++
 sdks/java/io/google-cloud-platform/build.gradle    |   2 +
 .../beam/sdk/io/gcp/bigquery/BigQueryIO.java       |  17 +-
 .../beam/sdk/io/gcp/bigquery/BigQueryUtils.java    |   6 +
 .../apache/beam/sdk/io/gcp/healthcare/FhirIO.java  |   6 +-
 .../apache/beam/sdk/io/gcp/healthcare/HL7v2IO.java |  17 +-
 .../apache/beam/sdk/io/gcp/pubsub/PubsubIO.java    | 130 +++
 .../changestreams/action/ActionFactory.java        |   7 +-
 .../action/QueryChangeStreamAction.java            |  44 +-
 .../dofn/ReadChangeStreamPartitionDoFn.java        |   4 +-
 .../sdk/io/gcp/bigquery/BigQueryUtilsTest.java     |  42 +-
 .../action/QueryChangeStreamActionTest.java        |  10 +-
 .../dofn/ReadChangeStreamPartitionDoFnTest.java    |   3 +-
 sdks/java/io/iceberg/build.gradle                  |  10 +-
 .../org/apache/beam/sdk/io/iceberg/IcebergIO.java  |  14 +-
 .../beam/sdk/io/iceberg/IcebergScanConfig.java     | 116 ++-
 .../apache/beam/sdk/io/iceberg/IcebergUtils.java   | 320 ++++---
 .../beam/sdk/io/iceberg/IncrementalScanSource.java |   4 +-
 .../apache/beam/sdk/io/iceberg/ReadFromTasks.java  |   4 +-
 .../org/apache/beam/sdk/io/iceberg/ScanSource.java |   4 +-
 .../apache/beam/sdk/io/iceberg/ScanTaskReader.java |   5 +-
 .../beam/sdk/io/iceberg/cdc/CdcOutputUtils.java    |   6 +-
 .../beam/sdk/io/iceberg/cdc/CdcReadUtils.java      |   2 +
 .../beam/sdk/io/iceberg/cdc/CdcResolver.java       | 180 ++++
 ...ngelogDescriptor.java => CdcRowDescriptor.java} |  57 +-
 .../beam/sdk/io/iceberg/cdc/ChangelogScanner.java  |  13 +-
 .../beam/sdk/io/iceberg/cdc/LocalResolveDoFn.java  | 245 ++++++
 .../beam/sdk/io/iceberg/cdc/OverlapRange.java      | 102 +++
 .../sdk/io/iceberg/cdc/ReadFromChangelogs.java     | 494 +++++++++++
 .../org/apache/beam/sdk/io/iceberg/AddFilesIT.java |  15 +-
 .../beam/sdk/io/iceberg/IcebergIOReadTest.java     |  48 ++
 .../beam/sdk/io/iceberg/IcebergUtilsTest.java      |  53 +-
 .../IcebergWriteSchemaTransformProviderTest.java   |  13 +-
 .../catalog/BigQueryMetastoreCatalogIT.java        |   1 +
 .../io/iceberg/catalog/IcebergCatalogBaseIT.java   |  11 +-
 .../beam/sdk/io/iceberg/cdc/CdcResolverTest.java   | 156 ++++
 .../sdk/io/iceberg/cdc/ChangelogScannerTest.java   |  20 +
 .../sdk/io/iceberg/cdc/LocalResolveDoFnTest.java   | 340 ++++++++
 .../beam/sdk/io/iceberg/cdc/OverlapRangeTest.java  | 161 ++++
 .../sdk/io/iceberg/cdc/ReadFromChangelogsTest.java | 366 ++++++++
 sdks/java/io/jms/build.gradle                      |   6 +
 ...r.java => BeamGenericJmsConnectionFactory.java} |  28 +-
 .../beam/sdk/io/jms/ConnectionConfiguration.java   | 252 ++++++
 .../java/org/apache/beam/sdk/io/jms/JmsIO.java     |  30 +
 .../io/jms/JmsReadSchemaTransformProvider.java}    | 125 ++-
 .../io/jms/JmsWriteSchemaTransformProvider.java    | 191 +++++
 .../sdk/io/jms/ConnectionConfigurationTest.java    | 159 ++++
 .../java/org/apache/beam/sdk/io/jms/JmsIOTest.java | 176 +---
 .../org/apache/beam/sdk/io/jms/JmsLocalTest.java   | 245 ++++++
 .../sdk/io/jms/JmsSchemaTransformProviderTest.java | 253 ++++++
 sdks/java/io/kafka/build.gradle                    |   3 +
 .../java/org/apache/beam/sdk/io/kafka/KafkaIO.java | 156 +++-
 .../KafkaIOReadImplementationCompatibility.java    |   6 +
 .../io/messaging-expansion-service/build.gradle    |   3 +
 .../beam/sdk/io/solace/broker/BrokerResponse.java  |  13 +-
 .../sdk/io/solace/broker/BrokerResponseTest.java   |  66 ++
 .../java/org/apache/beam/sdk/managed/Managed.java  |   5 +
 sdks/python/apache_beam/coders/row_coder_test.py   |  54 ++
 sdks/python/apache_beam/dataframe/io.py            |   2 +-
 sdks/python/apache_beam/dataframe/io_test.py       |  15 +-
 .../online_clustering/clustering_pipeline/setup.py |   2 +-
 .../transforms/elementwise/enrichment_test.py      |   2 +-
 .../io/external/xlang_debeziumio_it_test.py        |   2 +-
 .../apache_beam/io/external/xlang_jmsio_it_test.py | 383 +++++++++
 .../io/external/xlang_mqttio_it_test.py            |  66 +-
 sdks/python/apache_beam/io/filesystemio.py         |   5 +-
 sdks/python/apache_beam/io/gcp/__init__.py         |  20 -
 sdks/python/apache_beam/io/gcp/bigquery.py         |  26 +-
 .../apache_beam/io/gcp/bigquery_file_loads.py      |  83 +-
 .../apache_beam/io/gcp/bigquery_file_loads_test.py | 151 ++++
 .../io/gcp/bigquery_schema_tools_test.py           |  74 +-
 sdks/python/apache_beam/io/gcp/bigquery_test.py    |  53 +-
 .../apache_beam/io/gcp/bigquery_write_it_test.py   |  71 +-
 .../apache_beam/io/gcp/gcsfilesystem_test.py       |   4 +-
 sdks/python/apache_beam/io/watch.py                | 285 +++++--
 sdks/python/apache_beam/io/watch_test.py           | 406 ++++++++-
 .../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 +-
 sdks/python/apache_beam/ml/inference/base.py       |   2 +-
 .../python/apache_beam/options/pipeline_options.py |   5 +
 .../apache_beam/options/pipeline_options_test.py   |  25 +
 sdks/python/apache_beam/portability/common_urns.py |   1 +
 .../runners/dataflow/internal/apiclient.py         |   6 +-
 .../runners/dataflow/internal/apiclient_test.py    |  19 +
 .../runners/dataflow/internal/clients/README.txt   |  11 -
 .../runners/dataflow/internal/clients/__init__.py  |  16 -
 .../apache_beam/runners/dataflow/internal/names.py |   2 +-
 .../dataproc/dataproc_cluster_manager.py           |   3 +-
 .../runners/interactive/interactive_beam_test.py   |   8 +-
 .../runners/interactive/recording_manager.py       | 197 +++--
 .../runners/interactive/recording_manager_test.py  | 191 ++++-
 .../runners/portability/beam_plugins_it_test.py    |  70 ++
 .../apache_beam/runners/worker/sdk_worker_main.py  |  20 +-
 .../runners/worker/sdk_worker_main_test.py         |  49 ++
 .../testing/benchmarks/chicago_taxi/run_chicago.sh |   6 +-
 .../apache_beam/transforms/async_dofn_test.py      |  15 +-
 .../transforms/managed_iceberg_it_test.py          |   4 +-
 sdks/python/apache_beam/typehints/schemas.py       | 146 +++-
 sdks/python/apache_beam/typehints/schemas_test.py  | 110 +++
 sdks/python/apache_beam/utils/timestamp.py         | 327 ++++++--
 sdks/python/apache_beam/utils/timestamp_test.py    | 187 ++++-
 .../databases/debezium.yaml}                       |  41 +-
 sdks/python/apache_beam/yaml/integration_tests.py  |  73 +-
 sdks/python/apache_beam/yaml/standard_io.yaml      |  23 +
 sdks/python/apache_beam/yaml/yaml_io.py            |  14 +-
 sdks/python/apache_beam/yaml/yaml_io_test.py       |  43 +
 sdks/python/build.gradle                           |   4 +-
 .../container/base_image_requirements_manual.txt   |   1 +
 sdks/python/container/boot.go                      |  74 +-
 .../container/ml/py310/base_image_requirements.txt |   1 +
 .../container/ml/py310/gpu_image_requirements.txt  |   1 +
 .../container/ml/py311/base_image_requirements.txt |   1 +
 .../container/ml/py311/gpu_image_requirements.txt  |   1 +
 .../container/ml/py312/base_image_requirements.txt |   1 +
 .../container/ml/py312/gpu_image_requirements.txt  |   1 +
 .../container/ml/py313/base_image_requirements.txt |   1 +
 sdks/python/container/piputil.go                   |  39 +-
 sdks/python/container/profiler.go                  | 303 ++++++-
 sdks/python/container/profiler_test.go             | 125 +++
 .../container/py310/base_image_requirements.txt    |   1 +
 .../container/py311/base_image_requirements.txt    |   1 +
 .../container/py312/base_image_requirements.txt    |   1 +
 .../container/py313/base_image_requirements.txt    |   1 +
 .../container/py314/base_image_requirements.txt    |   1 +
 sdks/python/pyproject.toml                         |   3 +-
 sdks/python/scripts/run_snapshot_publish.sh        |   4 +-
 sdks/python/setup.py                               |  31 +-
 sdks/python/test-suites/dataflow/common.gradle     |   4 +-
 sdks/python/test-suites/direct/build.gradle        |   3 +-
 sdks/python/test-suites/direct/common.gradle       |  21 +-
 sdks/python/tox.ini                                |  11 +-
 sdks/standard_expansion_services.yaml              |   5 +
 sdks/standard_external_transforms.yaml             | 103 ++-
 settings.gradle.kts                                |   3 +-
 website/Dockerfile                                 |   2 +-
 website/build.gradle                               |   5 +-
 .../content/en/blog/apache-hop-with-dataflow.md    |  12 +-
 .../content/en/blog/beam-sql-with-notebooks.md     |   4 +-
 .../site/content/en/documentation/io/connectors.md |  26 +-
 .../site/content/en/documentation/io/managed-io.md | 104 +++
 .../site/content/en/documentation/runners/flink.md |  23 +-
 .../site/content/en/documentation/runners/spark.md |   4 +-
 .../sdks/python-multi-language-pipelines.md        |   2 +-
 .../sdks/python-pipeline-dependencies.md           |  41 +-
 .../site/content/en/get-started/quickstart-java.md |   4 +-
 343 files changed, 17235 insertions(+), 2080 deletions(-)
 delete mode 100644 .test-infra/metrics/influxdb/gsutil/.boto
 delete mode 100644 .test-infra/metrics/influxdb/gsutil/Dockerfile
 create mode 100644 contributor-docs/local-flink-python.md
 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%)
 create mode 100644 
runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/DataflowOutputCounterTest.java
 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/go/pkg/beam/core/graph/coder/sharded_key_test.go
 create mode 100644 sdks/go/pkg/beam/transforms/batch/batch.go
 create mode 100644 sdks/go/pkg/beam/transforms/batch/batch_prism_test.go
 copy sdks/go/pkg/beam/{core/runtime/xlangx/payload_test.go => 
transforms/batch/batch_test.go} (52%)
 create mode 100644 sdks/go/pkg/beam/transforms/batch/doc.go
 create mode 100644 sdks/go/pkg/beam/transforms/batch/size.go
 create mode 100644 sdks/go/pkg/beam/transforms/batch/size_test.go
 create mode 100644 
sdks/java/extensions/sql/src/test/java/org/apache/beam/sdk/extensions/sql/impl/BeamSqlEnvRegisterOperatorTest.java
 copy sdks/java/io/{synthetic => arrow-flight}/build.gradle (65%)
 create mode 100644 
sdks/java/io/arrow-flight/src/main/java/org/apache/beam/sdk/io/arrowflight/ArrowFlightIO.java
 copy 
sdks/java/{extensions/sql/udf/src/main/java/org/apache/beam/sdk/extensions/sql/udf
 => 
io/arrow-flight/src/main/java/org/apache/beam/sdk/io/arrowflight}/package-info.java
 (69%)
 create mode 100644 
sdks/java/io/arrow-flight/src/test/java/org/apache/beam/sdk/io/arrowflight/ArrowFlightIOTest.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/java/io/delta/src/main/java/org/apache/beam/sdk/io/delta/{DeltaReadSchemaTransformProvider.java
 => DeltaCdcReadSchemaTransformProvider.java} (55%)
 create mode 100644 
sdks/java/io/delta/src/test/java/org/apache/beam/sdk/io/delta/DeltaWriteTestUtils.java
 create mode 100644 
sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/cdc/CdcResolver.java
 copy 
sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/cdc/{ChangelogDescriptor.java
 => CdcRowDescriptor.java} (59%)
 create mode 100644 
sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/cdc/LocalResolveDoFn.java
 create mode 100644 
sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/cdc/OverlapRange.java
 create mode 100644 
sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/cdc/ReadFromChangelogs.java
 create mode 100644 
sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/cdc/CdcResolverTest.java
 create mode 100644 
sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/cdc/LocalResolveDoFnTest.java
 create mode 100644 
sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/cdc/OverlapRangeTest.java
 create mode 100644 
sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/cdc/ReadFromChangelogsTest.java
 copy 
sdks/java/io/jms/src/main/java/org/apache/beam/sdk/io/jms/{AutoScaler.java => 
BeamGenericJmsConnectionFactory.java} (52%)
 create mode 100644 
sdks/java/io/jms/src/main/java/org/apache/beam/sdk/io/jms/ConnectionConfiguration.java
 copy 
sdks/java/io/{mqtt/src/main/java/org/apache/beam/sdk/io/mqtt/MqttReadSchemaTransformProvider.java
 => 
jms/src/main/java/org/apache/beam/sdk/io/jms/JmsReadSchemaTransformProvider.java}
 (50%)
 create mode 100644 
sdks/java/io/jms/src/main/java/org/apache/beam/sdk/io/jms/JmsWriteSchemaTransformProvider.java
 create mode 100644 
sdks/java/io/jms/src/test/java/org/apache/beam/sdk/io/jms/ConnectionConfigurationTest.java
 create mode 100644 
sdks/java/io/jms/src/test/java/org/apache/beam/sdk/io/jms/JmsLocalTest.java
 create mode 100644 
sdks/java/io/jms/src/test/java/org/apache/beam/sdk/io/jms/JmsSchemaTransformProviderTest.java
 create mode 100644 
sdks/java/io/solace/src/test/java/org/apache/beam/sdk/io/solace/broker/BrokerResponseTest.java
 create mode 100644 sdks/python/apache_beam/io/external/xlang_jmsio_it_test.py
 delete mode 100644 
sdks/python/apache_beam/runners/dataflow/internal/clients/README.txt
 delete mode 100644 
sdks/python/apache_beam/runners/dataflow/internal/clients/__init__.py
 create mode 100644 
sdks/python/apache_beam/runners/portability/beam_plugins_it_test.py
 copy sdks/python/apache_beam/yaml/{tests/map.yaml => 
extended_tests/databases/debezium.yaml} (55%)
 create mode 100644 sdks/python/container/profiler_test.go

Reply via email to