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