This is an automated email from the ASF dual-hosted git repository.
dependabot[bot] pushed a change to branch
dependabot/go_modules/learning/tour-of-beam/backend/google.golang.org/grpc-1.82.1
in repository https://gitbox.apache.org/repos/asf/beam.git
omit b5ccbc73091 Bump google.golang.org/grpc in
/learning/tour-of-beam/backend
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 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 1a43d8de57d [Python] Refactor MatchContinuously onto the Watch
transform (#39461)
add ce13c65a4fb Set Apache Beam user agent for Python BigtableIO write
client (#39792)
add 68024f21b2b Add Delta Lake to Iceberg Yaml blueprint IT (#39549)
add fde5698dc99 test for column default values (#39739)
add a01a5cf6ef4 [GSoC 2026] Fix Duplicated Subscription Path in
stale_cleaner.py (#39794)
add af7f50e62e0 Fix stale broken cluster connection in CassandraIO (#39788)
add 4c0d2ab465d Bump golang.org/x/net from 0.57.0 to 0.58.0 in /sdks
(#39802)
add ac087db42b4 Bump github.com/aws/aws-sdk-go-v2 from 1.43.5 to 1.43.6 in
/sdks (#39797)
add 7bdd5d1a4ad Bump github.com/nats-io/nats.go from 1.52.0 to 1.53.1 in
/sdks (#39798)
add 719f811a03b Bump github.com/aws/aws-sdk-go-v2/service/s3 in /sdks
(#39799)
add ca6065508a3 Bump github.com/aws/aws-sdk-go-v2/config in /sdks (#39801)
add c08f2a89452 [#72] Fix Golang Zip Slip Vulnerability (#39638)
add 52702a15957 [python] Add Secret management module in
apache_beam.utils.secret (#39636)
add e1a72b3fc7b Updates Dataflow Python container (#39796)
add 36066509b81 Fix PeriodicImpulse/PeriodicSequence watermark regression
(#39026) (#39465)
add f9b15b711ae Refactor Java Secret classes to align with Python SDK
(#39806)
add e7aad65cf0e Bump github.com/moby/go-archive from 0.2.0 to 0.3.0 in
/sdks (#39810)
add 3889477e216 [GSoC-273] Fixing the github action Unmanaged Service
Account Keys (#39177)
add 6356a3c3d28 [Dataflow Streaming] Commit size validation for multi key
commits (#39473)
add b68a384657b Bump google.golang.org/protobuf from 1.36.11 to 1.36.12 in
/sdks (#39814)
add a99c2638cef Bump github.com/nats-io/nats-server/v2 from 2.14.4 to
2.14.5 in /sdks (#39813)
add 658397cde39 Bump github.com/aws/aws-sdk-go-v2/feature/s3/manager in
/sdks (#39812)
add d8935269847 Set getSession to synchronized to avoid race conditions
(#39817)
add 6a8eee94aa4 Bump ClassGraph from 4.8.162 to 4.8.192 (#39815)
add e4962bc90d9 [GSoC-273] Implement TestPubsubContext for Python GCP
Integration Tests (#39685)
add 94510075efe Update SKILLs based on review practice (#39805)
add 46e46f07f32 Make primary channel failover timeout configurable (#39646)
add 93f3e051c3e Require google-cloud-bigtable>=2.42.0 and test write error
surfacing (#39820)
add a26ecfd4857 [#39723] Implement model for Iceberg side input cache
(#39724)
add 67018be2650 [runners-spark] Support stateful ParDo in the Structured
Streaming batch runner (#39793)
add 2d94b11acaf normalize KinesisIO
add 02060be7b7a sync Kinesis
add 661862a4ba2 Add LocalStack Kinesis YAML integration test
add 2ff4f0eda94 Merge pull request #39763 from
aIbrahiim/yaml-normalize-kinesis
add 7420ad22aea Update google-cloud-bigtable
add 11e9f8b5598 Merge pull request #39837 from apache/update-bigtable
add e626f54690b Implement Vertex AI Model Monitoring v2 (#39738)
add 3b9e647dec1 [Prism] Schedule consumers of a self checkpointing source
(#39572)
add 88b3ee7b488 JmsIO yaml (#39818)
add 4e5ae91554e Support lakehouse PCNT format in BigQueryIO storage read.
(#39597)
add 6434f7411e7 [Docs] Document UnboundedSource in the Python I/O
connector guide (#39529)
add 61ed38fd5fd Support Secret Manager in JdbcIO for Java, Python and YAML
(#39834)
add 7de4d3ab73b Fix KafkaIO wrtie SchemaTransform parallelism (#39844)
add 015e823a1a9 [GSoC-273] Feat: Integrate TestPubsubContext to prevent
Pub/Sub resource leaks and expand stale cleaner scope (#39826)
add c1017953ac3 Adds support for reading at a given Delta Lake version or
timestamp (#39758)
add 2577a6cff4a Install Cloud Spanner emulator component in Go PreCommit
CI (#39850)
add 2e445307153 Add Kafka Streams runner skeleton module and portable
entry points
add 61c891a69ca Address review notes on KafkaStreamsPipelineResult and
state dir
add cef6544e792 Address review feedback on Kafka Streams Runner skeleton
add 0d445f79e0c Drop setRunner(null) suppression; make applicationId
required
add d4c7dae98b8 Catch Exception in KafkaStreamsRunner.run() to avoid
job-server leak
add 64ad482f9de Merge pull request #38534: [GSoC 2026] Kafka Streams
runner skeleton module + portable entry points
add faefc95e9f0 [GSoC 2026] Kafka Streams runner — translation framework +
Impulse translator (#38689)
add 0f875a07162 Add temporary feature-branch CI for Kafka Streams runner
(#38725)
add ff554f7438d [GSoC 2026] Kafka Streams runner — ExecutableStage
(stateless ParDo) translator (#38764)
add 4c02a7789a2 [GSoC 2026] Kafka Streams runner — Redistribute translator
+ ExecutableStage type-agnostic edge (#38843)
add 4636b1f333e #38957: Add in-memory WatermarkManager core
(per-source-partition tracking
add 1058e94027e [GSoC 2026] Kafka Streams runner #38987: Wire
WatermarkManager into ExecutableStageProcessor
add 3e963a68343 [GSoC 2026] Kafka Streams runner #39051: Add
KStreamsPayload Serde for crossing topic boundaries
add 5e65d4772fb [GSoC 2026] Kafka Streams runner #39141: Add GroupByKey
(GlobalWindow, fire at watermark)
add 4d5847f6bfa [GSoC 2026] Kafka Streams runner #39211: Add
KafkaStreamsTestRunner test harness
add 27c4522ae5e [GSoC 2026] Kafka Streams runner #39249: Support Create
add ff323ebd6db [GSoC 2026] Kafka Streams runner #39273: Kafka Streams
runner: Flatten support
add 75adf49ee8d [GSoC 2026] Kafka Streams runner: surface SDK-harness
metrics as MetricResults (#39341)
add 87911f00973 [GSoC 2026] Kafka Streams runner:
TestPipeline-dispatchable test runner (PAssert works) (#39362)
add d9cea589e5b [GSoC 2026] Kafka Streams runner: validatesRunner task;
Create and Flatten suites green (#39380)
add fca14356614 [GSoC 2026] Kafka Streams runner: multi-output executable
stages (#39410)
add b615ae85aed [GSoC 2026] Kafka Streams runner: enable ParDoTest in the
ValidatesRunner suite (#39451)
add 47614a943a7 [GSoC 2026] Kafka Streams runner: windowed GroupByKey via
ReduceFnRunner (#39494)
add 27198768f89 [GSoC 2026] Kafka Streams runner: run on a real broker,
correctly across partitions (#39546)
add cf1f10bdf82 [GSoC 2026] Kafka Streams runner: bound a bundle by
element count (#39578)
add fc36301e6c2 [GSoC 2026] Kafka Streams runner: CombineTest coverage and
two review follow-ups (#39610)
add 5861f31e8ac [GSoC 2026] Kafka Streams runner: read unbounded sources
(#39611)
add cb30afd092e [GSoC 2026] Kafka Streams runner: user documentation,
marked experimental (#39627)
add e051e06bf83 [GSoC 2026] Kafka Streams runner: Python wrapper that
starts its own job server (#39680)
add 4ff618047c3 [GSoC 2026] Kafka Streams runner: terminate a bounded
pipeline when it is drained (#39700)
add 5f1658f8b5e [GSoC 2026] Kafka Streams runner: portable ValidatesRunner
suite for Python (#39736)
add 511a40e4f20 [GSoC 2026] Kafka Streams runner: separate the source's
poll size from the bundle size, and expose the session timeout (#39748)
add 65a2e400c73 [GSoC 2026] Kafka Streams runner: bound a source poll in
time, not only in elements (#39761)
add 104dc272e8d [GSoC 2026] Kafka Streams runner: ask for primitive reads
in the Java wrapper (#39766)
add 49d459a243e [GSoC 2026] Kafka Streams runner: put the runner behind an
opt-in build flag (#39762)
add 10ff557d616 [GSoC 2026] Kafka Streams runner: an application for
measuring instances coming and going (#39752)
add 5b3702484f4 [GSoC 2026] Kafka Streams runner: license header and
Python formatting for master CI
add b12bfc159ba [GSoC 2026] Kafka Streams runner: shorten the explanation
comments
add 8d5f6511a6b Merge pull request #39781: [GSoC 2026] Kafka Streams
runner: shorten comments, and fix the license header and Python formatting
add e2681085c96 Build Kafka Streams runner during javaPreCommit (#18479)
add 8b3b08dc799 Merge pull request #39784: Build Kafka Streams runner
during javaPreCommit (#18479)
add b67cf234ead [GSoC 2026] Kafka Streams runner: update the CHANGES.md
entry
add 52dadbac433 Merge pull request #39786: [GSoC 2026] Kafka Streams
runner: update the CHANGES.md entry
add 89617f59b45 Merge branch 'master' of https://github.com/apache/beam
into feat/18479-kafka-streams-runner-skeleton
add 6b75a44f0da Removed feature branch build
add 52d6eeab74e [GSoC 2026] Kafka Streams runner: say what the Python side
is, and is not
add 473a5dd649c Merge pull request #39847: [GSoC 2026] Kafka Streams
runner: say what the Python side is, and is not
add e10e8c969b0 Merge pull request #39785: Merge Kafka Streams Runner
skeleton (#18479)
add 7056e67a285 [GSoC 2026] Publish the Kafka Streams runner to nightly
snapshots
add ec5ebf422f0 Merge pull request #39854: [GSoC 2026] Publish the Kafka
Streams runner to nightly snapshots
add 466cf6d265b Revert "[GSoC-273] Feat: Integrate TestPubsubContext to
prevent Pub/Sub resou…" (#39853)
add 59e916e5a11 AddFiles: regenerate name mapping, read footers once,
error helper (#39836)
add 1ef9cfd66fb Bump actions/checkout from 6 to 7 (#39859)
add 08522c7ca94 Bump github/codeql-action from 4.37.7 to 4.37.8 (#39864)
add f15bd00f526 use LocalStack TLS hostname in Kinesis YAML IT (#39868)
add f95421bbd90 Revert "Install Cloud Spanner emulator component in Go
PreCommit CI (#39850)" (#39877)
add da360800706 Bump github.com/aws/aws-sdk-go-v2/feature/s3/manager in
/sdks (#39860)
add 3e30bac4cf6 Bump google.golang.org/grpc from 1.83.0 to 1.83.1 in /sdks
(#39863)
add ffbe1d58390 Bump cloud.google.com/go/pubsub from 1.51.0 to 1.51.1 in
/sdks (#39862)
add ea7d39cb9ab Bump cloud.google.com/go/storage from 1.64.0 to 1.65.0 in
/sdks (#39861)
add af99c7ce5d4 [Bigtable] Fix Bigtable segment truncation when open end
key is startKey + null byte (#39842) (#39843)
add 66ecc31a22c Adds a new CoderTranslator for Java SchemaCoders. (#39594)
add 4bd8a86fd83 Bump github.com/aws/smithy-go from 1.27.8 to 1.27.9 in
/sdks (#39880)
add 13875fc6bd3 Bump cloud.google.com/go/bigquery from 1.80.0 to 1.81.0 in
/sdks (#39881)
add 73ea0b72c12 [Prism] Honor the resume delay of self-checkpointing SDF
residuals (#39849)
add c267ce6520f Bump github.com/fsouza/fake-gcs-server from 1.55.1 to
1.56.0 in /sdks (#39890)
add 571184a509e Add Vertex Model Monitoring to CHANGES.md (#39882)
add 87e31bd8a31 Update build.gradle.kts (#39892)
add 167b8649fc9 [Java] Bound the Watch deduplication state with a
timestamp cursor (#39746)
add b2013e09f91 Fix ZeroDivisionError at initial invocation of
monitoring_info when there is no work yet (#39885)
add c494cc647fb Feat: expanding the context of the second and third
security layers of testPubSubContext so that, in the event of failure,
subscriptions retain a 24-hour grace period before being deleted (#39856)
add 1fca1c2d371 prevent cloudML TFT extras from replacing the SDK under
test (#39884)
add 69f7fd48a0c Add Snowflake YAML write transform
add 85bad221c71 Add Snowflake Read and Streaming Write
add 90c93d7f541 Fix bytes
add d131e0e74b7 Fix imports
add f4e62fdafbf Refactoring
add 15efa5d2084 Add extended test
add b275c84f8e2 Add license
add 27c4e6395cb Move depends
add a93fdcdda75 Fix jms yaml
add da201e642a2 Fix import Nullable
add 497870ade99 Change to provided
add 8cfbd8e01c2 Merge pull request #39742 from apache/snowflakeio-yaml
add ad39be0f58b Clarify service account key issue reports (#39894)
add 4d2c83d9b66 Bump github.com/fsouza/fake-gcs-server from 1.56.0 to
1.56.1 in /sdks (#39902)
add 351f4c7eac9 Bump actions/setup-java from 5 to 6 (#39903)
add cf27972eb10 Update CHANGES.md with new known issues (#39898)
add 01559cc9127 Added schema_update_options to Python BigQuery writes
(#39078)
add 34a00c7781f AddFiles: extract bounded async task plumbing and Parquet
footer reads (#39896)
add 8d5a5300fa3 Remove nullness suppression in trigger implementation
add 73e5ecc2ea8 Merge pull request #39807: Remove nullness suppression in
trigger implementation
add 493a18f8217 Bump docker/setup-buildx-action from 4.2.0 to 4.3.0
(#39829)
add 67ed9595191 Bump google.golang.org/grpc in
/learning/tour-of-beam/backend
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 (b5ccbc73091)
\
N -- N -- N
refs/heads/dependabot/go_modules/learning/tour-of-beam/backend/google.golang.org/grpc-1.82.1
(67ed9595191)
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:
.agent/skills/beam-concepts/SKILL.md | 31 +
.agent/skills/contributing/SKILL.md | 3 +
.asf.yaml | 1 +
.../actions/setup-environment-action/action.yml | 24 +-
.../IO_Iceberg_Integration_Tests.json | 2 +-
.../IO_Iceberg_Integration_Tests_Dataflow.json | 2 +-
.../beam_CloudML_Benchmarks_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 +-
...am_PostCommit_Java_ValidatesRunner_Spark4.json} | 0
...a_ValidatesRunner_SparkStructuredStreaming.json | 3 +-
.github/trigger_files/beam_PostCommit_Python.json | 2 +-
..._Spark.json => beam_PostCommit_Python_Arm.json} | 0
.../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 +-
... beam_PreCommit_Java_Kafka_Streams_Runner.json} | 0
.github/trigger_files/beam_PreCommit_SQL.json | 2 +-
.github/workflows/README.md | 2 +-
.../beam_Infrastructure_AuditUnmanagedKeys.yml | 71 --
.../beam_Infrastructure_PolicyEnforcer.yml | 13 +-
.../beam_PerformanceTests_xlang_KafkaIO_Python.yml | 2 +-
.github/workflows/beam_PostCommit_Go.yml | 2 +-
.../workflows/beam_PostCommit_Go_Dataflow_ARM.yml | 2 +-
.../beam_PostCommit_Java_Delta_IO_Dataflow.yml | 9 +
.../beam_PostCommit_Java_Examples_Dataflow_ARM.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 +-
.../beam_PostCommit_XVR_GoUsingJava_Dataflow.yml | 2 +-
.../beam_PostCommit_Yaml_Xlang_Direct.yml | 4 +-
.../workflows/beam_PostRelease_NightlySnapshot.yml | 9 +-
.../workflows/beam_PreCommit_CommunityMetrics.yml | 2 +-
.github/workflows/beam_PreCommit_GHA.yml | 18 +-
.github/workflows/beam_PreCommit_Java.yml | 1 +
...> beam_PreCommit_Java_Kafka_Streams_Runner.yml} | 75 +-
.github/workflows/beam_PreCommit_PythonDocker.yml | 2 +-
.../workflows/beam_Publish_Beam_SDK_Snapshots.yml | 19 +-
.../workflows/beam_Publish_Python_VLLM_Image.yml | 2 +-
...beam_Python_ValidatesContainer_Dataflow_ARM.yml | 2 +-
.github/workflows/beam_Release_NightlySnapshot.yml | 1 +
.github/workflows/build_release_candidate.yml | 26 +-
.github/workflows/build_runner_image.yml | 2 +-
.github/workflows/build_wheels.yml | 12 +-
.github/workflows/code_completion_plugin_tests.yml | 2 +-
.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 +-
.github/workflows/python_dependency_tests.yml | 1 +
.../republish_released_docker_containers.yml | 4 +-
.../workflows/run_rc_validation_go_wordcount.yml | 12 +-
.../run_rc_validation_python_mobile_gaming.yml | 6 +-
.../workflows/run_rc_validation_python_yaml.yml | 6 +-
.github/workflows/stale.yml | 2 +-
.github/workflows/tour_of_beam_backend.yml | 4 +-
.../workflows/tour_of_beam_backend_integration.yml | 1 +
.github/workflows/typescript_tests.yml | 4 +-
.github/workflows/update_python_dependencies.yml | 2 +
.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/mock-apis/go.mod | 2 +-
.test-infra/mock-apis/go.sum | 4 +-
.test-infra/tools/stale_cleaner.py | 27 +-
.test-infra/tools/test_stale_cleaner.py | 68 +-
CHANGES.md | 91 +-
build.gradle.kts | 19 +-
.../org/apache/beam/gradle/BeamModulePlugin.groovy | 19 +-
contributor-docs/README.md | 1 +
contributor-docs/local-flink-python.md | 204 ++++
examples/java/iceberg/build.gradle | 4 +-
.../beam/examples/complete/game/UserScore.java | 2 +-
.../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/enforcement/README.md | 20 +-
infra/enforcement/account_keys.py | 77 +-
infra/enforcement/iam.py | 83 +-
infra/enforcement/sending.py | 67 +-
infra/enforcement/test_sending.py | 6 +-
infra/iam/users.yml | 13 +-
it/iceberg/build.gradle | 4 +-
it/mongodb/build.gradle | 1 +
.../beam/it/mongodb/MongoDBResourceManager.java | 11 +-
.../it/mongodb/MongoDBResourceManagerTest.java | 5 +
learning/tour-of-beam/backend/function.go | 6 +-
.../backend/integration_tests/client.go | 8 +-
.../backend/internal/fs_content/yaml.go | 10 +-
.../backend/internal/storage/datastore.go | 11 +-
.../tour-of-beam/backend/internal/storage/mock.go | 8 +-
.../beam/model/fnexecution/v1/standard_coders.yaml | 31 +
.../model/pipeline/v1/external_transforms.proto | 2 +
.../org/apache/beam/model/pipeline/v1/schema.proto | 13 +
playground/backend/go.mod | 7 +-
playground/backend/go.sum | 10 +-
playground/kafka-emulator/build.gradle | 11 +
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 +-
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 -
.../core/triggers/AfterAllStateMachine.java | 5 +-
.../AfterDelayFromFirstElementStateMachine.java | 5 +-
.../core/triggers/AfterEachStateMachine.java | 26 +-
.../core/triggers/AfterFirstStateMachine.java | 5 +-
.../core/triggers/AfterWatermarkStateMachine.java | 12 +-
.../core/triggers/DefaultTriggerStateMachine.java | 3 -
.../triggers/ExecutableTriggerStateMachine.java | 21 +-
.../core/triggers/OrFinallyStateMachine.java | 5 +-
.../core/triggers/RepeatedlyStateMachine.java | 5 +-
.../runners/core/triggers/TriggerStateMachine.java | 17 +-
.../TriggerStateMachineContextFactory.java | 23 +-
.../beam/runners/core/triggers/package-info.java | 4 -
.../runners/core/SplittableParDoProcessFnTest.java | 283 ++++-
.../core/triggers/AfterEachStateMachineTest.java | 54 +
.../runners/extensions/metrics/package-info.java | 4 -
runners/google-cloud-dataflow-java/build.gradle | 6 +-
.../dataflow/DataflowPipelineTranslator.java | 2 +-
.../beam/runners/dataflow/DataflowRunner.java | 54 +-
.../options/DataflowStreamingPipelineOptions.java | 2 +-
.../dataflow/DataflowPipelineTranslatorTest.java | 78 +-
.../google-cloud-dataflow-java/worker/build.gradle | 2 +
.../dataflow/worker/DataflowExecutionContext.java | 4 +
.../dataflow/worker/DataflowOutputCounter.java | 68 +-
.../dataflow/worker/DataflowWorkUnitClient.java | 12 +-
.../worker/IntrinsicMapTaskExecutorFactory.java | 11 +-
.../dataflow/worker/MultiKeyBundleOptions.java | 138 +++
.../dataflow/worker/SimpleParDoFnHelpers.java | 5 +-
.../dataflow/worker/StreamingDataflowWorker.java | 85 +-
.../StreamingGroupAlsoByWindowViaWindowSetFn.java | 2 +-
.../worker/StreamingModeExecutionContext.java | 129 ++-
.../dataflow/worker/UngroupedWindmillReader.java | 7 +-
.../dataflow/worker/WindmillKeyedWorkItem.java | 33 +-
.../WindmillOpenTelemetryContextPropagator.java | 22 +-
.../beam/runners/dataflow/worker/WindmillSink.java | 43 +-
...Exception.java => WorkCancellingException.java} | 33 +-
.../worker/WorkItemCancelledException.java | 29 +-
.../logging/DataflowWorkerLoggingHandler.java | 35 +-
.../logging/DataflowWorkerLoggingInitializer.java | 4 +
.../worker/logging/DataflowWorkerLoggingMDC.java | 15 +-
.../dataflow/worker/streaming/ActiveWorkState.java | 51 +-
.../streaming/BoundedQueueExecutorWorkHandle.java | 9 +-
.../worker/streaming/ComputationState.java | 16 +-
.../worker/streaming/ComputationWorkExecutor.java | 11 +-
.../dataflow/worker/streaming/ExecutableWork.java | 8 +
...cutorWorkHandle.java => FailedWorkHandler.java} | 10 +-
...java => MultiKeyCommitValidationException.java} | 10 +-
.../runners/dataflow/worker/streaming/Work.java | 80 +-
.../streaming/harness/MetricsDataProvider.java | 4 +-
.../harness/StreamingWorkerStatusReporter.java | 2 +-
.../dataflow/worker/util/BoundedQueueExecutor.java | 67 +-
.../dataflow/worker/util/KeyGroupWorkQueue.java | 12 +-
.../worker/windmill/client/commits/Commit.java | 8 +-
.../windmill/client/commits/CompleteCommit.java | 15 +-
.../commits/StreamingApplianceWorkCommitter.java | 3 +-
.../commits/StreamingEngineWorkCommitter.java | 29 +-
.../client/getdata/StreamGetDataClient.java | 5 +-
.../client/grpc/stubs/FailoverChannel.java | 66 +-
.../worker/windmill/state/WindmillStateReader.java | 9 +-
.../processing/ComputationWorkExecutorFactory.java | 9 +-
.../work/processing/StreamingWorkScheduler.java | 206 ++--
.../processing/failures/WorkFailureProcessor.java | 107 +-
.../windmill/work/refresh/ActiveWorkRefresher.java | 21 +-
.../dataflow/worker/DataflowOutputCounterTest.java | 108 ++
.../worker/DataflowWorkUnitClientTest.java | 8 +-
.../IntrinsicMapTaskExecutorFactoryTest.java | 14 +-
.../worker/KeyTokenInvalidExceptionTest.java | 39 -
.../worker/StreamingDataflowWorkerTest.java | 1079 ++++++++++++++++--
.../worker/StreamingModeExecutionContextTest.java | 405 ++++++-
.../worker/WindmillReaderIteratorBaseTest.java | 4 +-
.../worker/WindowingWindmillReaderTest.java | 3 +-
.../dataflow/worker/WorkerCustomSourcesTest.java | 19 +-
.../logging/DataflowWorkerLoggingHandlerTest.java | 107 +-
.../worker/streaming/ActiveWorkStateTest.java | 69 +-
.../streaming/ComputationStateCacheTest.java | 4 +-
.../worker/streaming/ComputationStateTest.java | 114 ++
.../dataflow/worker/streaming/WorkTest.java | 4 +-
.../worker/testing/RestoreDataflowLoggingMDC.java | 8 +-
.../testing/RestoreDataflowLoggingMDCTest.java | 10 +-
.../worker/util/BoundedQueueExecutorTest.java | 205 +++-
.../worker/util/KeyGroupWorkQueueTest.java | 31 +-
.../StreamingApplianceWorkCommitterTest.java | 4 +-
.../commits/StreamingEngineWorkCommitterTest.java | 64 +-
.../client/grpc/GrpcCommitWorkStreamTest.java | 188 ++++
.../client/grpc/stubs/FailoverChannelTest.java | 107 +-
.../windmill/state/WindmillStateReaderTest.java | 10 +-
.../failures/WorkFailureProcessorTest.java | 153 ++-
.../work/refresh/ActiveWorkRefresherTest.java | 73 +-
.../worker/windmill/src/main/proto/windmill.proto | 4 +
.../control/ProcessBundleDescriptorsTest.java | 6 +-
runners/kafka-streams/build.gradle | 201 ++++
runners/kafka-streams/job-server/build.gradle | 88 ++
runners/kafka-streams/measurement/build.gradle | 78 ++
.../kafka-streams/measurement/docker-compose.yml | 44 +
.../streams/measurement/RescalingMeasurement.java | 246 +++++
.../kafka/streams/measurement}/package-info.java | 13 +-
.../proto/build.gradle} | 32 +-
.../src/main/proto/kafka_streams_payload.proto | 58 +
.../kafka/streams/KafkaStreamsJobInvoker.java | 96 ++
.../kafka/streams/KafkaStreamsJobServerDriver.java | 106 ++
.../kafka/streams/KafkaStreamsPipelineOptions.java | 151 +++
.../kafka/streams/KafkaStreamsPipelineResult.java | 77 ++
.../kafka/streams/KafkaStreamsPipelineRunner.java | 185 ++++
.../KafkaStreamsPortablePipelineResult.java | 154 +++
.../runners/kafka/streams/KafkaStreamsRunner.java | 138 +++
.../kafka/streams/KafkaStreamsRunnerRegistrar.java | 48 +
.../kafka/streams/KafkaStreamsTopicManager.java | 171 +++
.../beam/runners/kafka/streams}/package-info.java | 8 +-
.../streams/translation/EmptyBoundedSource.java | 89 ++
.../translation/ExecutableStageProcessor.java | 365 +++++++
.../translation/ExecutableStageTranslator.java | 133 +++
.../streams/translation/FlattenProcessor.java | 124 +++
.../streams/translation/FlattenTranslator.java | 95 ++
.../GroupByKeyBroadcastPartitioner.java | 70 ++
.../streams/translation/GroupByKeyTranslator.java | 215 ++++
.../streams/translation/ImpulseProcessor.java | 145 +++
.../streams/translation/ImpulseTranslator.java | 77 ++
.../kafka/streams/translation/KStreamsPayload.java | 184 ++++
.../streams/translation/KStreamsPayloadSerde.java | 123 +++
.../KafkaStreamsExecutableStageContextFactory.java | 66 ++
.../KafkaStreamsPipelineTranslator.java | 208 ++++
.../translation/KafkaStreamsStateInternals.java | 463 ++++++++
.../translation/KafkaStreamsTimerInternals.java | 258 +++++
.../KafkaStreamsTranslationContext.java | 179 +++
.../streams/translation/PTransformTranslator.java | 41 +
.../kafka/streams/translation/ReadProcessor.java | 206 ++++
.../kafka/streams/translation/ReadTranslator.java | 272 +++++
.../translation/RedistributeTranslator.java | 57 +
.../streams/translation/ShuffleByKeyProcessor.java | 138 +++
.../streams/translation/StageOutputProcessor.java | 103 ++
.../kafka/streams/translation/StoreKeys.java | 105 ++
.../streams/translation/TerminationReporter.java | 105 ++
.../streams/translation/TerminationTracker.java | 176 +++
.../translation/UnboundedReadProcessor.java | 338 ++++++
.../streams/translation/WatermarkAggregator.java | 98 ++
.../streams/translation/WatermarkManager.java | 136 +++
.../streams/translation/WatermarkPayload.java | 49 +
.../translation/WindowedGroupByKeyProcessor.java | 337 ++++++
.../kafka/streams/translation}/package-info.java | 8 +-
.../streams/KafkaStreamsJobServerDriverTest.java | 71 ++
.../streams/KafkaStreamsPipelineOptionsTest.java | 90 ++
.../KafkaStreamsPipelineRunnerConfigTest.java | 86 ++
.../KafkaStreamsPortablePipelineResultTest.java | 96 ++
.../kafka/streams/KafkaStreamsRunnerBrokerIT.java | 372 +++++++
.../kafka/streams/KafkaStreamsRunnerTest.java | 159 +++
.../kafka/streams/KafkaStreamsTestRunner.java | 149 +++
.../kafka/streams/MultiOutputStageTest.java | 87 ++
.../kafka/streams/TestKafkaStreamsRunner.java | 151 +++
.../kafka/streams/TestKafkaStreamsRunnerTest.java | 82 ++
.../streams/translation/BundleBoundaryTest.java | 118 ++
.../translation/ChainedExecutableStageTest.java | 163 +++
.../kafka/streams/translation/CreateTest.java | 73 ++
.../ExecutableStageProcessorWatermarkTest.java | 145 +++
.../translation/ExecutableStageTranslatorTest.java | 82 ++
.../translation/FixedWindowGroupByKeyTest.java | 106 ++
.../translation/FlattenParallelismTest.java | 141 +++
.../kafka/streams/translation/FlattenTest.java | 212 ++++
.../kafka/streams/translation/GroupByKeyTest.java | 93 ++
.../streams/translation/ImpulseTranslatorTest.java | 144 +++
.../translation/KStreamsPayloadSerdeTest.java | 97 ++
.../KafkaStreamsPipelineTranslatorTest.java | 146 +++
.../KafkaStreamsTimerInternalsTest.java | 208 ++++
.../translation/MetricsAcrossBundlesTest.java | 81 ++
.../kafka/streams/translation/MetricsTest.java | 79 ++
.../kafka/streams/translation/ReadTest.java | 75 ++
.../streams/translation/SharedTestCollector.java | 92 ++
.../translation/ShuffleByKeyProcessorTest.java | 125 +++
.../translation/StageOutputProcessorTest.java | 108 ++
.../StandardWindowFnTranslationTest.java | 152 +++
.../translation/TerminationTrackerTest.java | 172 +++
.../streams/translation/UnboundedReadTest.java | 263 +++++
.../translation/WatermarkAggregatorTest.java | 160 +++
.../streams/translation/WatermarkManagerTest.java | 155 +++
.../translation/WatermarkPropagationTest.java | 97 ++
runners/prism/java/build.gradle | 6 -
runners/spark/job-server/spark_job_server.gradle | 5 +-
runners/spark/spark_runner.gradle | 15 +-
.../translation/batch/DoFnRunnerFactory.java | 26 +-
.../translation/batch/ParDoTranslatorBatch.java | 50 +-
.../translation/batch/PipelineTranslatorBatch.java | 19 +
.../batch/StatefulDoFnGroupFunction.java | 391 +++++++
.../batch/StatefulParDoTranslatorBatch.java | 282 +++++
.../SparkBatchPortablePipelineTranslator.java | 8 +-
.../translation/SparkExecutableStageFunction.java | 172 ++-
.../SparkStreamingPortablePipelineTranslator.java | 4 +-
.../batch/StatefulParDoExecutionTest.java | 357 ++++++
.../batch/StatefulParDoTranslatorBatchTest.java | 261 +++++
.../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 | 115 +-
sdks/go.sum | 241 +++--
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 +
.../core/runtime/xlangx/expansionx/download.go | 17 +-
.../runtime/xlangx/expansionx/download_test.go | 25 +
sdks/go/pkg/beam/core/typex/class.go | 4 +-
sdks/go/pkg/beam/core/typex/fulltype.go | 23 +
sdks/go/pkg/beam/core/typex/special.go | 18 +-
sdks/go/pkg/beam/core/util/reflectx/call.go | 31 +
sdks/go/pkg/beam/pcollection.go | 16 +
sdks/go/pkg/beam/runners/prism/internal/coders.go | 11 +
.../pkg/beam/runners/prism/internal/coders_test.go | 16 +
.../prism/internal/engine/elementmanager.go | 192 +++-
.../engine/elementmanager_continuation_test.go | 417 +++++++
sdks/go/pkg/beam/transforms/batch/batch.go | 677 ++++++++++++
.../pkg/beam/transforms/batch/batch_prism_test.go | 222 ++++
.../beam/transforms/batch/batch_test.go} | 53 +-
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 +-
.../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 -
.../apache/beam/sdk/options/SdkHarnessOptions.java | 7 +
.../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 +
.../org/apache/beam/sdk/schemas/SchemaUtils.java | 7 +
.../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 -
.../apache/beam/sdk/transforms/Redistribute.java | 23 +-
.../java/org/apache/beam/sdk/transforms/Reify.java | 4 +-
.../java/org/apache/beam/sdk/transforms/Watch.java | 239 +++-
.../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 -
.../beam/sdk/transforms/reflect/package-info.java | 3 -
.../transforms/splittabledofn/package-info.java | 4 -
.../sdk/transforms/windowing/package-info.java | 4 -
.../beam/sdk/util/GcpHsmGeneratedSecret.java | 65 +-
.../java/org/apache/beam/sdk/util/GcpSecret.java | 103 +-
.../java/org/apache/beam/sdk/util/RawSecret.java | 57 +
.../main/java/org/apache/beam/sdk/util/Secret.java | 212 ++--
.../sdk/util/construction/CoderTranslation.java | 62 +-
.../construction/CoderTranslatorRegistrar.java | 16 +
.../sdk/util/construction/CoderTranslators.java | 79 ++
.../sdk/util/construction/ModelCoderRegistrar.java | 28 +-
.../beam/sdk/util/construction/ModelCoders.java | 2 +
.../util/construction/RehydratedComponents.java | 3 +-
.../beam/sdk/util/construction/SdkComponents.java | 39 +-
.../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 +
.../apache/beam/sdk/schemas/SchemaUtilsTest.java | 43 +
.../sdk/transforms/GroupByEncryptedKeyTest.java | 2 +-
.../org/apache/beam/sdk/transforms/WatchTest.java | 348 +++++-
.../java/org/apache/beam/sdk/util/SecretTest.java | 179 ++-
.../util/construction/CoderTranslationTest.java | 38 +-
.../service/WindowIntoTransformProvider.java | 1 +
.../beam/sdk/extensions/arrow/ArrowConversion.java | 61 ++
.../sdk/extensions/arrow/ArrowConversionTest.java | 33 +
.../extensions/avro/AvroGenericCoderRegistrar.java | 18 +
.../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 +-
.../schemaio-expansion-service/build.gradle | 6 +
sdks/java/extensions/sql/iceberg/build.gradle | 4 +-
.../provider/iceberg/BeamSqlCliIcebergTest.java | 6 +-
.../meta/provider/iceberg/IcebergReadWriteIT.java | 3 -
.../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/BeamCalcRel.java | 40 +
.../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 -
.../extensions/sql/impl/utils/CalciteUtils.java | 8 +-
.../sql/meta/provider/mongodb/package-info.java | 4 -
.../sql/meta/provider/pubsub/package-info.java | 4 -
.../sdk/extensions/sql/BeamComplexTypeTest.java | 32 +
.../sql/impl/BeamSqlEnvRegisterOperatorTest.java | 100 ++
.../extensions/sql/impl/rel/BeamCalcRelTest.java | 20 +
.../beam/fn/harness/state/StateBackedIterable.java | 22 +
.../KinesisReadSchemaTransformProvider.java | 315 ++++++
.../KinesisWriteSchemaTransformProvider.java | 280 +++++
.../KinesisSchemaTransformProviderTest.java | 221 ++++
sdks/java/io/{jms => arrow-flight}/build.gradle | 38 +-
.../beam/sdk/io/arrowflight/ArrowFlightIO.java | 840 +++++++++++++++
.../beam/sdk/io/arrowflight}/package-info.java | 17 +-
.../beam/sdk/io/arrowflight/ArrowFlightIOTest.java | 330 ++++++
.../beam/sdk/io/cassandra/ConnectionManager.java | 17 +-
.../beam/sdk/io/cassandra/CassandraIOTest.java | 30 +
sdks/java/io/debezium/build.gradle | 17 +-
.../io/debezium/expansion-service/build.gradle | 6 +-
sdks/java/io/debezium/src/README.md | 8 +-
.../org/apache/beam/io/debezium/DebeziumIO.java | 15 +-
.../DebeziumReadSchemaTransformProvider.java | 96 +-
.../beam/io/debezium/KafkaSourceConsumerFn.java | 5 +
.../io/debezium/DebeziumIOMySqlConnectorIT.java | 6 +-
.../debezium/DebeziumIOPostgresSqlConnectorIT.java | 4 +-
.../apache/beam/io/debezium/DebeziumIOTest.java | 3 +-
.../DebeziumReadSchemaTransformProviderTest.java | 139 +++
.../debezium/DebeziumReadSchemaTransformTest.java | 29 +-
sdks/java/io/delta/build.gradle | 21 +-
.../beam/sdk/io/delta/CreateCDCReadTasksDoFn.java | 296 +++++
.../beam/sdk/io/delta/CreateReadTasksDoFn.java | 23 +-
.../apache/beam/sdk/io/delta/DeltaCDCReadTask.java | 125 +++
.../beam/sdk/io/delta/DeltaCDCSourceDoFn.java | 367 +++++++
...va => DeltaCdcReadSchemaTransformProvider.java} | 80 +-
.../java/org/apache/beam/sdk/io/delta/DeltaIO.java | 213 +++-
.../io/delta/DeltaReadSchemaTransformProvider.java | 12 +-
.../apache/beam/sdk/io/delta/DeltaSourceDoFn.java | 2 +-
.../org/apache/beam/sdk/io/delta/DeltaIOIT.java | 273 ++++-
.../io/delta/{DeltaIOIT.java => DeltaIOS3IT.java} | 132 ++-
.../org/apache/beam/sdk/io/delta/DeltaIOTest.java | 1139 ++++++++++++++++++--
.../DeltaReadSchemaTransformProviderTest.java | 51 +
.../beam/sdk/io/delta/DeltaWriteTestUtils.java | 401 +++++++
sdks/java/io/expansion-service/build.gradle | 1 +
sdks/java/io/google-cloud-platform/build.gradle | 39 +-
.../beam/sdk/io/gcp/bigquery/BigQueryHelpers.java | 100 +-
.../beam/sdk/io/gcp/bigquery/BigQueryIO.java | 30 +-
.../io/gcp/bigquery/BigQueryStorageSourceBase.java | 7 +-
.../gcp/bigquery/BigQueryStorageTableSource.java | 3 +-
.../sdk/io/gcp/bigquery/BigQueryTableSource.java | 5 +
.../beam/sdk/io/gcp/bigquery/BigQueryUtils.java | 6 +
.../sdk/io/gcp/bigtable/BigtableServiceImpl.java | 8 +
.../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 +-
.../changestreams/action/ActionFactory.java | 7 +-
.../action/QueryChangeStreamAction.java | 44 +-
.../gcp/spanner/changestreams/dao/DaoFactory.java | 16 +-
.../dofn/CleanUpReadChangeStreamDoFn.java | 7 +
.../dofn/DetectNewPartitionsDoFn.java | 5 +-
.../spanner/changestreams/dofn/InitializeDoFn.java | 7 +
.../dofn/ReadChangeStreamPartitionDoFn.java | 9 +-
.../apache/beam/sdk/io/gcp/GcpApiSurfaceTest.java | 1 +
.../sdk/io/gcp/bigquery/BigQueryHelpersTest.java | 348 +++++-
.../bigquery/BigQueryIOIcebergManagedTableIT.java | 408 +++++++
.../io/gcp/bigquery/BigQueryIOStorageReadTest.java | 63 ++
.../sdk/io/gcp/bigquery/BigQueryUtilsTest.java | 42 +-
....java => StorageApiSinkSchemaUpdateITBase.java} | 61 +-
...torageApiSinkSchemaUpdateWithInputSchemaIT.java | 50 +
...ageApiSinkSchemaUpdateWithoutInputSchemaIT.java | 51 +
.../io/gcp/bigtable/BigtableServiceImplTest.java | 81 ++
.../sdk/io/gcp/spanner/SpannerIOWriteTest.java | 14 +-
.../action/QueryChangeStreamActionTest.java | 10 +-
.../dofn/ReadChangeStreamPartitionDoFnTest.java | 6 +-
sdks/java/io/iceberg/build.gradle | 21 +-
.../org/apache/beam/sdk/io/iceberg/AddFiles.java | 255 ++---
.../beam/sdk/io/iceberg/BoundedAsyncTasks.java | 113 ++
.../beam/sdk/io/iceberg/DynamicDestinations.java | 12 +-
.../IcebergCdcReadSchemaTransformProvider.java | 31 +-
.../org/apache/beam/sdk/io/iceberg/IcebergIO.java | 102 +-
.../beam/sdk/io/iceberg/IcebergScanConfig.java | 115 +-
.../apache/beam/sdk/io/iceberg/IcebergUtils.java | 320 +++---
.../beam/sdk/io/iceberg/IncrementalScanSource.java | 100 --
.../beam/sdk/io/iceberg/NameMappingUtils.java | 215 ++++
.../io/iceberg/OneTableDynamicDestinations.java | 33 +-
.../apache/beam/sdk/io/iceberg/ParquetFooters.java | 51 +
.../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/SerializableTableSpec.java | 382 +++++++
.../apache/beam/sdk/io/iceberg/SideInputTable.java | 368 +++++++
.../beam/sdk/io/iceberg/WatchForSnapshots.java | 190 ----
.../io/iceberg/WritePartitionedRowsToFiles.java | 16 +-
.../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 ++++
.../beam/sdk/io/iceberg/cdc/CdcRowDescriptor.java | 89 ++
.../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 | 89 +-
.../apache/beam/sdk/io/iceberg/AddFilesTest.java | 255 ++++-
.../iceberg/BigQueryManagedTableCrossEngineIT.java | 183 ++++
.../beam/sdk/io/iceberg/BoundedAsyncTasksTest.java | 220 ++++
.../IcebergCdcReadSchemaTransformProviderTest.java | 114 ++
.../beam/sdk/io/iceberg/IcebergIOReadTest.java | 48 +
.../beam/sdk/io/iceberg/IcebergIOWriteTest.java | 54 +
.../beam/sdk/io/iceberg/IcebergScanConfigTest.java | 270 +++++
.../IcebergSchemaTransformTranslationTest.java | 1 +
.../beam/sdk/io/iceberg/IcebergUtilsTest.java | 53 +-
.../IcebergWriteSchemaTransformProviderTest.java | 76 +-
.../beam/sdk/io/iceberg/NameMappingUtilsTest.java | 425 ++++++++
.../beam/sdk/io/iceberg/ParquetFootersTest.java | 113 ++
.../beam/sdk/io/iceberg/PartitionUtilsTest.java | 10 +-
.../sdk/io/iceberg/RecordWriterManagerTest.java | 24 +-
.../sdk/io/iceberg/SerializableDataFileTest.java | 156 ++-
.../sdk/io/iceberg/SerializableTableSpecTest.java | 348 ++++++
.../beam/sdk/io/iceberg/SideInputTableTest.java | 239 ++++
.../catalog/BigQueryMetastoreCatalogIT.java | 11 +
.../io/iceberg/catalog/IcebergCatalogBaseIT.java | 551 +++++++++-
.../sdk/io/iceberg/catalog/RESTCatalogBLMSIT.java | 22 +-
.../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 | 600 +++++++++++
.../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 +-
.../java/org/apache/beam/sdk/io/jdbc/JdbcIO.java | 70 ++
.../io/jdbc/JdbcReadSchemaTransformProvider.java | 34 +-
.../beam/sdk/io/jdbc/JdbcSchemaIOProvider.java | 6 +
.../io/jdbc/JdbcWriteSchemaTransformProvider.java | 34 +-
.../org/apache/beam/sdk/io/jdbc/JdbcIOTest.java | 45 +
sdks/java/io/jms/build.gradle | 6 +
.../io/jms/BeamGenericJmsConnectionFactory.java | 42 +
.../beam/sdk/io/jms/ConnectionConfiguration.java | 252 +++++
.../java/org/apache/beam/sdk/io/jms/JmsIO.java | 30 +
.../sdk/io/jms/JmsReadSchemaTransformProvider.java | 227 ++++
.../io/jms/JmsWriteSchemaTransformProvider.java | 191 ++++
.../sdk/io/jms/ConnectionConfigurationTest.java | 159 +++
.../java/org/apache/beam/sdk/io/jms/JmsIOTest.java | 176 +--
.../org/apache/beam/sdk/io/jms/JmsLocalTest.java | 245 +++++
.../sdk/io/jms/JmsSchemaTransformProviderTest.java | 253 +++++
sdks/java/io/kafka/build.gradle | 3 +
.../java/org/apache/beam/sdk/io/kafka/KafkaIO.java | 163 ++-
.../KafkaIOReadImplementationCompatibility.java | 6 +
.../kafka/KafkaWriteSchemaTransformProvider.java | 41 +-
...KafkaIOReadImplementationCompatibilityTest.java | 22 +
.../KafkaWriteSchemaTransformProviderTest.java | 27 +-
.../io/messaging-expansion-service/build.gradle | 3 +
.../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 +-
sdks/java/io/snowflake/build.gradle | 2 +
.../apache/beam/sdk/io/snowflake/SnowflakeIO.java | 3 +-
.../SnowflakeReadSchemaTransformProvider.java | 289 +++++
.../snowflake/SnowflakeSchemaTransformUtils.java | 342 ++++++
.../io/snowflake/SnowflakeWriteConfiguration.java | 227 ++++
.../SnowflakeWriteSchemaTransformProvider.java | 179 +++
.../io/snowflake/crosslanguage/package-info.java | 4 -
.../SnowflakeReadSchemaTransformProviderTest.java | 311 ++++++
.../SnowflakeWriteSchemaTransformProviderTest.java | 494 +++++++++
.../beam/sdk/io/solace/broker/BrokerResponse.java | 13 +-
.../sdk/io/solace/broker/BrokerResponseTest.java | 66 ++
.../java/org/apache/beam/sdk/managed/Managed.java | 5 +
sdks/python/apache_beam/coders/row_coder_test.py | 54 +
sdks/python/apache_beam/dataframe/io.py | 3 +-
sdks/python/apache_beam/dataframe/io_test.py | 21 +-
.../anomaly_detection_pipeline/setup.py | 2 +-
.../online_clustering/clustering_pipeline/setup.py | 2 +-
.../transforms/elementwise/enrichment_test.py | 2 +-
.../io/external/xlang_debeziumio_it_test.py | 2 +-
.../io/external/xlang_jdbcio_it_test.py | 113 ++
.../apache_beam/io/external/xlang_jmsio_it_test.py | 376 +++++++
.../io/external/xlang_mqttio_it_test.py | 66 +-
sdks/python/apache_beam/io/fileio.py | 305 ++++--
sdks/python/apache_beam/io/fileio_test.py | 345 ++++++
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 | 116 +-
.../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 | 197 +++-
sdks/python/apache_beam/io/gcp/bigquery_tools.py | 8 +
.../apache_beam/io/gcp/bigquery_write_it_test.py | 71 +-
sdks/python/apache_beam/io/gcp/bigtableio.py | 8 +-
sdks/python/apache_beam/io/gcp/bigtableio_test.py | 55 +
.../apache_beam/io/gcp/gcsfilesystem_test.py | 4 +-
.../apache_beam/io/gcp/healthcare/dicomclient.py | 5 +-
.../apache_beam/io/gcp/pubsub_io_perf_test.py | 13 +-
sdks/python/apache_beam/io/gcp/pubsub_test.py | 42 +
sdks/python/apache_beam/io/iobase.py | 12 +-
sdks/python/apache_beam/io/iobase_test.py | 21 +
sdks/python/apache_beam/io/jdbc.py | 48 +-
sdks/python/apache_beam/io/textio_test.py | 7 +-
sdks/python/apache_beam/io/watch.py | 254 ++++-
sdks/python/apache_beam/io/watch_test.py | 466 +++++++-
.../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 | 30 +-
.../ml/inference/vertex_ai_model_monitoring_v2.py | 458 ++++++++
.../vertex_ai_model_monitoring_v2_it_test.py | 485 +++++++++
.../vertex_ai_model_monitoring_v2_test.py | 655 +++++++++++
.../python/apache_beam/options/pipeline_options.py | 33 +
.../apache_beam/options/pipeline_options_test.py | 25 +
sdks/python/apache_beam/portability/common_urns.py | 1 +
.../runners/dataflow/internal/apiclient.py | 6 +-
.../runners/dataflow/internal/apiclient_test.py | 19 +
.../runners/dataflow/internal/clients/README.txt | 11 -
.../runners/dataflow/internal/clients/__init__.py | 16 -
.../apache_beam/runners/dataflow/internal/names.py | 2 +-
.../runners/direct/transform_evaluator.py | 5 +
.../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/fn_api_runner/execution.py | 4 +-
.../portability/fn_api_runner/translations.py | 4 +-
.../kafka_streams_java_job_server_test.py | 140 +++
.../runners/portability/kafka_streams_runner.py | 137 +++
.../portability/kafka_streams_runner_test.py | 285 +++++
.../runners/portability/spark_runner_test.py | 20 -
.../apache_beam/runners/worker/sdk_worker_main.py | 20 +-
.../runners/worker/sdk_worker_main_test.py | 49 +
.../testing/benchmarks/chicago_taxi/run_chicago.sh | 6 +-
.../apache_beam/testing/pubsub_test_context.py | 179 +++
.../testing/pubsub_test_context_test.py | 149 +++
.../apache_beam/transforms/async_dofn_test.py | 15 +-
sdks/python/apache_beam/transforms/combiners.py | 58 +
.../apache_beam/transforms/combiners_test.py | 54 +
.../transforms/managed_iceberg_it_test.py | 7 +-
.../apache_beam/transforms/periodicsequence.py | 13 +-
.../transforms/periodicsequence_test.py | 61 ++
sdks/python/apache_beam/transforms/util.py | 245 +----
sdks/python/apache_beam/transforms/util_test.py | 157 +--
sdks/python/apache_beam/typehints/schemas.py | 152 ++-
sdks/python/apache_beam/typehints/schemas_test.py | 120 +++
sdks/python/apache_beam/utils/secret.py | 464 ++++++++
sdks/python/apache_beam/utils/secret_test.py | 454 ++++++++
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/debezium.yaml | 50 +
.../yaml/extended_tests/databases/iceberg.yaml | 3 +
.../databases/jdbc_secret_manager.yaml | 59 +
.../yaml/extended_tests/databases/snowflake.yaml | 90 ++
.../extended_tests/e2e/delta_lake_to_iceberg.yaml | 71 ++
.../yaml/extended_tests/messaging/kinesis.yaml | 83 ++
sdks/python/apache_beam/yaml/integration_tests.py | 375 ++++++-
sdks/python/apache_beam/yaml/standard_io.yaml | 184 ++++
sdks/python/apache_beam/yaml/tests/ibm_mq.yaml | 62 ++
sdks/python/apache_beam/yaml/tests/jms.yaml | 56 +
sdks/python/apache_beam/yaml/yaml_io.py | 143 ++-
sdks/python/apache_beam/yaml/yaml_io_test.py | 219 ++++
sdks/python/apache_beam/yaml/yaml_provider.py | 6 +-
sdks/python/build.gradle | 11 +-
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 | 6 +-
sdks/python/scripts/run_integration_test.sh | 11 +-
sdks/python/scripts/run_snapshot_publish.sh | 4 +-
sdks/python/setup.py | 37 +-
sdks/python/test-suites/dataflow/common.gradle | 18 +-
sdks/python/test-suites/direct/build.gradle | 3 +-
sdks/python/test-suites/direct/common.gradle | 21 +-
sdks/python/test-suites/portable/common.gradle | 39 +
sdks/python/test-suites/tox/py310/build.gradle | 27 +-
sdks/python/tox.ini | 16 +-
sdks/standard_expansion_services.yaml | 5 +
sdks/standard_external_transforms.yaml | 235 +++-
sdks/typescript/container/boot.go | 4 +-
sdks/typescript/package.json | 2 +-
settings.gradle.kts | 15 +-
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 | 44 +-
.../en/documentation/io/developing-io-python.md | 239 +++-
.../site/content/en/documentation/io/managed-io.md | 172 +++
.../site/content/en/documentation/runners/flink.md | 23 +-
.../en/documentation/runners/kafkastreams.md | 248 +++++
.../site/content/en/documentation/runners/spark.md | 4 +-
.../sdks/python-multi-language-pipelines.md | 2 +-
.../sdks/python-pipeline-dependencies.md | 41 +-
.../site/content/en/get-started/quickstart-java.md | 4 +-
.../site/content/en/get-started/quickstart-py.md | 17 +
website/www/site/data/capability_matrix.yaml | 155 +++
.../layouts/partials/section-menu/en/runners.html | 1 +
.../documentation/capability-matrix-big.html | 36 +-
.../documentation/capability-matrix-single.html | 37 +-
website/www/yarn.lock | 6 +-
815 files changed, 56201 insertions(+), 6765 deletions(-)
copy .github/trigger_files/{beam_PostCommit_Java_Nexmark_Spark.json =>
beam_PostCommit_Java_ValidatesRunner_Spark4.json} (100%)
copy .github/trigger_files/{beam_PostCommit_Java_ValidatesRunner_Spark.json =>
beam_PostCommit_Python_Arm.json} (100%)
copy .github/trigger_files/{beam_PostCommit_Java_Delta_IO_Dataflow.json =>
beam_PreCommit_Java_Kafka_Streams_Runner.json} (100%)
delete mode 100644 .github/workflows/beam_Infrastructure_AuditUnmanagedKeys.yml
copy .github/workflows/{beam_PostCommit_XVR_GoUsingJava_Dataflow.yml =>
beam_PreCommit_Java_Kafka_Streams_Runner.yml} (63%)
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
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%)
copy
runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/{WorkItemCancelledException.java
=> WorkCancellingException.java} (52%)
copy
runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/{BoundedQueueExecutorWorkHandle.java
=> FailedWorkHandler.java} (83%)
copy
runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/{BoundedQueueExecutorWorkHandle.java
=> MultiKeyCommitValidationException.java} (74%)
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 runners/kafka-streams/build.gradle
create mode 100644 runners/kafka-streams/job-server/build.gradle
create mode 100644 runners/kafka-streams/measurement/build.gradle
create mode 100644 runners/kafka-streams/measurement/docker-compose.yml
create mode 100644
runners/kafka-streams/measurement/src/main/java/org/apache/beam/runners/kafka/streams/measurement/RescalingMeasurement.java
copy {sdks/java/core/src/main/java/org/apache/beam/sdk/schemas/logicaltypes =>
runners/kafka-streams/measurement/src/main/java/org/apache/beam/runners/kafka/streams/measurement}/package-info.java
(67%)
rename
runners/{google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/KeyTokenInvalidException.java
=> kafka-streams/proto/build.gradle} (52%)
create mode 100644
runners/kafka-streams/proto/src/main/proto/kafka_streams_payload.proto
create mode 100644
runners/kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/KafkaStreamsJobInvoker.java
create mode 100644
runners/kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/KafkaStreamsJobServerDriver.java
create mode 100644
runners/kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/KafkaStreamsPipelineOptions.java
create mode 100644
runners/kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/KafkaStreamsPipelineResult.java
create mode 100644
runners/kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/KafkaStreamsPipelineRunner.java
create mode 100644
runners/kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/KafkaStreamsPortablePipelineResult.java
create mode 100644
runners/kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/KafkaStreamsRunner.java
create mode 100644
runners/kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/KafkaStreamsRunnerRegistrar.java
create mode 100644
runners/kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/KafkaStreamsTopicManager.java
copy
{sdks/java/io/snowflake/src/main/java/org/apache/beam/sdk/io/snowflake/crosslanguage
=>
runners/kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams}/package-info.java
(77%)
create mode 100644
runners/kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/translation/EmptyBoundedSource.java
create mode 100644
runners/kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/translation/ExecutableStageProcessor.java
create mode 100644
runners/kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/translation/ExecutableStageTranslator.java
create mode 100644
runners/kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/translation/FlattenProcessor.java
create mode 100644
runners/kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/translation/FlattenTranslator.java
create mode 100644
runners/kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/translation/GroupByKeyBroadcastPartitioner.java
create mode 100644
runners/kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/translation/GroupByKeyTranslator.java
create mode 100644
runners/kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/translation/ImpulseProcessor.java
create mode 100644
runners/kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/translation/ImpulseTranslator.java
create mode 100644
runners/kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/translation/KStreamsPayload.java
create mode 100644
runners/kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/translation/KStreamsPayloadSerde.java
create mode 100644
runners/kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/translation/KafkaStreamsExecutableStageContextFactory.java
create mode 100644
runners/kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/translation/KafkaStreamsPipelineTranslator.java
create mode 100644
runners/kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/translation/KafkaStreamsStateInternals.java
create mode 100644
runners/kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/translation/KafkaStreamsTimerInternals.java
create mode 100644
runners/kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/translation/KafkaStreamsTranslationContext.java
create mode 100644
runners/kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/translation/PTransformTranslator.java
create mode 100644
runners/kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/translation/ReadProcessor.java
create mode 100644
runners/kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/translation/ReadTranslator.java
create mode 100644
runners/kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/translation/RedistributeTranslator.java
create mode 100644
runners/kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/translation/ShuffleByKeyProcessor.java
create mode 100644
runners/kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/translation/StageOutputProcessor.java
create mode 100644
runners/kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/translation/StoreKeys.java
create mode 100644
runners/kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/translation/TerminationReporter.java
create mode 100644
runners/kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/translation/TerminationTracker.java
create mode 100644
runners/kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/translation/UnboundedReadProcessor.java
create mode 100644
runners/kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/translation/WatermarkAggregator.java
create mode 100644
runners/kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/translation/WatermarkManager.java
create mode 100644
runners/kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/translation/WatermarkPayload.java
create mode 100644
runners/kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/translation/WindowedGroupByKeyProcessor.java
copy
{sdks/java/io/snowflake/src/main/java/org/apache/beam/sdk/io/snowflake/crosslanguage
=>
runners/kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/translation}/package-info.java
(77%)
create mode 100644
runners/kafka-streams/src/test/java/org/apache/beam/runners/kafka/streams/KafkaStreamsJobServerDriverTest.java
create mode 100644
runners/kafka-streams/src/test/java/org/apache/beam/runners/kafka/streams/KafkaStreamsPipelineOptionsTest.java
create mode 100644
runners/kafka-streams/src/test/java/org/apache/beam/runners/kafka/streams/KafkaStreamsPipelineRunnerConfigTest.java
create mode 100644
runners/kafka-streams/src/test/java/org/apache/beam/runners/kafka/streams/KafkaStreamsPortablePipelineResultTest.java
create mode 100644
runners/kafka-streams/src/test/java/org/apache/beam/runners/kafka/streams/KafkaStreamsRunnerBrokerIT.java
create mode 100644
runners/kafka-streams/src/test/java/org/apache/beam/runners/kafka/streams/KafkaStreamsRunnerTest.java
create mode 100644
runners/kafka-streams/src/test/java/org/apache/beam/runners/kafka/streams/KafkaStreamsTestRunner.java
create mode 100644
runners/kafka-streams/src/test/java/org/apache/beam/runners/kafka/streams/MultiOutputStageTest.java
create mode 100644
runners/kafka-streams/src/test/java/org/apache/beam/runners/kafka/streams/TestKafkaStreamsRunner.java
create mode 100644
runners/kafka-streams/src/test/java/org/apache/beam/runners/kafka/streams/TestKafkaStreamsRunnerTest.java
create mode 100644
runners/kafka-streams/src/test/java/org/apache/beam/runners/kafka/streams/translation/BundleBoundaryTest.java
create mode 100644
runners/kafka-streams/src/test/java/org/apache/beam/runners/kafka/streams/translation/ChainedExecutableStageTest.java
create mode 100644
runners/kafka-streams/src/test/java/org/apache/beam/runners/kafka/streams/translation/CreateTest.java
create mode 100644
runners/kafka-streams/src/test/java/org/apache/beam/runners/kafka/streams/translation/ExecutableStageProcessorWatermarkTest.java
create mode 100644
runners/kafka-streams/src/test/java/org/apache/beam/runners/kafka/streams/translation/ExecutableStageTranslatorTest.java
create mode 100644
runners/kafka-streams/src/test/java/org/apache/beam/runners/kafka/streams/translation/FixedWindowGroupByKeyTest.java
create mode 100644
runners/kafka-streams/src/test/java/org/apache/beam/runners/kafka/streams/translation/FlattenParallelismTest.java
create mode 100644
runners/kafka-streams/src/test/java/org/apache/beam/runners/kafka/streams/translation/FlattenTest.java
create mode 100644
runners/kafka-streams/src/test/java/org/apache/beam/runners/kafka/streams/translation/GroupByKeyTest.java
create mode 100644
runners/kafka-streams/src/test/java/org/apache/beam/runners/kafka/streams/translation/ImpulseTranslatorTest.java
create mode 100644
runners/kafka-streams/src/test/java/org/apache/beam/runners/kafka/streams/translation/KStreamsPayloadSerdeTest.java
create mode 100644
runners/kafka-streams/src/test/java/org/apache/beam/runners/kafka/streams/translation/KafkaStreamsPipelineTranslatorTest.java
create mode 100644
runners/kafka-streams/src/test/java/org/apache/beam/runners/kafka/streams/translation/KafkaStreamsTimerInternalsTest.java
create mode 100644
runners/kafka-streams/src/test/java/org/apache/beam/runners/kafka/streams/translation/MetricsAcrossBundlesTest.java
create mode 100644
runners/kafka-streams/src/test/java/org/apache/beam/runners/kafka/streams/translation/MetricsTest.java
create mode 100644
runners/kafka-streams/src/test/java/org/apache/beam/runners/kafka/streams/translation/ReadTest.java
create mode 100644
runners/kafka-streams/src/test/java/org/apache/beam/runners/kafka/streams/translation/SharedTestCollector.java
create mode 100644
runners/kafka-streams/src/test/java/org/apache/beam/runners/kafka/streams/translation/ShuffleByKeyProcessorTest.java
create mode 100644
runners/kafka-streams/src/test/java/org/apache/beam/runners/kafka/streams/translation/StageOutputProcessorTest.java
create mode 100644
runners/kafka-streams/src/test/java/org/apache/beam/runners/kafka/streams/translation/StandardWindowFnTranslationTest.java
create mode 100644
runners/kafka-streams/src/test/java/org/apache/beam/runners/kafka/streams/translation/TerminationTrackerTest.java
create mode 100644
runners/kafka-streams/src/test/java/org/apache/beam/runners/kafka/streams/translation/UnboundedReadTest.java
create mode 100644
runners/kafka-streams/src/test/java/org/apache/beam/runners/kafka/streams/translation/WatermarkAggregatorTest.java
create mode 100644
runners/kafka-streams/src/test/java/org/apache/beam/runners/kafka/streams/translation/WatermarkManagerTest.java
create mode 100644
runners/kafka-streams/src/test/java/org/apache/beam/runners/kafka/streams/translation/WatermarkPropagationTest.java
create mode 100644
runners/spark/src/main/java/org/apache/beam/runners/spark/structuredstreaming/translation/batch/StatefulDoFnGroupFunction.java
create mode 100644
runners/spark/src/main/java/org/apache/beam/runners/spark/structuredstreaming/translation/batch/StatefulParDoTranslatorBatch.java
create mode 100644
runners/spark/src/test/java/org/apache/beam/runners/spark/structuredstreaming/translation/batch/StatefulParDoExecutionTest.java
create mode 100644
runners/spark/src/test/java/org/apache/beam/runners/spark/structuredstreaming/translation/batch/StatefulParDoTranslatorBatchTest.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/runners/prism/internal/engine/elementmanager_continuation_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/{container/tools/pipeline_options_test.go =>
pkg/beam/transforms/batch/batch_test.go} (51%)
create mode 100644 sdks/go/pkg/beam/transforms/batch/doc.go
create mode 100644 sdks/go/pkg/beam/transforms/batch/size.go
create mode 100644 sdks/go/pkg/beam/transforms/batch/size_test.go
create mode 100644
sdks/java/core/src/main/java/org/apache/beam/sdk/util/RawSecret.java
create mode 100644
sdks/java/extensions/sql/src/test/java/org/apache/beam/sdk/extensions/sql/impl/BeamSqlEnvRegisterOperatorTest.java
create mode 100644
sdks/java/io/amazon-web-services2/src/main/java/org/apache/beam/sdk/io/aws2/kinesis/KinesisReadSchemaTransformProvider.java
create mode 100644
sdks/java/io/amazon-web-services2/src/main/java/org/apache/beam/sdk/io/aws2/kinesis/KinesisWriteSchemaTransformProvider.java
create mode 100644
sdks/java/io/amazon-web-services2/src/test/java/org/apache/beam/sdk/io/aws2/kinesis/KinesisSchemaTransformProviderTest.java
copy sdks/java/io/{jms => arrow-flight}/build.gradle (56%)
create mode 100644
sdks/java/io/arrow-flight/src/main/java/org/apache/beam/sdk/io/arrowflight/ArrowFlightIO.java
copy sdks/java/{core/src/main/java/org/apache/beam/sdk/schemas/annotations =>
io/arrow-flight/src/main/java/org/apache/beam/sdk/io/arrowflight}/package-info.java
(64%)
create mode 100644
sdks/java/io/arrow-flight/src/test/java/org/apache/beam/sdk/io/arrowflight/ArrowFlightIOTest.java
create mode 100644
sdks/java/io/debezium/src/test/java/org/apache/beam/io/debezium/DebeziumReadSchemaTransformProviderTest.java
create mode 100644
sdks/java/io/delta/src/main/java/org/apache/beam/sdk/io/delta/CreateCDCReadTasksDoFn.java
create mode 100644
sdks/java/io/delta/src/main/java/org/apache/beam/sdk/io/delta/DeltaCDCReadTask.java
create mode 100644
sdks/java/io/delta/src/main/java/org/apache/beam/sdk/io/delta/DeltaCDCSourceDoFn.java
copy
sdks/java/io/delta/src/main/java/org/apache/beam/sdk/io/delta/{DeltaReadSchemaTransformProvider.java
=> DeltaCdcReadSchemaTransformProvider.java} (55%)
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
create mode 100644
sdks/java/io/google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/bigquery/BigQueryIOIcebergManagedTableIT.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
create mode 100644
sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/BoundedAsyncTasks.java
delete mode 100644
sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/IncrementalScanSource.java
create mode 100644
sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/NameMappingUtils.java
create mode 100644
sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/ParquetFooters.java
delete mode 100644
sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/ReadFromTasks.java
create mode 100644
sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/SerializableTableSpec.java
create mode 100644
sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/SideInputTable.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
create mode 100644
sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/cdc/CdcRowDescriptor.java
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/BigQueryManagedTableCrossEngineIT.java
create mode 100644
sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/BoundedAsyncTasksTest.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/NameMappingUtilsTest.java
create mode 100644
sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/ParquetFootersTest.java
create mode 100644
sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/SerializableTableSpecTest.java
create mode 100644
sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/SideInputTableTest.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
create mode 100644
sdks/java/io/jms/src/main/java/org/apache/beam/sdk/io/jms/BeamGenericJmsConnectionFactory.java
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/main/java/org/apache/beam/sdk/io/jms/JmsReadSchemaTransformProvider.java
create mode 100644
sdks/java/io/jms/src/main/java/org/apache/beam/sdk/io/jms/JmsWriteSchemaTransformProvider.java
create mode 100644
sdks/java/io/jms/src/test/java/org/apache/beam/sdk/io/jms/ConnectionConfigurationTest.java
create mode 100644
sdks/java/io/jms/src/test/java/org/apache/beam/sdk/io/jms/JmsLocalTest.java
create mode 100644
sdks/java/io/jms/src/test/java/org/apache/beam/sdk/io/jms/JmsSchemaTransformProviderTest.java
create mode 100644
sdks/java/io/snowflake/src/main/java/org/apache/beam/sdk/io/snowflake/SnowflakeReadSchemaTransformProvider.java
create mode 100644
sdks/java/io/snowflake/src/main/java/org/apache/beam/sdk/io/snowflake/SnowflakeSchemaTransformUtils.java
create mode 100644
sdks/java/io/snowflake/src/main/java/org/apache/beam/sdk/io/snowflake/SnowflakeWriteConfiguration.java
create mode 100644
sdks/java/io/snowflake/src/main/java/org/apache/beam/sdk/io/snowflake/SnowflakeWriteSchemaTransformProvider.java
create mode 100644
sdks/java/io/snowflake/src/test/java/org/apache/beam/sdk/io/snowflake/SnowflakeReadSchemaTransformProviderTest.java
create mode 100644
sdks/java/io/snowflake/src/test/java/org/apache/beam/sdk/io/snowflake/SnowflakeWriteSchemaTransformProviderTest.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/ml/inference/vertex_ai_model_monitoring_v2.py
create mode 100644
sdks/python/apache_beam/ml/inference/vertex_ai_model_monitoring_v2_it_test.py
create mode 100644
sdks/python/apache_beam/ml/inference/vertex_ai_model_monitoring_v2_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
create mode 100644
sdks/python/apache_beam/runners/portability/kafka_streams_java_job_server_test.py
create mode 100644
sdks/python/apache_beam/runners/portability/kafka_streams_runner.py
create mode 100644
sdks/python/apache_beam/runners/portability/kafka_streams_runner_test.py
create mode 100644 sdks/python/apache_beam/testing/pubsub_test_context.py
create mode 100644 sdks/python/apache_beam/testing/pubsub_test_context_test.py
create mode 100644 sdks/python/apache_beam/utils/secret.py
create mode 100644 sdks/python/apache_beam/utils/secret_test.py
create mode 100644
sdks/python/apache_beam/yaml/extended_tests/databases/debezium.yaml
create mode 100644
sdks/python/apache_beam/yaml/extended_tests/databases/jdbc_secret_manager.yaml
create mode 100644
sdks/python/apache_beam/yaml/extended_tests/databases/snowflake.yaml
create mode 100644
sdks/python/apache_beam/yaml/extended_tests/e2e/delta_lake_to_iceberg.yaml
create mode 100644
sdks/python/apache_beam/yaml/extended_tests/messaging/kinesis.yaml
create mode 100644 sdks/python/apache_beam/yaml/tests/ibm_mq.yaml
create mode 100644 sdks/python/apache_beam/yaml/tests/jms.yaml
create mode 100644 sdks/python/container/profiler_test.go
create mode 100644
website/www/site/content/en/documentation/runners/kafkastreams.md