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