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

Amar3tto pushed a change to branch snowflakeio-yaml
in repository https://gitbox.apache.org/repos/asf/beam.git


    omit 25a48019aea Add license
    omit 6f35d263316 Add extended test
    omit fe70ef495bb Refactoring
    omit fd9679d2804 Fix imports
    omit c2d8c71af61 Fix bytes
    omit 5bd11daeca4 Add Snowflake Read and Streaming Write
    omit 3f39865016f Add Snowflake YAML write transform
     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 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

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   (25a48019aea)
            \
             N -- N -- N   refs/heads/snowflakeio-yaml (27c4e6395cb)

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 +
 .../IO_Iceberg_Integration_Tests.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 +-
 ...am_PostCommit_Python_ValidatesRunner_Spark.json |   3 +-
 .../beam_PostCommit_Python_Xlang_Gcp_Direct.json   |   2 +-
 ...m_PostCommit_Python_Xlang_Messaging_Direct.json |   2 +-
 .../beam_PostCommit_Yaml_Xlang_Direct.json         |   2 +-
 ... beam_PreCommit_Java_Kafka_Streams_Runner.json} |   0
 .github/workflows/README.md                        |   2 +-
 .../beam_Infrastructure_AuditUnmanagedKeys.yml     |  76 --
 .../beam_Infrastructure_PolicyEnforcer.yml         |  13 +-
 .../beam_PostCommit_Yaml_Xlang_Direct.yml          |   4 +-
 .../workflows/beam_PostRelease_NightlySnapshot.yml |   9 +-
 .github/workflows/beam_PreCommit_Java.yml          |   1 +
 ...> beam_PreCommit_Java_Kafka_Streams_Runner.yml} |  87 ++-
 .github/workflows/beam_Release_NightlySnapshot.yml |   1 +
 .github/workflows/build_release_candidate.yml      |  12 +-
 .github/workflows/codeql.yml                       |   6 +-
 .../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 +-
 .test-infra/dataproc/flink_cluster.sh              |   6 +
 .test-infra/tools/stale_cleaner.py                 |   5 +-
 .test-infra/tools/test_stale_cleaner.py            |   7 +-
 CHANGES.md                                         |  20 +-
 build.gradle.kts                                   |   4 +
 .../org/apache/beam/gradle/BeamModulePlugin.groovy |   3 +-
 .../beam/examples/subprocess/utils/FileUtils.java  |  30 +-
 .../examples/subprocess/utils/FileUtilsTest.java   |  80 +++
 infra/enforcement/README.md                        |  18 +-
 infra/enforcement/account_keys.py                  |  76 +-
 infra/enforcement/iam.py                           |  83 ++-
 infra/enforcement/sending.py                       |  57 +-
 infra/iam/users.yml                                |  13 +-
 it/mongodb/build.gradle                            |   1 +
 .../beam/it/mongodb/MongoDBResourceManager.java    |  11 +-
 .../it/mongodb/MongoDBResourceManagerTest.java     |   5 +
 release/build.gradle.kts                           |   2 +-
 release/src/main/groovy/TestScripts.groovy         | 156 +++-
 .../main/groovy/mobilegaming-java-dataflow.groovy  |  75 +-
 .../main/groovy/mobilegaming-java-direct.groovy    |  73 +-
 .../main/groovy/quickstart-java-flinklocal.groovy  |   4 +-
 .../src/main/groovy/quickstart-java-spark.groovy   |  20 +-
 runners/google-cloud-dataflow-java/build.gradle    |   2 +
 .../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 |   1 +
 .../dataflow/worker/StreamingDataflowWorker.java   |  49 +-
 .../worker/StreamingModeExecutionContext.java      |  67 +-
 .../logging/DataflowWorkerLoggingHandler.java      |  30 +-
 .../logging/DataflowWorkerLoggingInitializer.java  |   4 +
 .../dataflow/worker/streaming/ActiveWorkState.java |  46 +-
 .../streaming/BoundedQueueExecutorWorkHandle.java  |   8 +-
 .../worker/streaming/ComputationState.java         |   5 +-
 .../worker/streaming/ComputationWorkExecutor.java  |   6 +-
 .../dataflow/worker/streaming/ExecutableWork.java  |   8 +
 ...cutorWorkHandle.java => FailedWorkHandler.java} |  13 +-
 ...java => MultiKeyCommitValidationException.java} |  13 +-
 .../runners/dataflow/worker/streaming/Work.java    |  20 +
 .../dataflow/worker/util/BoundedQueueExecutor.java |  37 +-
 .../dataflow/worker/util/KeyGroupWorkQueue.java    |  12 +-
 .../client/grpc/stubs/FailoverChannel.java         |  66 +-
 .../work/processing/StreamingWorkScheduler.java    |  51 +-
 .../processing/failures/WorkFailureProcessor.java  |  42 +-
 .../windmill/work/refresh/ActiveWorkRefresher.java |  21 +-
 .../worker/StreamingDataflowWorkerTest.java        | 787 +++++++++++++++++++--
 .../worker/StreamingModeExecutionContextTest.java  | 141 +++-
 .../dataflow/worker/WorkerCustomSourcesTest.java   |   9 +-
 .../logging/DataflowWorkerLoggingHandlerTest.java  | 103 ++-
 .../worker/streaming/ActiveWorkStateTest.java      |  42 +-
 .../worker/util/BoundedQueueExecutorTest.java      | 129 +++-
 .../worker/util/KeyGroupWorkQueueTest.java         |  21 +
 .../client/grpc/stubs/FailoverChannelTest.java     | 107 ++-
 .../failures/WorkFailureProcessorTest.java         |  72 +-
 .../work/refresh/ActiveWorkRefresherTest.java      |  69 --
 .../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}   |  15 +-
 .../kafka-streams/proto/build.gradle               |  31 +-
 .../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}  |  13 +-
 .../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}   |  13 +-
 .../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/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/ci/pr-bot/processNewPrs.ts                 |  20 +
 scripts/ci/pr-bot/shared/githubUtils.ts            |  21 +
 sdks/go.mod                                        |  92 +--
 sdks/go.sum                                        | 184 ++---
 sdks/go/container/boot.go                          |  40 +-
 sdks/go/container/boot_test.go                     |  17 +-
 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 --
 .../core/runtime/xlangx/expansionx/download.go     |  17 +-
 .../runtime/xlangx/expansionx/download_test.go     |  25 +
 .../prism/internal/engine/elementmanager.go        |  91 ++-
 .../engine/elementmanager_continuation_test.go     | 299 ++++++++
 sdks/java/container/boot.go                        |   4 +-
 .../apache/beam/sdk/options/SdkHarnessOptions.java |   7 +
 .../org/apache/beam/sdk/schemas/SchemaUtils.java   |   7 +
 .../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 +-
 .../apache/beam/sdk/schemas/SchemaUtilsTest.java   |  43 ++
 .../sdk/transforms/GroupByEncryptedKeyTest.java    |   2 +-
 .../java/org/apache/beam/sdk/util/SecretTest.java  | 179 ++++-
 .../util/construction/CoderTranslationTest.java    |  38 +-
 .../extensions/avro/AvroGenericCoderRegistrar.java |  18 +
 .../schemaio-expansion-service/build.gradle        |   6 +
 .../beam/fn/harness/state/StateBackedIterable.java |  22 +
 .../KinesisReadSchemaTransformProvider.java        | 315 +++++++++
 .../KinesisWriteSchemaTransformProvider.java       | 280 ++++++++
 .../KinesisSchemaTransformProviderTest.java        | 221 ++++++
 .../beam/sdk/io/cassandra/ConnectionManager.java   |  17 +-
 .../beam/sdk/io/cassandra/CassandraIOTest.java     |  30 +
 .../beam/sdk/io/delta/CreateReadTasksDoFn.java     |  23 +-
 .../beam/sdk/io/delta/DeltaCDCSourceDoFn.java      |  24 +-
 .../java/org/apache/beam/sdk/io/delta/DeltaIO.java |  30 +-
 .../io/delta/DeltaReadSchemaTransformProvider.java |   6 +-
 .../org/apache/beam/sdk/io/delta/DeltaIOIT.java    | 141 ++--
 .../org/apache/beam/sdk/io/delta/DeltaIOTest.java  | 200 ++++++
 .../DeltaReadSchemaTransformProviderTest.java      |  51 ++
 .../beam/sdk/io/delta/DeltaWriteTestUtils.java     |  30 +
 sdks/java/io/google-cloud-platform/build.gradle    |  36 +-
 .../beam/sdk/io/gcp/bigquery/BigQueryHelpers.java  | 100 ++-
 .../beam/sdk/io/gcp/bigquery/BigQueryIO.java       |  13 +
 .../io/gcp/bigquery/BigQueryStorageSourceBase.java |   7 +-
 .../gcp/bigquery/BigQueryStorageTableSource.java   |   3 +-
 .../sdk/io/gcp/bigquery/BigQueryTableSource.java   |   5 +
 .../sdk/io/gcp/bigtable/BigtableServiceImpl.java   |   8 +
 .../sdk/io/gcp/bigquery/BigQueryHelpersTest.java   | 348 ++++++++-
 .../bigquery/BigQueryIOIcebergManagedTableIT.java  | 408 +++++++++++
 .../io/gcp/bigquery/BigQueryIOStorageReadTest.java |  63 ++
 ....java => StorageApiSinkSchemaUpdateITBase.java} |  61 +-
 ...torageApiSinkSchemaUpdateWithInputSchemaIT.java |  50 ++
 ...ageApiSinkSchemaUpdateWithoutInputSchemaIT.java |  51 ++
 .../io/gcp/bigtable/BigtableServiceImplTest.java   |  81 +++
 sdks/java/io/iceberg/build.gradle                  |  11 +
 .../org/apache/beam/sdk/io/iceberg/AddFiles.java   | 168 +++--
 .../org/apache/beam/sdk/io/iceberg/IcebergIO.java  |  10 +
 .../beam/sdk/io/iceberg/NameMappingUtils.java      | 215 ++++++
 .../apache/beam/sdk/io/iceberg/PartitionUtils.java |  10 +-
 .../beam/sdk/io/iceberg/RecordWriterManager.java   |   9 +-
 .../beam/sdk/io/iceberg/SerializableDataFile.java  |  88 ++-
 .../beam/sdk/io/iceberg/SerializableTableSpec.java | 382 ++++++++++
 .../apache/beam/sdk/io/iceberg/SideInputTable.java | 368 ++++++++++
 .../io/iceberg/WritePartitionedRowsToFiles.java    |   3 +-
 .../beam/sdk/io/iceberg/cdc/CdcReadUtils.java      |   2 +-
 .../io/iceberg/cdc/SerializableChangelogTask.java  |   9 +-
 .../org/apache/beam/sdk/io/iceberg/AddFilesIT.java |  74 +-
 .../apache/beam/sdk/io/iceberg/AddFilesTest.java   | 255 ++++++-
 .../iceberg/BigQueryManagedTableCrossEngineIT.java | 183 +++++
 .../IcebergWriteSchemaTransformProviderTest.java   |  63 ++
 .../beam/sdk/io/iceberg/NameMappingUtilsTest.java  | 425 +++++++++++
 .../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        |  10 +
 .../io/iceberg/catalog/IcebergCatalogBaseIT.java   | 259 +++++++
 .../sdk/io/iceberg/catalog/RESTCatalogBLMSIT.java  |  22 +-
 .../beam/sdk/io/iceberg/cdc/CdcReadUtilsTest.java  |   2 +-
 .../sdk/io/iceberg/cdc/ChangelogScannerTest.java   |   2 +-
 .../cdc/IncrementalChangelogSourceTest.java        |  86 +++
 .../sdk/io/iceberg/cdc/LocalResolveDoFnTest.java   |   2 +-
 .../sdk/io/iceberg/cdc/ReadFromChangelogsTest.java |   2 +-
 .../iceberg/cdc/SerializableChangelogTaskTest.java |   2 +-
 .../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 ++
 .../java/org/apache/beam/sdk/io/kafka/KafkaIO.java |   7 +-
 .../kafka/KafkaWriteSchemaTransformProvider.java   |  41 +-
 ...KafkaIOReadImplementationCompatibilityTest.java |  22 +
 .../KafkaWriteSchemaTransformProviderTest.java     |  27 +-
 .../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 +-
 .../io/external/xlang_jdbcio_it_test.py            | 113 +++
 .../apache_beam/io/external/xlang_jmsio_it_test.py |  13 +-
 sdks/python/apache_beam/io/fileio.py               | 305 +++++---
 sdks/python/apache_beam/io/fileio_test.py          | 345 +++++++++
 sdks/python/apache_beam/io/gcp/bigquery_tools.py   |   8 +
 sdks/python/apache_beam/io/gcp/bigtableio.py       |   8 +-
 sdks/python/apache_beam/io/gcp/bigtableio_test.py  |  55 ++
 .../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/jdbc.py                 |  48 +-
 sdks/python/apache_beam/io/textio_test.py          |   7 +-
 sdks/python/apache_beam/io/watch.py                | 265 +++----
 sdks/python/apache_beam/io/watch_test.py           | 222 +++---
 sdks/python/apache_beam/ml/inference/base.py       |  28 +
 .../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 |  28 +
 .../apache_beam/runners/dataflow/internal/names.py |   2 +-
 .../runners/direct/transform_evaluator.py          |   5 +
 .../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/testing/pubsub_test_context.py     | 179 +++++
 .../testing/pubsub_test_context_test.py            | 149 ++++
 sdks/python/apache_beam/transforms/combiners.py    |  58 ++
 .../apache_beam/transforms/combiners_test.py       |  54 ++
 .../transforms/managed_iceberg_it_test.py          |   3 +-
 .../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       |  10 +-
 sdks/python/apache_beam/typehints/schemas_test.py  |  10 +
 sdks/python/apache_beam/utils/secret.py            | 464 ++++++++++++
 sdks/python/apache_beam/utils/secret_test.py       | 454 ++++++++++++
 sdks/python/apache_beam/utils/timestamp.py         |   5 +-
 sdks/python/apache_beam/utils/timestamp_test.py    |  17 +
 .../yaml/extended_tests/databases/iceberg.yaml     |   3 +
 .../databases/jdbc_secret_manager.yaml             |  59 ++
 .../extended_tests/e2e/delta_lake_to_iceberg.yaml  |  71 ++
 .../yaml/extended_tests/messaging/kinesis.yaml     |  83 +++
 sdks/python/apache_beam/yaml/integration_tests.py  | 269 +++++++
 sdks/python/apache_beam/yaml/standard_io.yaml      | 105 +++
 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            | 129 ++++
 sdks/python/apache_beam/yaml/yaml_io_test.py       | 176 +++++
 sdks/python/apache_beam/yaml/yaml_provider.py      |   6 +-
 sdks/python/build.gradle                           |  12 +-
 sdks/python/container/Dockerfile                   |  14 +-
 sdks/python/container/boot.go                      |  57 +-
 .../license_scripts/upgrade_bundled_pip.py         |  15 +-
 .../container/ml/py310/base_image_requirements.txt | 155 ++--
 .../container/ml/py310/gpu_image_requirements.txt  | 211 +++---
 .../container/ml/py311/base_image_requirements.txt | 158 ++---
 .../container/ml/py311/gpu_image_requirements.txt  | 214 +++---
 .../container/ml/py312/base_image_requirements.txt | 158 ++---
 .../container/ml/py312/gpu_image_requirements.txt  | 212 +++---
 .../container/ml/py313/base_image_requirements.txt | 158 ++---
 sdks/python/container/profiler.go                  |  69 +-
 sdks/python/container/profiler_test.go             |  19 +-
 .../container/py310/base_image_requirements.txt    | 143 ++--
 .../container/py311/base_image_requirements.txt    | 146 ++--
 .../container/py312/base_image_requirements.txt    | 146 ++--
 .../container/py313/base_image_requirements.txt    | 146 ++--
 .../container/py314/base_image_requirements.txt    | 146 ++--
 sdks/python/pyproject.toml                         |   3 +-
 sdks/python/setup.py                               |   6 +-
 sdks/python/test-suites/portable/common.gradle     |  39 +
 sdks/python/tox.ini                                |   5 +
 sdks/standard_external_transforms.yaml             | 134 +++-
 sdks/typescript/container/boot.go                  |   4 +-
 settings.gradle.kts                                |  12 +
 .../www/site/assets/scss/_capability-matrix.scss   | 123 +---
 .../www/site/assets/scss/capability-matrix.scss    | 121 +---
 .../en/documentation/io/developing-io-python.md    | 239 ++++++-
 .../site/content/en/documentation/io/managed-io.md |   4 +-
 .../en/documentation/runners/kafkastreams.md       | 248 +++++++
 .../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 +-
 382 files changed, 31459 insertions(+), 3695 deletions(-)
 copy .github/trigger_files/{beam_PostCommit_Java_Nexmark_Spark.json => 
beam_PostCommit_Java_ValidatesRunner_Spark4.json} (100%)
 copy .github/trigger_files/{IO_Iceberg_Integration_Tests_Dataflow.json => 
beam_PreCommit_Java_Kafka_Streams_Runner.json} (100%)
 delete mode 100644 .github/workflows/beam_Infrastructure_AuditUnmanagedKeys.yml
 copy .github/workflows/{beam_PostCommit_Yaml_Xlang_Direct.yml => 
beam_PreCommit_Java_Kafka_Streams_Runner.yml} (54%)
 create mode 100644 
