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


    from ede83944734 add test set name to run step
     add ffed89059a2 change name from blueprints to e2e
     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 7ec43af58eb fix PostCommit Python Dependency
     add 7712ac62dd7 Merge pull request #39637 from 
aIbrahiim/fix-postcommit-pydep-pyarrow
     add 48c9e1f9b38 fix(dataframe): claim remaining restriction range on 
empty/header-only CSV reads (#39581)
     add c60b021cde4 [Iceberg CDC] Finish wiring CDC source together and add 
external API (#39600)
     add cb75d1b773f Updates CHANGES.md to include Delta Lake CDC
     add 8a5d5d2f468 Merge pull request #39644 from 
chamikaramj/update_change_log
     add e793e8abab4 add Timestamp.MICROS for iceberg timestamptz (#39592)
     add e2ae447bef7 Add google-api-python-client to Python 3.14 container 
(#39640)
     add 8bf709c0106 add aws hadoop to DeltaIO (#39617)
     add 02ff2978792 [Dataflow Streaming] Remove finalizeCommits from 
processWork (#39648)
     add 71a7efe31c6 Update CHANGES.md for new release
     add cc822c0d6b1 Moving to 2.77.0-SNAPSHOT on master branch.
     add 92de1e434a2 Bump github/codeql-action from 4.37.4 to 4.37.5 (#39651)
     add eea1e03cf8e [Dataflow Streaming] Remove redundant onKeyTransition call 
(#39652)
     add c7a8f93413f Fix Python 3.14 Container Build, Streamline Installation 
(#39659)
     add 0b91ed1d432 add mention of managed iceberg read breakage (#39660)
     add 111c9c35d8b Bump github/codeql-action from 4.37.5 to 4.37.6 (#39673)
     add 9a03f7222c8 Bump cloud.google.com/go/bigtable from 1.51.0 to 1.52.0 in 
/sdks (#39672)
     add 9f2d498234a Bump github.com/aws/aws-sdk-go-v2/service/s3 in /sdks 
(#39671)
     add 57e1f8eb651 Bump github.com/aws/aws-sdk-go-v2/feature/s3/manager in 
/sdks (#39670)
     add 1a21c6183a8 fix iceberg CDC test (#39675)
     add f9ca09b3183 [Spark] Support splittable DoFn self-checkpointing in 
portable batch (#39331)
     add 837590e5426 Fix Row.toString NPE on a null nested inside an array, map 
or row (#39587)
     add d0cff01a766 [BigQueryIO] Parallelize schema update integration tests 
(#39622)
     add 1e4c093445a [examples] Atomically publish subprocess executables 
(#39621)
     add 1fe64b99e8d Bump h2 from 4.3.0 to 4.4.1 in 
/sdks/python/container/py314 (#39677)
     add 2e48a718712 Add ml and interactive extras to quickstart-py doc (#39679)
     add d3d6e484a74 add closing dependabot step (#39649)
     add 0aacd0375f0 Add helpers to interact with pipeline options in boot 
entrypoints (#39595)
     add 1318bfe2039 Fix runner compatibility matrix (#39682)
     add 559d22c498b fix ensurepip bundled pip cleanup for Python 3.12+ 
containers (#39683)
     add d97899b7ab8 Add Sample.Any to the Python SDK to match Java's 
Sample.any (#39442)
     add 0b40089ffd1 Fix mobile gaming release validation background process 
cleanup for Java 21 compatibility (#39658)
     add d6a865d2e4d Persist credentials for build_release_candidate.yml
     add f0da6f36657 Enable OpenTelemetry stiching with Logs for Dataflow 
worker, both for direct logging and file based (#39625)
     add 8d24582beab (IcebergIO) document writeProperties param more clearly 
(#39645)
     add 367f46d4014 [Python] Create temporary dataset with a 24 hour ttl. 
(#39615)
     add 4731dbc5a38 Fix RequestResponseIO parseAndThrow to preserve retryable 
exception types (#37342)
     add c91aa1c1d89 feat: add MongoDB driver handshake metadata for Java-based 
client connections (#39504)
     add eab1bceed4f Bump dorny/paths-filter from 4.0.1 to 4.0.3 (#39693)
     add f55c10b33e5 Pin grpcio-tools==1.78.0 for python 3.14
     add 67d6402cfb3 Merge pull request #39715 from apache/fix-python314
     add f353f12b543 [Dataflow Streaming] Mark worker as unhealthy in presence 
of stuck commits (#39666)
     add 680229c7fa7 normalize io.gcp.DicomSearch
     add 569933cc905 Merge pull request #39655 from 
aIbrahiim/yaml-normalize-dicom-search
     add 53b03f6329a [Python] Deflake TextIO footer test (#39668)
     add 33b40fbf5c4 Bump github.com/aws/aws-sdk-go-v2/config in /sdks (#39692)
     add 7dda2cdb72d Enable Apache Iceberg REST Metrics Reporting for Lakehouse 
(#39650)
     add ce45298a609 Fixes to delta CDC read (#39713)
     add e344ec03fb7 Python timestamp fixes. (#39722)
     add 524036fb523 bump FnAPI container to beam-master-20260811 (#39721)
     add deb5e274778 Bump js-yaml from 3.15.0 to 3.15.1 in /website/www (#39678)
     add 22b73bb4fdb Bump github.com/aws/aws-sdk-go-v2/feature/s3/manager in 
/sdks (#39695)
     add a0e27149ee1 [KafkaIO] Remove beam_fn_api requirement for dynamic reads 
(#39735)
     add 369409ea492 [IcebergIO] Serialize using json partition (#39705)
     add 6ddc7fec3b8 Log the System name in more places instead of the 
computationId (#39665)
     add e818a0c4ad5 Feat: implementing active cleanup of orphaned 
subscriptions for the `taxirides` topic. (#39728)
     add dd896e2b239 [Dataflow Streaming] [Multi Key] Drop failed work in 
BoundedQueueExecutor::pollWork (#38920)
     add 630c751b23d Restore go CoGBK load test parameter (#39753)
     add d872d0a5e7f [GSoC 2026] Requesting permissions for the 
TestPubSubContext cleanup handler tests (#39757)
     add 078798646d6 Bump github.com/testcontainers/testcontainers-go in /sdks 
(#39740)
     add befa812ecc5 Fix: Removing users who do not have a valid Google account 
from the list. (#39769)
     add d507f1bb3e2 Bump github/codeql-action from 4.37.6 to 4.37.7 (#39775)
     add 649a9004e25 Bump google.golang.org/api from 0.291.0 to 0.293.0 in 
/sdks (#39777)
     add 9ea7c97b611 Bump cloud.google.com/go/bigquery from 1.79.0 to 1.80.0 in 
/sdks (#39776)
     add 151318591e1 Bump github.com/aws/aws-sdk-go-v2/config in /sdks (#39774)
     add cfc35b76610 Bump cryptography from 48.0.1 to 50.0.0 in Python SDK
     add a5f5f49c1e9 Merge pull request #39756: Bump cryptography from 48.0.1 
to 50.0.0 in Python SDK
     add e0336ce4dae Bump github.com/aws/aws-sdk-go-v2/service/s3 in /sdks 
(#39747)
     add bd4f9df02c6 fix merge conflict

No new revisions were added by this update.

Summary of changes:
 .asf.yaml                                          |    1 +
 .../IO_Iceberg_Integration_Tests_Dataflow.json     |    2 +-
 .../beam_PostCommit_Java_Delta_IO_Dataflow.json    |    2 +-
 .../beam_PostCommit_Java_PVR_Spark3_Streaming.json |    2 +-
 .../beam_PostCommit_Java_PVR_Spark_Batch.json      |    2 +-
 ...beam_PostCommit_Java_ValidatesRunner_Spark.json |    2 +-
 .../beam_PostCommit_Python_Dependency.json         |    4 +-
 ...am_PostCommit_Python_ValidatesRunner_Spark.json |    3 +-
 .../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 +-
 .github/trigger_files/beam_PostCommit_SQL.json     |    2 +-
 .../beam_PostCommit_Yaml_Xlang_Direct.json         |    2 +-
 .github/trigger_files/beam_PreCommit_SQL.json      |    2 +-
 .../beam_PostCommit_Java_Delta_IO_Dataflow.yml     |    9 +
 .../beam_PostCommit_Yaml_Xlang_Direct.yml          |    2 +-
 .../workflows/beam_PostRelease_NightlySnapshot.yml |    9 +-
 .github/workflows/beam_PreCommit_GHA.yml           |   18 +-
 .../workflows/beam_Publish_Beam_SDK_Snapshots.yml  |   17 +-
 .github/workflows/build_release_candidate.yml      |   14 +-
 .github/workflows/build_wheels.yml                 |   12 +-
 .github/workflows/codeql.yml                       |    6 +-
 .github/workflows/cut_release_branch.yml           |    4 +-
 .github/workflows/finalize_release.yml             |    2 +-
 .../go_CoGBK_Flink_Batch_MultipleKey.txt           |    4 +-
 .../go_CoGBK_Flink_Batch_Reiteration_10KB.txt      |    4 +-
 .../go_CoGBK_Flink_Batch_Reiteration_2MB.txt       |    4 +-
 .../go_GBK_Flink_Batch_100kb.txt                   |    2 +-
 .../go_GBK_Flink_Batch_Fanout_4.txt                |    2 +-
 .../go_GBK_Flink_Batch_Fanout_8.txt                |    2 +-
 .../go_GBK_Flink_Batch_Reiteration_10KB.txt        |    2 +-
 .../workflows/run_rc_validation_go_wordcount.yml   |   12 +-
 .../run_rc_validation_python_mobile_gaming.yml     |    6 +-
 .../workflows/run_rc_validation_python_yaml.yml    |    6 +-
 .test-infra/dataproc/flink_cluster.sh              |    8 +-
 .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                 |   17 +-
 .test-infra/tools/test_stale_cleaner.py            |   61 +-
 CHANGES.md                                         |   83 +-
 .../org/apache/beam/gradle/BeamModulePlugin.groovy |    4 +-
 contributor-docs/README.md                         |    1 +
 contributor-docs/local-flink-python.md             |  204 ++++
 .../beam/examples/complete/game/UserScore.java     |    2 +-
 .../beam/examples/subprocess/utils/FileUtils.java  |   30 +-
 .../examples/subprocess/utils/FileUtilsTest.java   |   80 ++
 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 +-
 gradle.properties                                  |    4 +-
 infra/iam/users.yml                                |   13 +-
 it/mongodb/build.gradle                            |    1 +
 .../beam/it/mongodb/MongoDBResourceManager.java    |   11 +-
 .../it/mongodb/MongoDBResourceManagerTest.java     |    5 +
 .../beam/model/fnexecution/v1/standard_coders.yaml |   31 +
 .../model/pipeline/v1/external_transforms.proto    |    2 +
 .../org/apache/beam/model/pipeline/v1/schema.proto |   13 +
 release/build.gradle.kts                           |    2 +-
 release/src/main/groovy/TestScripts.groovy         |  156 ++-
 .../main/groovy/mobilegaming-java-dataflow.groovy  |   87 +-
 .../groovy/mobilegaming-java-dataflowbom.groovy    |   12 +-
 .../main/groovy/mobilegaming-java-direct.groovy    |   73 +-
 .../main/groovy/quickstart-java-dataflow.groovy    |    8 +-
 .../main/groovy/quickstart-java-flinklocal.groovy  |    4 +-
 .../src/main/groovy/quickstart-java-spark.groovy   |   20 +-
 .../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 +-
 runners/google-cloud-dataflow-java/build.gradle    |    4 +-
 .../options/DataflowStreamingPipelineOptions.java  |    2 +-
 .../google-cloud-dataflow-java/worker/build.gradle |    1 +
 .../dataflow/worker/DataflowOutputCounter.java     |   68 +-
 .../dataflow/worker/DataflowWorkUnitClient.java    |   12 +-
 .../worker/IntrinsicMapTaskExecutorFactory.java    |   11 +-
 .../dataflow/worker/SimpleParDoFnHelpers.java      |    5 +-
 .../dataflow/worker/StreamingDataflowWorker.java   |   34 +-
 .../StreamingGroupAlsoByWindowViaWindowSetFn.java  |    2 +-
 .../worker/StreamingModeExecutionContext.java      |   24 +-
 .../dataflow/worker/WindmillKeyedWorkItem.java     |   28 +-
 .../beam/runners/dataflow/worker/WindmillSink.java |   14 +-
 .../logging/DataflowWorkerLoggingHandler.java      |   35 +-
 .../logging/DataflowWorkerLoggingInitializer.java  |    4 +
 .../worker/logging/DataflowWorkerLoggingMDC.java   |   15 +-
 .../dataflow/worker/streaming/ActiveWorkState.java |   46 +-
 .../worker/streaming/ComputationState.java         |    9 +-
 .../worker/streaming/ComputationWorkExecutor.java  |    6 +-
 .../worker/streaming/FailedWorkHandler.java        |    8 +-
 .../streaming/harness/MetricsDataProvider.java     |    4 +-
 .../harness/StreamingWorkerStatusReporter.java     |    2 +-
 .../dataflow/worker/util/BoundedQueueExecutor.java |   26 +-
 .../worker/windmill/client/commits/Commit.java     |    8 +-
 .../work/processing/StreamingWorkScheduler.java    |   38 +-
 .../processing/failures/WorkFailureProcessor.java  |   31 +-
 .../windmill/work/refresh/ActiveWorkRefresher.java |   21 +-
 .../dataflow/worker/DataflowOutputCounterTest.java |  108 ++
 .../worker/DataflowWorkUnitClientTest.java         |    8 +-
 .../IntrinsicMapTaskExecutorFactoryTest.java       |   14 +-
 .../worker/StreamingDataflowWorkerTest.java        |  192 +++-
 .../worker/StreamingModeExecutionContextTest.java  |   97 +-
 .../dataflow/worker/WorkerCustomSourcesTest.java   |    9 +-
 .../logging/DataflowWorkerLoggingHandlerTest.java  |  107 +-
 .../worker/streaming/ActiveWorkStateTest.java      |   42 +-
 .../worker/testing/RestoreDataflowLoggingMDC.java  |    8 +-
 .../testing/RestoreDataflowLoggingMDCTest.java     |   10 +-
 .../worker/util/BoundedQueueExecutorTest.java      |  129 ++-
 .../failures/WorkFailureProcessorTest.java         |   49 +-
 .../work/refresh/ActiveWorkRefresherTest.java      |   69 --
 runners/spark/job-server/spark_job_server.gradle   |    5 +-
 .../SparkBatchPortablePipelineTranslator.java      |    8 +-
 .../translation/SparkExecutableStageFunction.java  |  172 +++-
 .../SparkStreamingPortablePipelineTranslator.java  |    4 +-
 .../SparkExecutableStageFunctionTest.java          |  138 ++-
 scripts/beam-sql.sh                                |    2 +-
 scripts/ci/pr-bot/processNewPrs.ts                 |   20 +
 scripts/ci/pr-bot/shared/githubUtils.ts            |   21 +
 sdks/go.mod                                        |   76 +-
 sdks/go.sum                                        |  152 +--
 sdks/go/README.md                                  |    2 +-
 sdks/go/container/boot.go                          |   40 +-
 sdks/go/container/boot_test.go                     |   17 +-
 sdks/go/container/tools/buffered_logging.go        |   64 +-
 sdks/go/container/tools/buffered_logging_test.go   |  168 +++-
 sdks/go/container/tools/pipeline_options.go        |  218 ++++
 sdks/go/container/tools/pipeline_options_test.go   |  243 +++++
 sdks/go/pkg/beam/artifact/options.go               |   48 -
 sdks/go/pkg/beam/artifact/options_test.go          |   78 --
 sdks/go/pkg/beam/coder.go                          |   17 +
 sdks/go/pkg/beam/core/core.go                      |    2 +-
 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 +
 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 +-
 sdks/java/container/boot.go                        |    4 +-
 .../apache/beam/sdk/options/SdkHarnessOptions.java |    7 +
 .../apache/beam/sdk/schemas/SchemaTranslation.java |    2 +
 .../org/apache/beam/sdk/schemas/SchemaUtils.java   |    7 +
 .../beam/sdk/schemas/logicaltypes/Timestamp.java   |    9 +-
 .../apache/beam/sdk/transforms/Redistribute.java   |   23 +-
 .../java/org/apache/beam/sdk/transforms/Reify.java |    4 +-
 .../java/org/apache/beam/sdk/io/FileIOTest.java    |   22 +-
 .../beam/sdk/schemas/SchemaTranslationTest.java    |    7 +
 .../apache/beam/sdk/schemas/SchemaUtilsTest.java   |   43 +
 .../service/WindowIntoTransformProvider.java       |    1 +
 .../provider/iceberg/BeamSqlCliIcebergTest.java    |    6 +-
 .../meta/provider/iceberg/IcebergReadWriteIT.java  |    3 -
 .../sdk/extensions/sql/impl/rel/BeamCalcRel.java   |   40 +
 .../extensions/sql/impl/utils/CalciteUtils.java    |    8 +-
 .../sdk/extensions/sql/BeamComplexTypeTest.java    |   32 +
 .../extensions/sql/impl/rel/BeamCalcRelTest.java   |   20 +
 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    |   31 +-
 .../beam/io/debezium/KafkaSourceConsumerFn.java    |    5 +
 .../io/debezium/DebeziumIOMySqlConnectorIT.java    |    6 +-
 .../debezium/DebeziumIOPostgresSqlConnectorIT.java |    4 +-
 .../apache/beam/io/debezium/DebeziumIOTest.java    |    3 +-
 .../debezium/DebeziumReadSchemaTransformTest.java  |   29 +-
 sdks/java/io/delta/build.gradle                    |   21 +-
 .../beam/sdk/io/delta/DeltaCDCSourceDoFn.java      |   32 +-
 ...va => DeltaCdcReadSchemaTransformProvider.java} |   78 +-
 .../java/org/apache/beam/sdk/io/delta/DeltaIO.java |   52 +-
 .../org/apache/beam/sdk/io/delta/DeltaIOIT.java    |  168 +++-
 .../io/delta/{DeltaIOIT.java => DeltaIOS3IT.java}  |  132 +--
 .../org/apache/beam/sdk/io/delta/DeltaIOTest.java  | 1056 ++++++++++++--------
 .../beam/sdk/io/delta/DeltaWriteTestUtils.java     |  371 +++++++
 sdks/java/io/expansion-service/build.gradle        |    1 +
 sdks/java/io/google-cloud-platform/build.gradle    |    4 +-
 .../beam/sdk/io/gcp/bigquery/BigQueryUtils.java    |    6 +
 .../changestreams/action/ActionFactory.java        |    7 +-
 .../action/QueryChangeStreamAction.java            |   44 +-
 .../dofn/ReadChangeStreamPartitionDoFn.java        |    4 +-
 .../sdk/io/gcp/bigquery/BigQueryUtilsTest.java     |   42 +-
 ....java => StorageApiSinkSchemaUpdateITBase.java} |   61 +-
 ...torageApiSinkSchemaUpdateWithInputSchemaIT.java |   50 +
 ...ageApiSinkSchemaUpdateWithoutInputSchemaIT.java |   51 +
 .../action/QueryChangeStreamActionTest.java        |   10 +-
 .../dofn/ReadChangeStreamPartitionDoFnTest.java    |    3 +-
 sdks/java/io/iceberg/build.gradle                  |    6 +-
 .../org/apache/beam/sdk/io/iceberg/AddFiles.java   |    2 +-
 .../IcebergCdcReadSchemaTransformProvider.java     |   31 +-
 .../org/apache/beam/sdk/io/iceberg/IcebergIO.java  |   69 +-
 .../beam/sdk/io/iceberg/IcebergScanConfig.java     |  115 ++-
 .../apache/beam/sdk/io/iceberg/IcebergUtils.java   |  320 +++---
 .../beam/sdk/io/iceberg/IncrementalScanSource.java |  100 --
 .../apache/beam/sdk/io/iceberg/PartitionUtils.java |   10 +-
 .../apache/beam/sdk/io/iceberg/ReadFromTasks.java  |   96 --
 .../beam/sdk/io/iceberg/RecordWriterManager.java   |    9 +-
 .../org/apache/beam/sdk/io/iceberg/ScanSource.java |    4 +-
 .../apache/beam/sdk/io/iceberg/ScanTaskReader.java |    5 +-
 .../beam/sdk/io/iceberg/SerializableDataFile.java  |   88 +-
 .../beam/sdk/io/iceberg/WatchForSnapshots.java     |  190 ----
 .../io/iceberg/WritePartitionedRowsToFiles.java    |    3 +-
 .../sdk/io/iceberg/cdc/ApplyWatermarkColumn.java   |   99 ++
 .../beam/sdk/io/iceberg/cdc/CdcOutputUtils.java    |   25 +-
 .../beam/sdk/io/iceberg/cdc/CdcReadUtils.java      |    4 +-
 .../beam/sdk/io/iceberg/cdc/CdcResolver.java       |  191 ++++
 ...ngelogDescriptor.java => CdcRowDescriptor.java} |   57 +-
 .../beam/sdk/io/iceberg/cdc/ChangelogScanner.java  |   13 +-
 .../io/iceberg/cdc/IncrementalChangelogSource.java |  211 ++++
 .../beam/sdk/io/iceberg/cdc/LocalResolveDoFn.java  |  249 +++++
 .../beam/sdk/io/iceberg/cdc/OverlapRange.java      |  102 ++
 .../sdk/io/iceberg/cdc/ReadFromChangelogs.java     |  499 +++++++++
 .../beam/sdk/io/iceberg/cdc/ResolveChanges.java    |  168 ++++
 .../io/iceberg/cdc/SerializableChangelogTask.java  |    9 +-
 .../beam/sdk/io/iceberg/cdc/SnapshotWindowFn.java  |   87 ++
 .../sdk/io/iceberg/cdc/WatchForSnapshotsSdf.java   |   57 +-
 .../org/apache/beam/sdk/io/iceberg/AddFilesIT.java |   18 +-
 .../IcebergCdcReadSchemaTransformProviderTest.java |  114 +++
 .../beam/sdk/io/iceberg/IcebergIOReadTest.java     |   48 +
 .../beam/sdk/io/iceberg/IcebergScanConfigTest.java |  270 +++++
 .../IcebergSchemaTransformTranslationTest.java     |    1 +
 .../beam/sdk/io/iceberg/IcebergUtilsTest.java      |   53 +-
 .../IcebergWriteSchemaTransformProviderTest.java   |   76 +-
 .../beam/sdk/io/iceberg/PartitionUtilsTest.java    |   10 +-
 .../sdk/io/iceberg/RecordWriterManagerTest.java    |   24 +-
 .../sdk/io/iceberg/SerializableDataFileTest.java   |  156 ++-
 .../catalog/BigQueryMetastoreCatalogIT.java        |    1 +
 .../io/iceberg/catalog/IcebergCatalogBaseIT.java   |  292 +++++-
 .../sdk/io/iceberg/catalog/RESTCatalogBLMSIT.java  |    1 -
 .../io/iceberg/cdc/ApplyWatermarkColumnTest.java   |  158 +++
 .../beam/sdk/io/iceberg/cdc/CdcReadUtilsTest.java  |    2 +-
 .../beam/sdk/io/iceberg/cdc/CdcResolverTest.java   |  156 +++
 .../sdk/io/iceberg/cdc/ChangelogScannerTest.java   |   22 +-
 .../cdc/IncrementalChangelogSourceTest.java        |  514 ++++++++++
 .../sdk/io/iceberg/cdc/LocalResolveDoFnTest.java   |  340 +++++++
 .../beam/sdk/io/iceberg/cdc/OverlapRangeTest.java  |  161 +++
 .../sdk/io/iceberg/cdc/ReadFromChangelogsTest.java |  366 +++++++
 .../sdk/io/iceberg/cdc/ResolveChangesTest.java     |  222 ++++
 .../iceberg/cdc/SerializableChangelogTaskTest.java |    2 +-
 .../sdk/io/iceberg/cdc/SnapshotWindowFnTest.java   |   93 ++
 .../io/iceberg/cdc/WatchForSnapshotsSdfTest.java   |   49 +-
 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     |  115 ---
 .../sdk/io/jms/JmsReadSchemaTransformProvider.java |    1 -
 .../io/jms/JmsWriteSchemaTransformProvider.java    |    1 -
 .../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 |   30 +-
 sdks/java/io/kafka/build.gradle                    |    1 +
 .../java/org/apache/beam/sdk/io/kafka/KafkaIO.java |    7 +-
 ...KafkaIOReadImplementationCompatibilityTest.java |   22 +
 .../beam/sdk/io/mongodb/MongoDbGridFSIO.java       |    8 +-
 .../org/apache/beam/sdk/io/mongodb/MongoDbIO.java  |   14 +-
 .../org/apache/beam/io/requestresponse/Call.java   |   11 +-
 .../apache/beam/io/requestresponse/CallTest.java   |   51 +-
 .../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            |    3 +-
 sdks/python/apache_beam/dataframe/io_test.py       |   21 +-
 .../io/external/xlang_debeziumio_it_test.py        |    2 +-
 .../apache_beam/io/external/xlang_jmsio_it_test.py |  296 ++++--
 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    |   49 +
 sdks/python/apache_beam/io/gcp/bigquery_tools.py   |    8 +
 .../apache_beam/io/gcp/bigquery_write_it_test.py   |   71 +-
 .../apache_beam/io/gcp/gcsfilesystem_test.py       |    4 +-
 .../apache_beam/io/gcp/healthcare/dicomclient.py   |    5 +-
 sdks/python/apache_beam/io/gcp/pubsub_test.py      |   42 +
 sdks/python/apache_beam/io/textio_test.py          |    7 +-
 sdks/python/apache_beam/io/watch.py                |  264 ++++-
 sdks/python/apache_beam/io/watch_test.py           |  387 ++++++-
 .../python/apache_beam/options/pipeline_options.py |    5 +
 .../apache_beam/options/pipeline_options_test.py   |    9 +
 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 +-
 .../runners/direct/transform_evaluator.py          |    5 +
 .../runners/interactive/recording_manager_test.py  |    5 +-
 .../runners/portability/spark_runner_test.py       |   20 -
 .../testing/benchmarks/chicago_taxi/run_chicago.sh |    6 +-
 sdks/python/apache_beam/transforms/combiners.py    |   58 ++
 .../apache_beam/transforms/combiners_test.py       |   54 +
 .../transforms/managed_iceberg_it_test.py          |    7 +-
 sdks/python/apache_beam/typehints/schemas.py       |  152 ++-
 sdks/python/apache_beam/typehints/schemas_test.py  |  120 +++
 sdks/python/apache_beam/utils/timestamp.py         |  332 ++++--
 sdks/python/apache_beam/utils/timestamp_test.py    |  204 +++-
 sdks/python/apache_beam/version.py                 |    2 +-
 .../yaml/extended_tests/databases/iceberg.yaml     |    3 +
 .../{blueprints => e2e}/delta_lake_to_iceberg.yaml |    0
 sdks/python/apache_beam/yaml/standard_io.yaml      |    1 +
 sdks/python/apache_beam/yaml/yaml_io.py            |  143 ++-
 sdks/python/apache_beam/yaml/yaml_io_test.py       |  219 ++++
 sdks/python/build.gradle                           |    8 +-
 sdks/python/container/Dockerfile                   |   14 +-
 .../container/base_image_requirements_manual.txt   |    3 +-
 sdks/python/container/boot.go                      |  131 +--
 .../license_scripts/upgrade_bundled_pip.py         |   15 +-
 .../container/ml/py310/base_image_requirements.txt |  154 ++-
 .../container/ml/py310/gpu_image_requirements.txt  |  210 ++--
 .../container/ml/py311/base_image_requirements.txt |  157 ++-
 .../container/ml/py311/gpu_image_requirements.txt  |  213 ++--
 .../container/ml/py312/base_image_requirements.txt |  157 ++-
 .../container/ml/py312/gpu_image_requirements.txt  |  211 ++--
 .../container/ml/py313/base_image_requirements.txt |  157 ++-
 sdks/python/container/piputil.go                   |   39 +-
 sdks/python/container/profiler.go                  |  372 ++++++-
 sdks/python/container/profiler_test.go             |  134 +++
 .../container/py310/base_image_requirements.txt    |  142 ++-
 .../container/py311/base_image_requirements.txt    |  145 ++-
 .../container/py312/base_image_requirements.txt    |  145 ++-
 .../container/py313/base_image_requirements.txt    |  145 ++-
 .../container/py314/base_image_requirements.txt    |  147 ++-
 sdks/python/pyproject.toml                         |    3 +-
 sdks/python/scripts/run_snapshot_publish.sh        |    4 +-
 sdks/python/setup.py                               |    2 +-
 sdks/python/test-suites/dataflow/common.gradle     |    4 +-
 sdks/python/test-suites/direct/build.gradle        |    3 +-
 sdks/python/test-suites/tox/py310/build.gradle     |   27 +-
 sdks/standard_expansion_services.yaml              |    1 +
 sdks/typescript/container/boot.go                  |    4 +-
 sdks/typescript/package.json                       |    2 +-
 settings.gradle.kts                                |    2 +-
 website/Dockerfile                                 |    2 +-
 website/build.gradle                               |    5 +-
 .../www/site/assets/scss/_capability-matrix.scss   |  123 +--
 .../www/site/assets/scss/capability-matrix.scss    |  121 +--
 .../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 |   10 +-
 .../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 +-
 .../site/content/en/get-started/quickstart-java.md |    4 +-
 .../site/content/en/get-started/quickstart-py.md   |   17 +
 .../documentation/capability-matrix-big.html       |   36 +-
 .../documentation/capability-matrix-single.html    |   37 +-
 website/www/yarn.lock                              |    6 +-
 371 files changed, 17274 insertions(+), 4451 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 
examples/java/src/test/java/org/apache/beam/examples/subprocess/utils/FileUtilsTest.java
 copy 
sdks/java/core/src/main/java/org/apache/beam/sdk/util/ThrowingRunnable.java => 
runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/FailedWorkHandler.java
 (83%)
 create mode 100644 
runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/DataflowOutputCounterTest.java
 delete mode 100644 sdks/go/pkg/beam/artifact/options.go
 delete mode 100644 sdks/go/pkg/beam/artifact/options_test.go
 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
 copy 
sdks/java/io/delta/src/main/java/org/apache/beam/sdk/io/delta/{DeltaReadSchemaTransformProvider.java
 => DeltaCdcReadSchemaTransformProvider.java} (56%)
 copy 
sdks/java/io/delta/src/test/java/org/apache/beam/sdk/io/delta/{DeltaIOIT.java 
=> DeltaIOS3IT.java} (64%)
 create mode 100644 
sdks/java/io/delta/src/test/java/org/apache/beam/sdk/io/delta/DeltaWriteTestUtils.java
 rename 
sdks/java/io/google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/bigquery/{StorageApiSinkSchemaUpdateIT.java
 => StorageApiSinkSchemaUpdateITBase.java} (94%)
 create mode 100644 
sdks/java/io/google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/bigquery/StorageApiSinkSchemaUpdateWithInputSchemaIT.java
 create mode 100644 
sdks/java/io/google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/bigquery/StorageApiSinkSchemaUpdateWithoutInputSchemaIT.java
 delete mode 100644 
sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/IncrementalScanSource.java
 delete mode 100644 
sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/ReadFromTasks.java
 delete mode 100644 
sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/WatchForSnapshots.java
 create mode 100644 
sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/cdc/ApplyWatermarkColumn.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/IncrementalChangelogSource.java
 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/main/java/org/apache/beam/sdk/io/iceberg/cdc/ResolveChanges.java
 create mode 100644 
sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/cdc/SnapshotWindowFn.java
 create mode 100644 
sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/IcebergScanConfigTest.java
 create mode 100644 
sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/cdc/ApplyWatermarkColumnTest.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/IncrementalChangelogSourceTest.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
 create mode 100644 
sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/cdc/ResolveChangesTest.java
 create mode 100644 
sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/cdc/SnapshotWindowFnTest.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
 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
 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
 rename sdks/python/apache_beam/yaml/extended_tests/{blueprints => 
e2e}/delta_lake_to_iceberg.yaml (100%)
 create mode 100644 sdks/python/container/profiler_test.go

Reply via email to