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

derrickaw pushed a change to branch 20260722_removeGsutil
in repository https://gitbox.apache.org/repos/asf/beam.git


    omit 6351fa15658 fix influxdb errors
    omit 91f90c1aaf2 add workflow to table
    omit a5a5f2fa6da add comment
    omit 4fced04bee3 remove gsutil usage
     add abe3fc77c3d Fix silent per-char iteration when DoFn returns 
str/bytes/dict (#38429)
     add a15b880c53c Merge pull request #39043: Improve WithKeys coder 
inference context
     add b04ca1e0cee Fix pyrefly check bad-specialization (#39419)
     add c8a689f7c0d Fix pyrefly check bad-argument-count (#39416)
     add da7693c3dd1 [Gemini] Fix pyrefly check invalid-yield (#39417)
     add f1d9c3e122e docs: Document pipeline runner API protos (#39292)
     add adbb11ddb46 Fix gcs endpoint pipeline option wiring (#39435)
     add 5e499388670 KeyCommitTooLarge logging improvements (#39316)
     add 7356d4be850 Replace port with listeners in LocalKafka configuration
     add 4f0044b5de9 Merge pull request #39438 from 
sjvanrossum/testing-kafka-service-fix
     add c7ab90ebad3 Bump google.golang.org/api from 0.289.0 to 0.290.0 in 
/sdks (#39449)
     add b0a7099e1f1 Fix lz4 conflicts (#39454)
     add 1c9fd1bc1b9 Remove isort dependency, migrate to ruff equivalent 
(#39411)
     add c8aa28f7e2e docs: Document Fn Execution API protos (#39293)
     add d688e1b9505 [Java SDK] Warn when ValueState contains collection types 
(#37530)
     add a72d6cba37a [KafkaIO] Use consumer position and lag to estimate end 
offsets (#39285)
     add 0140552d6d5 [#30019] Fix type checking failures around unions with 
nested complex types (#39413)
     add 9833974946c Remove legacy HttpError usage in Bigtable I/O IT test 
(#39441)
     add 9c3f51e56e0 [Python] Add Watch transform with growth_of polling SDF 
(#39023)
     add eb9e7a3dacd Override default fadvise to fix regression from 
gcs-connector v3 upgrade (#39445)
     add a3703d38100 [Dataflow Streaming][Multikey] Support MultiKey commits in 
windmill clients (#38768)
     add b15e5f47d40 Add equal_to_approx matcher for approximate numeric 
assertions (#39443)
     add 7f61891553c Bump cloud.google.com/go/datastore from 1.25.0 to 1.26.0 
in /sdks (#39468)
     add 3cf04be0dff Bump cloud.google.com/go/bigtable from 1.50.0 to 1.51.0 in 
/sdks (#39469)
     add 1b1b70645fb Bump docker/login-action from 4.4.0 to 4.5.0 (#39470)
     add 2232808083b Bump zizmorcore/zizmor-action from 0.6.0 to 0.6.1 (#39471)
     add 60455c3568a Bump setuptools from 78.1.1 to 83.0.0 in 
/.test-infra/mock-apis (#39482)
     add 0ffb108ab79 Revert "Merge pull request #39043: Improve WithKeys coder 
inference context" (#39483)
     add 72d7ae3b414 update tour of beam workflow go version (#39486)
     add d6ab9929a4e Bump docker/login-action from 4.5.0 to 4.5.1 (#39497)
     add cb7329c3273 Use prebuilt Snapshots SDK images for PostCommit Python 
Arm (#39498)
     add 6741628d3af Pin Playground kafka-emulator to kafka-clients 2.4.1 
(#39503)
     add 226ad87e0d9 enable otel context propagation - runner v1 sink, source 
changes, doFnRunner changes for per element propagation (#39152)
     add 6d752f569d9 OTEL in spanner. (#39149)
     add 70a522364e9 (IcebergIO) Support PartitionSpec/SortOrder on dynamic 
table creation via IcebergIO (#39408)
     add faa3ad81a95 Skip IcebergPerformanceTest until next release
     add d9d897a37f1 Merge pull request #39502 from apache/skip-iceberg-perf
     add 214863b2e04 Bump golang.org/x/net from 0.54.0 to 0.55.0 in 
/.test-infra/mock-apis (#39510)
     add de234c72b88 Fix inconsistent AvroSchema type and value for 
SqlType.Date values (#39414)
     add ac5262370d9 Adds documentation for the Delta Lake Read Managed I/O 
(#39495)
     add 901fccd7cf6 sdks/java: remove DefaultAnnotation(NonNull) from 
package-info.java files
     add 1d25d3b75a4 Merge pull request #39463: Remove 
@DefaultAnnotation(NonNull.class) throughout project - it is already default
     add fdac189a032 Fix nullness for PubsubIO
     add 405438b83fe Merge pull request #39444: Fix nullness for PubsubIO
     add dec8d23717a Bump torch (#39512)
     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 89a3d5d0dc1 remove gsutil usage
     add 0fb49e6ffc8 add comment
     add 51c9b1d8923 add workflow to table
     add 7560d432083 fix influxdb errors

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   (6351fa15658)
            \
             N -- N -- N   refs/heads/20260722_removeGsutil (7560d432083)

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 +-
 .github/trigger_files/beam_PostCommit_Python.json  |   2 +-
 ..._Spark.json => beam_PostCommit_Python_Arm.json} |   0
 .../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_Java_IO_Performance_Tests.yml  |   2 +-
 .github/workflows/beam_PostCommit_Python_Arm.yml   |  21 +-
 .../beam_PostCommit_Python_Xlang_IO_Dataflow.yml   |   2 +-
 .../beam_PostCommit_Python_Xlang_IO_Direct.yml     |   2 +-
 .github/workflows/beam_PreCommit_GHA.yml           |   2 +-
 .github/workflows/build_release_candidate.yml      |   2 +-
 .github/workflows/codeql.yml                       |   4 +-
 .github/workflows/finalize_release.yml             |   2 +-
 .github/workflows/python_dependency_tests.yml      |   1 +
 .github/workflows/stale.yml                        |   2 +-
 .github/workflows/tour_of_beam_backend.yml         |  10 +-
 .../workflows/tour_of_beam_backend_integration.yml |   1 +
 .github/workflows/update_python_dependencies.yml   |   2 +
 .test-infra/mock-apis/go.mod                       |   2 +-
 .test-infra/mock-apis/go.sum                       |   4 +-
 .test-infra/mock-apis/poetry.lock                  |  18 +-
 CHANGES.md                                         |  18 +-
 build.gradle.kts                                   |   1 +
 .../org/apache/beam/gradle/BeamModulePlugin.groovy |  13 +-
 examples/java/iceberg/build.gradle                 |   4 +-
 it/build.gradle                                    |   2 +-
 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/fn_execution/v1/beam_fn_api.proto   |  63 +-
 .../model/fn_execution/v1/beam_provision_api.proto |   2 +-
 .../beam/model/fnexecution/v1/standard_coders.yaml |  31 +
 .../beam/model/pipeline/v1/beam_runner_api.proto   |  94 ++-
 .../apache/beam/model/pipeline/v1/endpoints.proto  |   3 +
 .../model/pipeline/v1/external_transforms.proto    |  20 +
 .../apache/beam/model/pipeline/v1/metrics.proto    |  73 +-
 .../org/apache/beam/model/pipeline/v1/schema.proto | 107 +++
 playground/backend/go.mod                          |   7 +-
 playground/backend/go.sum                          |  10 +-
 playground/kafka-emulator/build.gradle             |  11 +
 runners/core-java/build.gradle                     |   1 +
 .../core/GroupAlsoByWindowViaWindowSetNewDoFn.java |   2 +-
 .../apache/beam/runners/core/KeyedWorkItem.java    |   9 +
 .../apache/beam/runners/core/ReduceFnRunner.java   |  18 +-
 .../apache/beam/runners/core/SimpleDoFnRunner.java |  13 +
 .../core/SplittableParDoViaKeyedWorkItems.java     |  70 +-
 .../runners/core/construction/package-info.java    |   4 -
 .../beam/runners/core/metrics/package-info.java    |   4 -
 .../org/apache/beam/runners/core/package-info.java |   4 -
 .../beam/runners/core/triggers/package-info.java   |   4 -
 .../runners/core/SplittableParDoProcessFnTest.java | 283 ++++++-
 .../runners/extensions/metrics/package-info.java   |   4 -
 runners/google-cloud-dataflow-java/build.gradle    |   4 +-
 .../google-cloud-dataflow-java/worker/build.gradle |   1 +
 .../dataflow/worker/DataflowExecutionContext.java  |   4 +
 .../dataflow/worker/DataflowOutputCounter.java     |  68 +-
 .../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      | 127 +++-
 .../dataflow/worker/UngroupedWindmillReader.java   |   7 +-
 .../dataflow/worker/WindmillKeyedWorkItem.java     |  31 +-
 .../WindmillOpenTelemetryContextPropagator.java    |  22 +-
 .../beam/runners/dataflow/worker/WindmillSink.java |  29 +-
 ...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 +-
 .../streaming/KeyCommitTooLargeException.java      |  20 +-
 .../runners/dataflow/worker/streaming/Work.java    |  60 +-
 .../dataflow/worker/util/BoundedQueueExecutor.java |  36 +-
 .../worker/windmill/client/WindmillStream.java     |   5 +
 .../worker/windmill/client/commits/Commit.java     |  72 +-
 .../worker/windmill/client/commits/Commits.java    |   2 +-
 .../windmill/client/commits/CompleteCommit.java    |  30 +-
 .../commits/StreamingApplianceWorkCommitter.java   |  13 +-
 .../commits/StreamingEngineWorkCommitter.java      | 113 ++-
 .../client/getdata/StreamGetDataClient.java        |   5 +-
 .../windmill/client/grpc/GrpcCommitWorkStream.java |  93 ++-
 .../worker/windmill/state/WindmillStateReader.java |   9 +-
 .../processing/ComputationWorkExecutorFactory.java |  23 +-
 .../work/processing/StreamingWorkScheduler.java    | 234 +++---
 .../processing/failures/WorkFailureProcessor.java  |  73 +-
 .../dataflow/worker/DataflowOutputCounterTest.java | 108 +++
 .../dataflow/worker/FakeWindmillServer.java        |  67 +-
 .../IntrinsicMapTaskExecutorFactoryTest.java       |  14 +-
 .../worker/KeyTokenInvalidExceptionTest.java       |  39 -
 .../worker/StreamingDataflowWorkerTest.java        | 383 ++++++++--
 .../worker/StreamingModeExecutionContextTest.java  | 291 ++++++-
 .../worker/WindmillReaderIteratorBaseTest.java     |   4 +-
 .../worker/WindowingWindmillReaderTest.java        |   3 +-
 .../dataflow/worker/WorkerCustomSourcesTest.java   |  18 +-
 .../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       |  15 +-
 .../commits/StreamingEngineWorkCommitterTest.java  | 323 +++++++-
 .../client/grpc/GrpcCommitWorkStreamTest.java      | 273 +++++++
 .../windmill/state/WindmillStateReaderTest.java    |  10 +-
 .../failures/WorkFailureProcessorTest.java         | 101 ++-
 .../work/refresh/ActiveWorkRefresherTest.java      |   4 +-
 .../worker/windmill/src/main/proto/windmill.proto  |  23 +
 sdks/go.mod                                        |  68 +-
 sdks/go.sum                                        | 136 ++--
 sdks/go/container/tools/buffered_logging.go        |  64 +-
 sdks/go/container/tools/buffered_logging_test.go   | 168 +++-
 .../prism/internal/engine/elementmanager.go        |   1 +
 .../org/apache/beam/sdk/jmh/util/package-info.java |   4 -
 .../apache/beam/sdk/annotations/package-info.java  |   4 -
 .../org/apache/beam/sdk/coders/package-info.java   |   4 -
 .../apache/beam/sdk/expansion/package-info.java    |   4 -
 .../beam/sdk/fn/splittabledofn/package-info.java   |   4 -
 .../org/apache/beam/sdk/harness/package-info.java  |   3 -
 .../org/apache/beam/sdk/io/fs/package-info.java    |   4 -
 .../java/org/apache/beam/sdk/io/package-info.java  |   4 -
 .../org/apache/beam/sdk/io/range/package-info.java |   4 -
 .../org/apache/beam/sdk/metrics/package-info.java  |   4 -
 .../java/org/apache/beam/sdk/package-info.java     |   4 -
 .../org/apache/beam/sdk/runners/package-info.java  |   3 -
 .../apache/beam/sdk/schemas/SchemaTranslation.java |   2 +
 .../beam/sdk/schemas/annotations/package-info.java |   4 -
 .../apache/beam/sdk/schemas/io/package-info.java   |   4 -
 .../beam/sdk/schemas/io/payloads/package-info.java |   4 -
 .../beam/sdk/schemas/logicaltypes/Timestamp.java   |   9 +-
 .../sdk/schemas/logicaltypes/package-info.java     |   4 -
 .../org/apache/beam/sdk/schemas/package-info.java  |   4 -
 .../sdk/schemas/parser/generated/package-info.java |   4 -
 .../beam/sdk/schemas/parser/package-info.java      |   4 -
 .../beam/sdk/schemas/transforms/package-info.java  |   4 -
 .../schemas/transforms/providers/package-info.java |   4 -
 .../beam/sdk/schemas/utils/package-info.java       |   4 -
 .../org/apache/beam/sdk/state/package-info.java    |   4 -
 .../org/apache/beam/sdk/testing/package-info.java  |   4 -
 .../org/apache/beam/sdk/transforms/WithKeys.java   |  27 +-
 .../beam/sdk/transforms/display/package-info.java  |   4 -
 .../sdk/transforms/errorhandling/package-info.java |   4 -
 .../beam/sdk/transforms/join/package-info.java     |   4 -
 .../apache/beam/sdk/transforms/package-info.java   |   4 -
 .../sdk/transforms/reflect/DoFnSignatures.java     |  49 ++
 .../beam/sdk/transforms/reflect/package-info.java  |   3 -
 .../transforms/splittabledofn/package-info.java    |   4 -
 .../sdk/transforms/windowing/package-info.java     |   4 -
 .../sdk/util/construction/graph/package-info.java  |   4 -
 .../beam/sdk/util/construction/package-info.java   |   4 -
 .../sdk/values/OpenTelemetryContextPropagator.java |   8 +-
 .../org/apache/beam/sdk/values/WindowedValues.java |   8 +-
 .../org/apache/beam/sdk/values/package-info.java   |   4 -
 .../java/org/apache/beam/sdk/io/FileIOTest.java    |  22 +-
 .../beam/sdk/schemas/SchemaTranslationTest.java    |   7 +
 .../sdk/transforms/reflect/DoFnSignaturesTest.java | 117 +++
 .../beam/sdk/extensions/arrow/ArrowConversion.java |  61 ++
 .../sdk/extensions/arrow/ArrowConversionTest.java  |  33 +
 .../sdk/extensions/avro/coders/package-info.java   |   4 -
 .../beam/sdk/extensions/avro/io/package-info.java  |   4 -
 .../beam/sdk/extensions/avro/package-info.java     |   4 -
 .../avro/schemas/io/payloads/package-info.java     |   4 -
 .../sdk/extensions/avro/schemas/package-info.java  |   4 -
 .../extensions/avro/schemas/utils/AvroUtils.java   |   2 +-
 .../avro/schemas/utils/package-info.java           |   4 -
 .../avro/schemas/utils/AvroUtilsTest.java          |   4 +-
 .../google-cloud-platform-core/build.gradle        |   1 +
 .../sdk/extensions/gcp/options/GcsOptions.java     |  85 ++-
 .../beam/sdk/extensions/gcp/util/GcsUtilV1.java    |  17 +-
 .../sdk/extensions/gcp/GcpCoreApiSurfaceTest.java  |   2 +
 .../sdk/extensions/gcp/options/GcsOptionsTest.java |  47 ++
 .../beam/sdk/extensions/gcp/util/GcsUtilTest.java  |  13 +
 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 +
 .../sdk/extensions/sql/impl/rel/package-info.java  |   4 -
 .../sdk/extensions/sql/impl/rule/package-info.java |   4 -
 .../sql/impl/transform/agg/package-info.java       |   4 -
 .../sql/meta/provider/mongodb/package-info.java    |   4 -
 .../sql/meta/provider/pubsub/package-info.java     |   4 -
 .../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 ++++++++
 .../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 | 157 +++-
 .../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/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 +-
 .../sdk/io/gcp/pubsub/AddTimestampAttribute.java   |  13 +-
 .../beam/sdk/io/gcp/pubsub/ExternalWrite.java      |  15 +-
 .../beam/sdk/io/gcp/pubsub/NestedRowToMessage.java |   8 +-
 .../io/gcp/pubsub/PubSubPayloadTranslation.java    |  46 +-
 .../beam/sdk/io/gcp/pubsub/PubsubClient.java       |  59 +-
 .../beam/sdk/io/gcp/pubsub/PubsubGrpcClient.java   |  47 +-
 .../apache/beam/sdk/io/gcp/pubsub/PubsubIO.java    | 301 ++++++--
 .../beam/sdk/io/gcp/pubsub/PubsubJsonClient.java   | 100 ++-
 .../beam/sdk/io/gcp/pubsub/PubsubMessage.java      |  23 +-
 .../beam/sdk/io/gcp/pubsub/PubsubMessageToRow.java |  42 +-
 ...hAttributesAndMessageIdAndOrderingKeyCoder.java |  19 +-
 ...bsubMessageWithAttributesAndMessageIdCoder.java |  14 +-
 .../pubsub/PubsubMessageWithAttributesCoder.java   |   9 +-
 .../pubsub/PubsubMessageWithMessageIdCoder.java    |   9 +-
 .../pubsub/PubsubReadSchemaTransformProvider.java  |  34 +-
 .../beam/sdk/io/gcp/pubsub/PubsubRowToMessage.java |  58 +-
 .../sdk/io/gcp/pubsub/PubsubSchemaIOProvider.java  |  51 +-
 .../beam/sdk/io/gcp/pubsub/PubsubTestClient.java   | 118 +--
 .../sdk/io/gcp/pubsub/PubsubUnboundedSink.java     |  62 +-
 .../sdk/io/gcp/pubsub/PubsubUnboundedSource.java   | 190 +++--
 .../pubsub/PubsubWriteSchemaTransformProvider.java |  59 +-
 .../apache/beam/sdk/io/gcp/pubsub/TestPubsub.java  |  92 ++-
 .../beam/sdk/io/gcp/pubsub/TestPubsubSignal.java   |  79 +-
 .../beam/sdk/io/gcp/spanner/BatchSpannerRead.java  |  13 +-
 .../sdk/io/gcp/spanner/CreateTransactionFn.java    |   8 +-
 .../beam/sdk/io/gcp/spanner/NaiveSpannerRead.java  |   8 +-
 .../beam/sdk/io/gcp/spanner/ReadSpannerSchema.java |   8 +-
 .../beam/sdk/io/gcp/spanner/SpannerAccessor.java   |  32 +-
 .../beam/sdk/io/gcp/spanner/SpannerConfig.java     |  15 +
 .../apache/beam/sdk/io/gcp/spanner/SpannerIO.java  |  42 +-
 .../gcp/spanner/changestreams/dao/DaoFactory.java  |  16 +-
 .../dofn/CleanUpReadChangeStreamDoFn.java          |   7 +
 .../dofn/DetectNewPartitionsDoFn.java              |   5 +-
 .../spanner/changestreams/dofn/InitializeDoFn.java |   7 +
 .../dofn/ReadChangeStreamPartitionDoFn.java        |   5 +-
 .../apache/beam/sdk/io/gcp/GcpApiSurfaceTest.java  |   1 +
 .../sdk/io/gcp/bigquery/BigQueryUtilsTest.java     |  42 +-
 .../sdk/io/gcp/spanner/SpannerIOWriteTest.java     |  14 +-
 .../dofn/ReadChangeStreamPartitionDoFnTest.java    |   3 +-
 sdks/java/io/hadoop-format/build.gradle            |   1 +
 sdks/java/io/iceberg/build.gradle                  |  10 +-
 .../beam/sdk/io/iceberg/DynamicDestinations.java   |  12 +-
 .../org/apache/beam/sdk/io/iceberg/IcebergIO.java  |  47 +-
 .../beam/sdk/io/iceberg/IcebergScanConfig.java     |   9 +-
 .../apache/beam/sdk/io/iceberg/IcebergUtils.java   |  76 +-
 .../beam/sdk/io/iceberg/IncrementalScanSource.java |   4 +-
 .../io/iceberg/OneTableDynamicDestinations.java    |  33 +-
 .../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 +-
 .../io/iceberg/WritePartitionedRowsToFiles.java    |  13 +-
 .../org/apache/beam/sdk/io/iceberg/AddFilesIT.java |  15 +-
 .../beam/sdk/io/iceberg/IcebergIOReadTest.java     |  48 ++
 .../beam/sdk/io/iceberg/IcebergIOWriteTest.java    |  54 ++
 .../beam/sdk/io/iceberg/IcebergUtilsTest.java      |  43 +-
 .../IcebergWriteSchemaTransformProviderTest.java   |  13 +-
 .../catalog/BigQueryMetastoreCatalogIT.java        |   1 +
 .../io/iceberg/catalog/IcebergCatalogBaseIT.java   |  11 +-
 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 ++++
 .../sdk/io/jms/JmsSchemaTransformProviderTest.java | 253 +++++++
 sdks/java/io/kafka/build.gradle                    |   2 +
 .../java/org/apache/beam/sdk/io/kafka/KafkaIO.java | 156 +++-
 .../KafkaIOReadImplementationCompatibility.java    |   6 +
 .../beam/sdk/io/kafka/ReadFromKafkaDoFn.java       | 114 +--
 .../io/messaging-expansion-service/build.gradle    |   3 +
 .../io/snowflake/crosslanguage/package-info.java   |   4 -
 .../beam/sdk/io/solace/broker/BrokerResponse.java  |  13 +-
 .../sdk/io/solace/broker/BrokerResponseTest.java   |  66 ++
 sdks/java/io/sparkreceiver/3/build.gradle          |   6 +
 .../apache/beam/sdk/testing/kafka/LocalKafka.java  |   2 +-
 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 +-
 .../anomaly_detection_pipeline/setup.py            |   2 +-
 .../online_clustering/clustering_pipeline/setup.py |   2 +-
 ...ytorch_image_classification_with_side_inputs.py |   2 +-
 .../transforms/elementwise/enrichment_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 -
 .../apache_beam/io/gcp/bigquery_change_history.py  |  10 +-
 .../apache_beam/io/gcp/bigquery_file_loads.py      |  83 +-
 .../apache_beam/io/gcp/bigquery_file_loads_test.py | 151 ++++
 sdks/python/apache_beam/io/gcp/bigquery_test.py    |   4 +-
 .../apache_beam/io/gcp/bigquery_write_it_test.py   |  71 +-
 .../apache_beam/io/gcp/bigtableio_it_test.py       |  10 +-
 .../apache_beam/io/gcp/gcsfilesystem_test.py       |   4 +-
 sdks/python/apache_beam/io/watch.py                | 737 ++++++++++++++++++
 sdks/python/apache_beam/io/watch_test.py           | 475 ++++++++++++
 sdks/python/apache_beam/ml/anomaly/specifiable.py  |   2 +-
 .../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       |   5 +-
 .../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 +
 sdks/python/apache_beam/runners/common.pxd         |   1 +
 sdks/python/apache_beam/runners/common.py          |  35 +
 sdks/python/apache_beam/runners/common_test.py     |  58 ++
 .../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 ++
 .../runners/portability/prism_runner.py            |   5 +-
 .../apache_beam/runners/worker/sdk_worker_main.py  |  20 +-
 .../runners/worker/sdk_worker_main_test.py         |  49 ++
 sdks/python/apache_beam/testing/util.py            |  43 ++
 sdks/python/apache_beam/testing/util_test.py       |  44 ++
 .../apache_beam/transforms/async_dofn_test.py      |  15 +-
 .../transforms/managed_iceberg_it_test.py          |   4 +-
 sdks/python/apache_beam/typehints/schemas.py       |  94 ++-
 sdks/python/apache_beam/typehints/schemas_test.py  |  94 +++
 sdks/python/apache_beam/typehints/typehints.py     |  13 +-
 .../python/apache_beam/typehints/typehints_test.py |  18 +
 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/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                         | 111 ++-
 sdks/python/scripts/run_integration_test.sh        |  11 +-
 sdks/python/scripts/run_lint.sh                    |  37 +-
 sdks/python/setup.py                               |  32 +-
 sdks/python/test-suites/dataflow/common.gradle     |  12 +-
 sdks/python/test-suites/direct/common.gradle       |  21 +-
 sdks/python/tox.ini                                |  11 +-
 sdks/standard_expansion_services.yaml              |   4 +
 sdks/standard_external_transforms.yaml             | 103 ++-
 settings.gradle.kts                                |   1 +
 .../site/content/en/documentation/io/connectors.md |  34 +-
 .../site/content/en/documentation/io/managed-io.md |  68 ++
 .../site/content/en/documentation/runners/flink.md |  23 +-
 .../sdks/python-pipeline-dependencies.md           |  41 +-
 374 files changed, 14764 insertions(+), 2543 deletions(-)
 copy .github/trigger_files/{beam_PostCommit_Java_ValidatesRunner_Spark.json => 
beam_PostCommit_Python_Arm.json} (100%)
 create mode 100644 
runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/MultiKeyBundleOptions.java
 copy 
sdks/java/core/src/main/java/org/apache/beam/sdk/values/OpenTelemetryContextPropagator.java
 => 
runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/WindmillOpenTelemetryContextPropagator.java
 (76%)
 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/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/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/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
 create mode 100644 sdks/python/apache_beam/io/watch.py
 create mode 100644 sdks/python/apache_beam/io/watch_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