examples/java/src/test/java/org/apache/beam/examples/subprocess/utils/FileUtilsTest.java
 copy 
runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/{BoundedQueueExecutorWorkHandle.java
 => FailedWorkHandler.java} (70%)
 copy 
runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/{BoundedQueueExecutorWorkHandle.java
 => MultiKeyCommitValidationException.java} (71%)
 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 
runners/{google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/BoundedQueueExecutorWorkHandle.java
 => 
kafka-streams/measurement/src/main/java/org/apache/beam/runners/kafka/streams/measurement/package-info.java}
 (66%)
 copy 
sdks/java/core/src/main/java/org/apache/beam/sdk/util/construction/CoderTranslatorRegistrar.java
 => runners/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 
runners/{google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/BoundedQueueExecutorWorkHandle.java
 => 
kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/package-info.java}
 (65%)
 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 
runners/{google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/BoundedQueueExecutorWorkHandle.java
 => 
kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/translation/package-info.java}
 (65%)
 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/runners/prism/internal/engine/elementmanager_continuation_test.go
 create mode 100644 
sdks/java/core/src/main/java/org/apache/beam/sdk/util/RawSecret.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
 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/NameMappingUtils.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
 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/NameMappingUtilsTest.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/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
 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/jdbc_secret_manager.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 
website/www/site/content/en/documentation/runners/kafkastreams.md

Reply via email to