This is an automated email from the ASF dual-hosted git repository.
Amar3tto pushed a change to branch updates_managed_io_docs_2.76.0_rc4
in repository https://gitbox.apache.org/repos/asf/beam.git
from 5f8c4cfa748 Update managed-io.md for release 2.76.0-RC4.
add e626f54690b Implement Vertex AI Model Monitoring v2 (#39738)
add 3b9e647dec1 [Prism] Schedule consumers of a self checkpointing source
(#39572)
add 88b3ee7b488 JmsIO yaml (#39818)
add 4e5ae91554e Support lakehouse PCNT format in BigQueryIO storage read.
(#39597)
add 6434f7411e7 [Docs] Document UnboundedSource in the Python I/O
connector guide (#39529)
add 61ed38fd5fd Support Secret Manager in JdbcIO for Java, Python and YAML
(#39834)
add 7de4d3ab73b Fix KafkaIO wrtie SchemaTransform parallelism (#39844)
add 015e823a1a9 [GSoC-273] Feat: Integrate TestPubsubContext to prevent
Pub/Sub resource leaks and expand stale cleaner scope (#39826)
add c1017953ac3 Adds support for reading at a given Delta Lake version or
timestamp (#39758)
add 2577a6cff4a Install Cloud Spanner emulator component in Go PreCommit
CI (#39850)
add 2e445307153 Add Kafka Streams runner skeleton module and portable
entry points
add 61c891a69ca Address review notes on KafkaStreamsPipelineResult and
state dir
add cef6544e792 Address review feedback on Kafka Streams Runner skeleton
add 0d445f79e0c Drop setRunner(null) suppression; make applicationId
required
add d4c7dae98b8 Catch Exception in KafkaStreamsRunner.run() to avoid
job-server leak
add 64ad482f9de Merge pull request #38534: [GSoC 2026] Kafka Streams
runner skeleton module + portable entry points
add faefc95e9f0 [GSoC 2026] Kafka Streams runner — translation framework +
Impulse translator (#38689)
add 0f875a07162 Add temporary feature-branch CI for Kafka Streams runner
(#38725)
add ff554f7438d [GSoC 2026] Kafka Streams runner — ExecutableStage
(stateless ParDo) translator (#38764)
add 4c02a7789a2 [GSoC 2026] Kafka Streams runner — Redistribute translator
+ ExecutableStage type-agnostic edge (#38843)
add 4636b1f333e #38957: Add in-memory WatermarkManager core
(per-source-partition tracking
add 1058e94027e [GSoC 2026] Kafka Streams runner #38987: Wire
WatermarkManager into ExecutableStageProcessor
add 3e963a68343 [GSoC 2026] Kafka Streams runner #39051: Add
KStreamsPayload Serde for crossing topic boundaries
add 5e65d4772fb [GSoC 2026] Kafka Streams runner #39141: Add GroupByKey
(GlobalWindow, fire at watermark)
add 4d5847f6bfa [GSoC 2026] Kafka Streams runner #39211: Add
KafkaStreamsTestRunner test harness
add 27c4522ae5e [GSoC 2026] Kafka Streams runner #39249: Support Create
add ff323ebd6db [GSoC 2026] Kafka Streams runner #39273: Kafka Streams
runner: Flatten support
add 75adf49ee8d [GSoC 2026] Kafka Streams runner: surface SDK-harness
metrics as MetricResults (#39341)
add 87911f00973 [GSoC 2026] Kafka Streams runner:
TestPipeline-dispatchable test runner (PAssert works) (#39362)
add d9cea589e5b [GSoC 2026] Kafka Streams runner: validatesRunner task;
Create and Flatten suites green (#39380)
add fca14356614 [GSoC 2026] Kafka Streams runner: multi-output executable
stages (#39410)
add b615ae85aed [GSoC 2026] Kafka Streams runner: enable ParDoTest in the
ValidatesRunner suite (#39451)
add 47614a943a7 [GSoC 2026] Kafka Streams runner: windowed GroupByKey via
ReduceFnRunner (#39494)
add 27198768f89 [GSoC 2026] Kafka Streams runner: run on a real broker,
correctly across partitions (#39546)
add cf1f10bdf82 [GSoC 2026] Kafka Streams runner: bound a bundle by
element count (#39578)
add fc36301e6c2 [GSoC 2026] Kafka Streams runner: CombineTest coverage and
two review follow-ups (#39610)
add 5861f31e8ac [GSoC 2026] Kafka Streams runner: read unbounded sources
(#39611)
add cb30afd092e [GSoC 2026] Kafka Streams runner: user documentation,
marked experimental (#39627)
add e051e06bf83 [GSoC 2026] Kafka Streams runner: Python wrapper that
starts its own job server (#39680)
add 4ff618047c3 [GSoC 2026] Kafka Streams runner: terminate a bounded
pipeline when it is drained (#39700)
add 5f1658f8b5e [GSoC 2026] Kafka Streams runner: portable ValidatesRunner
suite for Python (#39736)
add 511a40e4f20 [GSoC 2026] Kafka Streams runner: separate the source's
poll size from the bundle size, and expose the session timeout (#39748)
add 65a2e400c73 [GSoC 2026] Kafka Streams runner: bound a source poll in
time, not only in elements (#39761)
add 104dc272e8d [GSoC 2026] Kafka Streams runner: ask for primitive reads
in the Java wrapper (#39766)
add 49d459a243e [GSoC 2026] Kafka Streams runner: put the runner behind an
opt-in build flag (#39762)
add 10ff557d616 [GSoC 2026] Kafka Streams runner: an application for
measuring instances coming and going (#39752)
add 5b3702484f4 [GSoC 2026] Kafka Streams runner: license header and
Python formatting for master CI
add b12bfc159ba [GSoC 2026] Kafka Streams runner: shorten the explanation
comments
add 8d5f6511a6b Merge pull request #39781: [GSoC 2026] Kafka Streams
runner: shorten comments, and fix the license header and Python formatting
add e2681085c96 Build Kafka Streams runner during javaPreCommit (#18479)
add 8b3b08dc799 Merge pull request #39784: Build Kafka Streams runner
during javaPreCommit (#18479)
add b67cf234ead [GSoC 2026] Kafka Streams runner: update the CHANGES.md
entry
add 52dadbac433 Merge pull request #39786: [GSoC 2026] Kafka Streams
runner: update the CHANGES.md entry
add 89617f59b45 Merge branch 'master' of https://github.com/apache/beam
into feat/18479-kafka-streams-runner-skeleton
add 6b75a44f0da Removed feature branch build
add 52d6eeab74e [GSoC 2026] Kafka Streams runner: say what the Python side
is, and is not
add 473a5dd649c Merge pull request #39847: [GSoC 2026] Kafka Streams
runner: say what the Python side is, and is not
add e10e8c969b0 Merge pull request #39785: Merge Kafka Streams Runner
skeleton (#18479)
add 7056e67a285 [GSoC 2026] Publish the Kafka Streams runner to nightly
snapshots
add ec5ebf422f0 Merge pull request #39854: [GSoC 2026] Publish the Kafka
Streams runner to nightly snapshots
add 466cf6d265b Revert "[GSoC-273] Feat: Integrate TestPubsubContext to
prevent Pub/Sub resou…" (#39853)
add 59e916e5a11 AddFiles: regenerate name mapping, read footers once,
error helper (#39836)
add 1ef9cfd66fb Bump actions/checkout from 6 to 7 (#39859)
add 08522c7ca94 Bump github/codeql-action from 4.37.7 to 4.37.8 (#39864)
add f15bd00f526 use LocalStack TLS hostname in Kinesis YAML IT (#39868)
add f95421bbd90 Revert "Install Cloud Spanner emulator component in Go
PreCommit CI (#39850)" (#39877)
add da360800706 Bump github.com/aws/aws-sdk-go-v2/feature/s3/manager in
/sdks (#39860)
add 3e30bac4cf6 Bump google.golang.org/grpc from 1.83.0 to 1.83.1 in /sdks
(#39863)
add ffbe1d58390 Bump cloud.google.com/go/pubsub from 1.51.0 to 1.51.1 in
/sdks (#39862)
add ea7d39cb9ab Bump cloud.google.com/go/storage from 1.64.0 to 1.65.0 in
/sdks (#39861)
add af99c7ce5d4 [Bigtable] Fix Bigtable segment truncation when open end
key is startKey + null byte (#39842) (#39843)
add 66ecc31a22c Adds a new CoderTranslator for Java SchemaCoders. (#39594)
add 4bd8a86fd83 Bump github.com/aws/smithy-go from 1.27.8 to 1.27.9 in
/sdks (#39880)
add 13875fc6bd3 Bump cloud.google.com/go/bigquery from 1.80.0 to 1.81.0 in
/sdks (#39881)
add 73ea0b72c12 [Prism] Honor the resume delay of self-checkpointing SDF
residuals (#39849)
add c267ce6520f Bump github.com/fsouza/fake-gcs-server from 1.55.1 to
1.56.0 in /sdks (#39890)
add 571184a509e Add Vertex Model Monitoring to CHANGES.md (#39882)
add 87e31bd8a31 Update build.gradle.kts (#39892)
add 167b8649fc9 [Java] Bound the Watch deduplication state with a
timestamp cursor (#39746)
add b2013e09f91 Fix ZeroDivisionError at initial invocation of
monitoring_info when there is no work yet (#39885)
add c494cc647fb Feat: expanding the context of the second and third
security layers of testPubSubContext so that, in the event of failure,
subscriptions retain a 24-hour grace period before being deleted (#39856)
add 1fca1c2d371 prevent cloudML TFT extras from replacing the SDK under
test (#39884)
add 69f7fd48a0c Add Snowflake YAML write transform
add 85bad221c71 Add Snowflake Read and Streaming Write
add 90c93d7f541 Fix bytes
add d131e0e74b7 Fix imports
add f4e62fdafbf Refactoring
add 15efa5d2084 Add extended test
add b275c84f8e2 Add license
add 27c4e6395cb Move depends
add a93fdcdda75 Fix jms yaml
add da201e642a2 Fix import Nullable
add 497870ade99 Change to provided
add 8cfbd8e01c2 Merge pull request #39742 from apache/snowflakeio-yaml
add ad39be0f58b Clarify service account key issue reports (#39894)
add 4d2c83d9b66 Bump github.com/fsouza/fake-gcs-server from 1.56.0 to
1.56.1 in /sdks (#39902)
add 351f4c7eac9 Bump actions/setup-java from 5 to 6 (#39903)
add cf27972eb10 Update CHANGES.md with new known issues (#39898)
add 01559cc9127 Added schema_update_options to Python BigQuery writes
(#39078)
add 34a00c7781f AddFiles: extract bounded async task plumbing and Parquet
footer reads (#39896)
add 8d5a5300fa3 Remove nullness suppression in trigger implementation
add 73e5ecc2ea8 Merge pull request #39807: Remove nullness suppression in
trigger implementation
add 493a18f8217 Bump docker/setup-buildx-action from 4.2.0 to 4.3.0
(#39829)
add d14b6483d5e Pin googleapis-common-protos<1.70 for py310-tensorflow-212
(#39904)
add d7557af6a55 [Spark 4] Add streaming dispatch seam and lifecycle hooks
(#39906)
add 470454db370 Bump cloud.google.com/go/spanner from 1.94.0 to 1.95.0 in
/sdks (#39913)
add a057f722cda Try fix miniCluster jar (#39912)
add 13044d31d78 Fix MongoDbIO read splitting to preserve non-ObjectId _id
types (#39901)
add 8288d2f15d0 Update commons-pool2 dependency to 2.13.1 in JDBC IO
(#39914)
add ade938f14f3 Update Beam website roadmap pages (#39899)
add 85f7a48781c datastore: Same key in mutation batch
add 2c2036fe168 Merge branch 'master' into repeated-keys-datastore-writes
add 6dbbdda8c6c datastore: assert return value of adding to set
add 52d74b914e3 Merge pull request #37751: datastore: Same key in mutation
batch
add 25cb8b76a96 fix(beam): guard jinja_variable_flags collision with
pipeline options (#39888)
add bd9620f61c0 [Spark] Fire processing-time timers in timestamp order
(#39825)
add aa35a98a681 Deflake JmsIO authentication tests (#39920)
add 4748eb678a9 ClickHouseIO: Add Decimal(P, S) support for fixed-point
columns (#39846)
add 6dcf56a4e96 Release HBaseIO resources even when an earlier close()
throws (#39712)
add cbe82e50bc3 Tag the TFRecord write error output with the schema it
actually emits (#39759)
add 6a1efc90e59 SolaceIO - fix for data loss during scaling/rebalancing
(#36991) (#38603)
add bb068025e2d Fix flaky Playground local cache tests (#39918)
add 7bbb7f7c0e3 Bump github/codeql-action from 4.37.8 to 4.37.9 (#39928)
add 72240290c7a Bump github.com/aws/aws-sdk-go-v2/config from 1.32.38 to
1.33.1 in /sdks (#39925)
add 19095e367ed Bump google.golang.org/api from 0.293.0 to 0.294.0 in
/sdks (#39926)
add 7fa7ccc0fde Bump github.com/nats-io/nats-server/v2 from 2.14.5 to
2.14.6 in /sdks (#39930)
add 2c11a5fbc04 Fail fast in artifact staging when storing an artifact
fails (#39367)
add 7ada585d853 [Java] Preserve CoderTranslatorRegistrar binary
compatibility (#39919)
add 1955788db62 Persist credentials for finalize_release.yml
add a3e7524aec7 Merge remote-tracking branch 'refs/remotes/origin/master'
into updates_managed_io_docs_2.76.0_rc4
No new revisions were added by this update.
Summary of changes:
.../IO_Iceberg_Integration_Tests.json | 2 +-
.../beam_CloudML_Benchmarks_Dataflow.json | 2 +-
.../beam_PostCommit_Java_Delta_IO_Dataflow.json | 2 +-
.../beam_PostCommit_Java_PVR_Spark3_Streaming.json | 2 +-
...beam_PostCommit_Java_ValidatesRunner_Spark.json | 2 +-
.github/trigger_files/beam_PostCommit_Python.json | 2 +-
.../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 | 1 +
.github/workflows/beam_PostCommit_Go.yml | 2 +-
.../workflows/beam_PostCommit_Go_Dataflow_ARM.yml | 2 +-
.../beam_PostCommit_Java_Examples_Dataflow_ARM.yml | 2 +-
.../beam_PostCommit_XVR_GoUsingJava_Dataflow.yml | 2 +-
.../workflows/beam_PreCommit_CommunityMetrics.yml | 2 +-
.github/workflows/beam_PreCommit_Java.yml | 1 +
...> beam_PreCommit_Java_Kafka_Streams_Runner.yml} | 75 ++-
.github/workflows/beam_PreCommit_PythonDocker.yml | 2 +-
.../workflows/beam_Publish_Beam_SDK_Snapshots.yml | 2 +-
.../workflows/beam_Publish_Python_VLLM_Image.yml | 2 +-
...beam_Python_ValidatesContainer_Dataflow_ARM.yml | 2 +-
.github/workflows/beam_Release_NightlySnapshot.yml | 1 +
.github/workflows/build_release_candidate.yml | 12 +-
.github/workflows/build_runner_image.yml | 2 +-
.github/workflows/code_completion_plugin_tests.yml | 2 +-
.github/workflows/codeql.yml | 4 +-
.github/workflows/finalize_release.yml | 4 +-
.../republish_released_docker_containers.yml | 4 +-
.github/workflows/typescript_tests.yml | 4 +-
.test-infra/tools/stale_cleaner.py | 7 +
CHANGES.md | 18 +-
build.gradle.kts | 18 +-
infra/enforcement/README.md | 4 +-
infra/enforcement/account_keys.py | 9 +-
infra/enforcement/sending.py | 10 +-
infra/enforcement/test_sending.py | 6 +-
.../internal/cache/local/local_cache_test.go | 43 +-
.../core/triggers/AfterAllStateMachine.java | 5 +-
.../AfterDelayFromFirstElementStateMachine.java | 5 +-
.../core/triggers/AfterEachStateMachine.java | 26 +-
.../core/triggers/AfterFirstStateMachine.java | 5 +-
.../core/triggers/AfterWatermarkStateMachine.java | 12 +-
.../core/triggers/DefaultTriggerStateMachine.java | 3 -
.../triggers/ExecutableTriggerStateMachine.java | 21 +-
.../core/triggers/OrFinallyStateMachine.java | 5 +-
.../core/triggers/RepeatedlyStateMachine.java | 5 +-
.../runners/core/triggers/TriggerStateMachine.java | 17 +-
.../TriggerStateMachineContextFactory.java | 23 +-
.../core/triggers/AfterEachStateMachineTest.java | 54 ++
runners/flink/job-server/flink_job_server.gradle | 8 +-
runners/google-cloud-dataflow-java/build.gradle | 2 +
.../dataflow/DataflowPipelineTranslator.java | 2 +-
.../beam/runners/dataflow/DataflowRunner.java | 54 +-
.../dataflow/DataflowPipelineTranslatorTest.java | 78 ++-
.../artifact/ArtifactStagingService.java | 34 +-
.../artifact/ArtifactStagingServiceTest.java | 68 +++
.../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 | 27 +-
.../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 | 21 +-
.../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 | 21 +-
.../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/portability/test_flink_uber_jar.sh | 12 +-
runners/prism/java/build.gradle | 6 -
runners/spark/4/build.gradle | 8 +-
.../spark/stateful/SparkTimerInternals.java | 52 +-
.../SparkStructuredStreamingPipelineOptions.java | 34 ++
.../SparkStructuredStreamingPipelineResult.java | 16 +-
.../SparkStructuredStreamingRunner.java | 26 +-
.../translation/EvaluationContext.java | 18 +-
.../translation/PipelineTranslator.java | 22 +-
.../translation/PipelineTranslatorFactory.java | 41 ++
.../translation/SparkSessionFactory.java | 6 +
.../apache/beam/runners/spark/util/TimerUtils.java | 7 +-
.../spark/stateful/SparkTimerInternalsTest.java | 106 ++++
.../beam/runners/spark/util/TimerUtilsTest.java | 20 +
sdks/go.mod | 75 ++-
sdks/go.sum | 173 +++---
.../prism/internal/engine/elementmanager.go | 193 +++++-
.../engine/elementmanager_continuation_test.go | 417 +++++++++++++
.../io/TFRecordWriteSchemaTransformProvider.java | 3 +-
.../java/org/apache/beam/sdk/transforms/Watch.java | 239 +++++++-
.../sdk/util/construction/CoderTranslation.java | 62 +-
.../construction/CoderTranslatorRegistrar.java | 26 +
.../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 +-
.../io/TFRecordSchemaTransformProviderTest.java | 29 +
.../org/apache/beam/sdk/transforms/WatchTest.java | 348 ++++++++++-
.../util/construction/CoderTranslationTest.java | 78 ++-
.../extensions/avro/AvroGenericCoderRegistrar.java | 18 +
.../schemaio-expansion-service/build.gradle | 6 +
.../beam/fn/harness/state/StateBackedIterable.java | 22 +
.../beam/sdk/io/clickhouse/ClickHouseIO.java | 7 +
.../beam/sdk/io/clickhouse/ClickHouseWriter.java | 47 ++
.../apache/beam/sdk/io/clickhouse/TableSchema.java | 61 +-
.../clickhouse/src/main/javacc/ColumnTypeParser.jj | 43 ++
.../beam/sdk/io/clickhouse/ClickHouseIOIT.java | 182 ++++++
.../sdk/io/clickhouse/ClickHouseWriterTest.java | 116 ++++
.../beam/sdk/io/clickhouse/TableSchemaTest.java | 152 +++++
.../beam/sdk/io/delta/CreateReadTasksDoFn.java | 23 +-
.../java/org/apache/beam/sdk/io/delta/DeltaIO.java | 26 +-
.../io/delta/DeltaReadSchemaTransformProvider.java | 6 +-
.../org/apache/beam/sdk/io/delta/DeltaIOIT.java | 141 +++--
.../org/apache/beam/sdk/io/delta/DeltaIOTest.java | 102 ++++
.../DeltaReadSchemaTransformProviderTest.java | 51 ++
.../beam/sdk/io/delta/DeltaWriteTestUtils.java | 30 +
sdks/java/io/google-cloud-platform/build.gradle | 33 ++
.../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 +
.../beam/sdk/io/gcp/datastore/DatastoreV1.java | 6 +-
.../sdk/io/gcp/bigquery/BigQueryHelpersTest.java | 348 ++++++++++-
.../bigquery/BigQueryIOIcebergManagedTableIT.java | 408 +++++++++++++
.../io/gcp/bigquery/BigQueryIOStorageReadTest.java | 63 ++
.../io/gcp/bigtable/BigtableServiceImplTest.java | 81 +++
.../beam/sdk/io/gcp/datastore/DatastoreV1Test.java | 46 ++
sdks/java/io/hbase/build.gradle | 1 +
.../java/org/apache/beam/sdk/io/hbase/HBaseIO.java | 91 ++-
.../apache/beam/sdk/io/hbase/HBaseIOCloseTest.java | 160 +++++
sdks/java/io/iceberg/build.gradle | 11 +
.../org/apache/beam/sdk/io/iceberg/AddFiles.java | 255 ++++----
.../beam/sdk/io/iceberg/BoundedAsyncTasks.java | 113 ++++
.../beam/sdk/io/iceberg/NameMappingUtils.java | 215 +++++++
.../apache/beam/sdk/io/iceberg/ParquetFooters.java | 51 ++
.../org/apache/beam/sdk/io/iceberg/AddFilesIT.java | 71 ++-
.../apache/beam/sdk/io/iceberg/AddFilesTest.java | 255 +++++++-
.../iceberg/BigQueryManagedTableCrossEngineIT.java | 183 ++++++
.../beam/sdk/io/iceberg/BoundedAsyncTasksTest.java | 220 +++++++
.../beam/sdk/io/iceberg/NameMappingUtilsTest.java | 425 +++++++++++++
.../beam/sdk/io/iceberg/ParquetFootersTest.java | 113 ++++
.../catalog/BigQueryMetastoreCatalogIT.java | 10 +
.../io/iceberg/catalog/IcebergCatalogBaseIT.java | 259 ++++++++
.../sdk/io/iceberg/catalog/RESTCatalogBLMSIT.java | 21 +-
sdks/java/io/jdbc/build.gradle | 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/jms/JmsIOTest.java | 5 +-
.../kafka/KafkaWriteSchemaTransformProvider.java | 41 +-
.../KafkaWriteSchemaTransformProviderTest.java | 27 +-
.../org/apache/beam/sdk/io/mongodb/MongoDbIO.java | 69 +--
.../apache/beam/sdk/io/mongodb/MongoDbIOTest.java | 52 +-
sdks/java/io/snowflake/build.gradle | 2 +
.../apache/beam/sdk/io/snowflake/SnowflakeIO.java | 3 +-
.../SnowflakeReadSchemaTransformProvider.java | 289 +++++++++
.../snowflake/SnowflakeSchemaTransformUtils.java | 342 +++++++++++
.../io/snowflake/SnowflakeWriteConfiguration.java | 227 +++++++
.../SnowflakeWriteSchemaTransformProvider.java | 179 ++++++
.../SnowflakeReadSchemaTransformProviderTest.java | 311 ++++++++++
.../SnowflakeWriteSchemaTransformProviderTest.java | 494 ++++++++++++++++
.../org/apache/beam/sdk/io/solace/SolaceIO.java | 49 +-
.../sdk/io/solace/read/ActiveReadersRegistry.java | 46 ++
.../sdk/io/solace/read/SolaceCheckpointMark.java | 50 +-
.../sdk/io/solace/read/UnboundedSolaceReader.java | 175 +++++-
.../sdk/io/solace/read/UnboundedSolaceSource.java | 24 +-
.../beam/sdk/io/solace/SolaceIOReadTest.java | 103 +++-
.../beam/sdk/io/solace/data/SolaceDataUtils.java | 30 +-
.../io/solace/read/UnboundedSolaceReaderTest.java | 215 +++++++
.../io/external/xlang_jdbcio_it_test.py | 113 ++++
.../apache_beam/io/external/xlang_jmsio_it_test.py | 13 +-
sdks/python/apache_beam/io/gcp/bigquery.py | 90 ++-
sdks/python/apache_beam/io/gcp/bigquery_test.py | 144 ++++-
sdks/python/apache_beam/io/iobase.py | 12 +-
sdks/python/apache_beam/io/iobase_test.py | 21 +
sdks/python/apache_beam/io/jdbc.py | 48 +-
sdks/python/apache_beam/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 | 27 +
.../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 +++++++++
.../databases/jdbc_secret_manager.yaml | 59 ++
.../yaml/extended_tests/databases/snowflake.yaml | 90 +++
sdks/python/apache_beam/yaml/integration_tests.py | 199 ++++++-
sdks/python/apache_beam/yaml/main.py | 10 +
sdks/python/apache_beam/yaml/main_test.py | 14 +
sdks/python/apache_beam/yaml/standard_io.yaml | 116 ++++
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_provider.py | 6 +-
sdks/python/build.gradle | 3 +-
sdks/python/test-suites/dataflow/common.gradle | 6 +-
sdks/python/test-suites/portable/common.gradle | 39 ++
sdks/python/tox.ini | 6 +
settings.gradle.kts | 12 +
.../en/documentation/io/developing-io-python.md | 239 +++++++-
.../site/content/en/documentation/io/managed-io.md | 6 +-
.../en/documentation/runners/kafkastreams.md | 248 ++++++++
website/www/site/content/en/roadmap/_index.md | 40 +-
.../site/content/en/roadmap/connectors-go-sdk.md | 27 -
.../site/content/en/roadmap/connectors-java-sdk.md | 36 --
.../content/en/roadmap/connectors-multi-sdk.md | 92 +--
.../content/en/roadmap/connectors-python-sdk.md | 29 -
.../www/site/content/en/roadmap/dataflow-runner.md | 9 +-
website/www/site/content/en/roadmap/euphoria.md | 2 +
.../www/site/content/en/roadmap/flink-runner.md | 20 +-
website/www/site/content/en/roadmap/go-sdk.md | 15 +-
website/www/site/content/en/roadmap/java-sdk.md | 9 +-
.../www/site/content/en/roadmap/prism-runner.md | 2 +
website/www/site/content/en/roadmap/python-sdk.md | 35 +-
.../www/site/content/en/roadmap/spark-runner.md | 18 +-
website/www/site/data/capability_matrix.yaml | 155 +++++
.../layouts/partials/section-menu/en/roadmap.html | 16 +-
.../layouts/partials/section-menu/en/runners.html | 1 +
296 files changed, 25081 insertions(+), 1180 deletions(-)
copy .github/trigger_files/{IO_Iceberg_Integration_Tests_Dataflow.json =>
beam_PreCommit_Java_Kafka_Streams_Runner.json} (100%)
copy .github/workflows/{beam_PostCommit_XVR_GoUsingJava_Dataflow.yml =>
beam_PreCommit_Java_Kafka_Streams_Runner.yml} (63%)
create mode 100644 runners/kafka-streams/build.gradle
create mode 100644 runners/kafka-streams/job-server/build.gradle
create mode 100644 runners/kafka-streams/measurement/build.gradle
create mode 100644 runners/kafka-streams/measurement/docker-compose.yml
create mode 100644
runners/kafka-streams/measurement/src/main/java/org/apache/beam/runners/kafka/streams/measurement/RescalingMeasurement.java
copy
sdks/java/core/src/main/java/org/apache/beam/sdk/util/construction/CoderTranslatorRegistrar.java
=>
runners/kafka-streams/measurement/src/main/java/org/apache/beam/runners/kafka/streams/measurement/package-info.java
(54%)
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
sdks/java/core/src/main/java/org/apache/beam/sdk/util/construction/CoderTranslatorRegistrar.java
=>
runners/kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/package-info.java
(54%)
create mode 100644
runners/kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/translation/EmptyBoundedSource.java
create mode 100644
runners/kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/translation/ExecutableStageProcessor.java
create mode 100644
runners/kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/translation/ExecutableStageTranslator.java
create mode 100644
runners/kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/translation/FlattenProcessor.java
create mode 100644
runners/kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/translation/FlattenTranslator.java
create mode 100644
runners/kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/translation/GroupByKeyBroadcastPartitioner.java
create mode 100644
runners/kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/translation/GroupByKeyTranslator.java
create mode 100644
runners/kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/translation/ImpulseProcessor.java
create mode 100644
runners/kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/translation/ImpulseTranslator.java
create mode 100644
runners/kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/translation/KStreamsPayload.java
create mode 100644
runners/kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/translation/KStreamsPayloadSerde.java
create mode 100644
runners/kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/translation/KafkaStreamsExecutableStageContextFactory.java
create mode 100644
runners/kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/translation/KafkaStreamsPipelineTranslator.java
create mode 100644
runners/kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/translation/KafkaStreamsStateInternals.java
create mode 100644
runners/kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/translation/KafkaStreamsTimerInternals.java
create mode 100644
runners/kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/translation/KafkaStreamsTranslationContext.java
create mode 100644
runners/kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/translation/PTransformTranslator.java
create mode 100644
runners/kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/translation/ReadProcessor.java
create mode 100644
runners/kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/translation/ReadTranslator.java
create mode 100644
runners/kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/translation/RedistributeTranslator.java
create mode 100644
runners/kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/translation/ShuffleByKeyProcessor.java
create mode 100644
runners/kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/translation/StageOutputProcessor.java
create mode 100644
runners/kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/translation/StoreKeys.java
create mode 100644
runners/kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/translation/TerminationReporter.java
create mode 100644
runners/kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/translation/TerminationTracker.java
create mode 100644
runners/kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/translation/UnboundedReadProcessor.java
create mode 100644
runners/kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/translation/WatermarkAggregator.java
create mode 100644
runners/kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/translation/WatermarkManager.java
create mode 100644
runners/kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/translation/WatermarkPayload.java
create mode 100644
runners/kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/translation/WindowedGroupByKeyProcessor.java
copy
sdks/java/core/src/main/java/org/apache/beam/sdk/util/construction/CoderTranslatorRegistrar.java
=>
runners/kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/translation/package-info.java
(54%)
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/PipelineTranslatorFactory.java
create mode 100644
runners/spark/src/test/java/org/apache/beam/runners/spark/stateful/SparkTimerInternalsTest.java
create mode 100644
sdks/go/pkg/beam/runners/prism/internal/engine/elementmanager_continuation_test.go
create mode 100644
sdks/java/io/google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/bigquery/BigQueryIOIcebergManagedTableIT.java
create mode 100644
sdks/java/io/hbase/src/test/java/org/apache/beam/sdk/io/hbase/HBaseIOCloseTest.java
create mode 100644
sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/BoundedAsyncTasks.java
create mode 100644
sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/NameMappingUtils.java
create mode 100644
sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/ParquetFooters.java
create mode 100644
sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/BigQueryManagedTableCrossEngineIT.java
create mode 100644
sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/BoundedAsyncTasksTest.java
create mode 100644
sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/NameMappingUtilsTest.java
create mode 100644
sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/ParquetFootersTest.java
create mode 100644
sdks/java/io/snowflake/src/main/java/org/apache/beam/sdk/io/snowflake/SnowflakeReadSchemaTransformProvider.java
create mode 100644
sdks/java/io/snowflake/src/main/java/org/apache/beam/sdk/io/snowflake/SnowflakeSchemaTransformUtils.java
create mode 100644
sdks/java/io/snowflake/src/main/java/org/apache/beam/sdk/io/snowflake/SnowflakeWriteConfiguration.java
create mode 100644
sdks/java/io/snowflake/src/main/java/org/apache/beam/sdk/io/snowflake/SnowflakeWriteSchemaTransformProvider.java
create mode 100644
sdks/java/io/snowflake/src/test/java/org/apache/beam/sdk/io/snowflake/SnowflakeReadSchemaTransformProviderTest.java
create mode 100644
sdks/java/io/snowflake/src/test/java/org/apache/beam/sdk/io/snowflake/SnowflakeWriteSchemaTransformProviderTest.java
create mode 100644
sdks/java/io/solace/src/main/java/org/apache/beam/sdk/io/solace/read/ActiveReadersRegistry.java
create mode 100644
sdks/java/io/solace/src/test/java/org/apache/beam/sdk/io/solace/read/UnboundedSolaceReaderTest.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/yaml/extended_tests/databases/jdbc_secret_manager.yaml
create mode 100644
sdks/python/apache_beam/yaml/extended_tests/databases/snowflake.yaml
create mode 100644 sdks/python/apache_beam/yaml/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
delete mode 100644 website/www/site/content/en/roadmap/connectors-go-sdk.md
delete mode 100644 website/www/site/content/en/roadmap/connectors-java-sdk.md
delete mode 100644 website/www/site/content/en/roadmap/connectors-python-sdk.md