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

je-ik pushed a commit to branch 
feat/18479-kafka-streams-runner-skeleton-master-merged
in repository https://gitbox.apache.org/repos/asf/beam.git

commit fd8a62671dfa3aad10f9b55b10d51c3d8f73b85e
Merge: 52dadbac433 d507f1bb3e2
Author: Jan Lukavský <[email protected]>
AuthorDate: Mon Aug 17 13:57:16 2026 +0200

    Merge remote-tracking branch 'origin/master' into 
feat/18479-kafka-streams-runner-skeleton

 .agent/skills/README.md                            |    4 +
 .../skills/developing-new-io-connectors/SKILL.md   |  347 +
 .agent/skills/io-connectors/SKILL.md               |    3 +-
 .agent/skills/yaml-development/SKILL.md            |  107 +
 .asf.yaml                                          |    4 +
 .../test-properties.json                           |    2 +-
 .../actions/setup-environment-action/action.yml    |   32 +-
 .github/actions/setup-k8s-access/action.yml        |    2 +-
 .github/autolabeler.yml                            |    1 +
 .github/build.gradle                               |   10 +-
 .../arc/images/Dockerfile                          |   18 +-
 .../IO_Iceberg_Integration_Tests_Dataflow.json     |    2 +-
 .../beam_CloudML_Benchmarks_Dataflow.json          |    2 +-
 .github/trigger_files/beam_PostCommit_Go.json      |    2 +-
 .../trigger_files/beam_PostCommit_Go_VR_Flink.json |    2 +-
 .../trigger_files/beam_PostCommit_Go_VR_Spark.json |    4 +-
 .../beam_PostCommit_Java_DataflowV1.json           |    2 +-
 .../beam_PostCommit_Java_DataflowV2.json           |    2 +-
 ...=> beam_PostCommit_Java_Delta_IO_Dataflow.json} |    2 +-
 .../beam_PostCommit_Java_Hadoop_Versions.json      |    4 +-
 .../beam_PostCommit_Java_Nexmark_Flink.json        |    2 +-
 ...son => beam_PostCommit_Java_Nexmark_Spark.json} |    0
 .../beam_PostCommit_Java_PVR_Flink_Streaming.json  |    2 +-
 .../beam_PostCommit_Java_PVR_Spark3_Streaming.json |    2 +-
 .../beam_PostCommit_Java_PVR_Spark_Batch.json      |    2 +-
 ....json => beam_PostCommit_Java_Tpcds_Flink.json} |    0
 ....json => beam_PostCommit_Java_Tpcds_Spark.json} |    0
 ...m_PostCommit_Java_ValidatesRunner_Dataflow.json |    2 +-
 ...it_Java_ValidatesRunner_Dataflow_Streaming.json |    2 +-
 ...ostCommit_Java_ValidatesRunner_Dataflow_V2.json |    2 +-
 ...beam_PostCommit_Java_ValidatesRunner_Spark.json |   18 +-
 ...json => beam_PostCommit_PortableJar_Spark.json} |    0
 .github/trigger_files/beam_PostCommit_Python.json  |    4 +-
 ...taflow.json => beam_PostCommit_Python_Arm.json} |    2 +-
 .../beam_PostCommit_Python_Dependency.json         |    4 +-
 .../beam_PostCommit_Python_Examples_Dataflow.json  |    4 +-
 ... => beam_PostCommit_Python_Examples_Spark.json} |    0
 ...tCommit_Python_ValidatesContainer_Dataflow.json |    2 +-
 ...PostCommit_Python_ValidatesRunner_Dataflow.json |    3 +-
 ...am_PostCommit_Python_ValidatesRunner_Flink.json |    3 +-
 ...am_PostCommit_Python_ValidatesRunner_Spark.json |    5 +-
 .../beam_PostCommit_Python_Versions.json           |    2 +-
 .../beam_PostCommit_Python_Xlang_Gcp_Direct.json   |    2 +-
 .../beam_PostCommit_Python_Xlang_IO_Dataflow.json  |    2 +-
 .../beam_PostCommit_Python_Xlang_IO_Direct.json    |    2 +-
 ..._PostCommit_Python_Xlang_Messaging_Direct.json} |    2 +-
 .github/trigger_files/beam_PostCommit_SQL.json     |    2 +-
 .../trigger_files/beam_PostCommit_XVR_Flink.json   |    2 +-
 .../trigger_files/beam_PostCommit_XVR_Samza.json   |    1 -
 .../trigger_files/beam_PostCommit_XVR_Spark3.json  |    2 +-
 .../beam_PostCommit_Yaml_Xlang_Direct.json         |    2 +-
 .../beam_PreCommit_Flink_Container.json            |    2 +-
 .github/trigger_files/beam_PreCommit_Java.json     |    2 +-
 ...on => beam_PreCommit_Java_Delta_IO_Direct.json} |    0
 .../trigger_files/beam_PreCommit_Python_ML.json    |    2 +-
 ...t.json => beam_PreCommit_Python_PVR_Flink.json} |    0
 .github/trigger_files/beam_PreCommit_SQL.json      |    2 +-
 .../beam_PreCommit_Yaml_Xlang_Direct.json          |    2 +-
 .github/workflows/IO_Iceberg_Integration_Tests.yml |    6 +-
 .../IO_Iceberg_Integration_Tests_Dataflow.yml      |    6 +-
 ..._Iceberg_Managed_Integration_Tests_Dataflow.yml |    6 +-
 .github/workflows/IO_Iceberg_Performance_Tests.yml |    6 +-
 .github/workflows/IO_Iceberg_Unit_Tests.yml        |    4 +-
 .github/workflows/README.md                        |   17 +-
 .github/workflows/assign_milestone.yml             |    3 +-
 .github/workflows/beam_CancelStaleDataflowJobs.yml |    4 +-
 .../workflows/beam_CleanUpDataprocResources.yml    |    6 +-
 .github/workflows/beam_CleanUpGCPResources.yml     |    4 +-
 .../workflows/beam_CleanUpPrebuiltSDKImages.yml    |    6 +-
 .../workflows/beam_CloudML_Benchmarks_Dataflow.yml |    4 +-
 .../beam_IODatastoresCredentialsRotation.yml       |    4 +-
 .../beam_Inference_Python_Benchmarks_Dataflow.yml  |  207 +-
 .../beam_Infrastructure_AuditUnmanagedKeys.yml     |   76 +
 .../beam_Infrastructure_PolicyEnforcer.yml         |    6 +-
 .../beam_Infrastructure_SecurityLogging.yml        |    6 +-
 .../beam_Infrastructure_ServiceAccountKeys.yml     |    6 +-
 .../beam_Infrastructure_UsersPermissions.yml       |    5 +-
 .github/workflows/beam_Java_JMH.yml                |    6 +-
 .../beam_Java_LoadTests_Combine_Smoke.yml          |    6 +-
 .../beam_LoadTests_Go_CoGBK_Dataflow_Batch.yml     |    6 +-
 .../beam_LoadTests_Go_CoGBK_Flink_batch.yml        |   12 +-
 .../beam_LoadTests_Go_Combine_Dataflow_Batch.yml   |    6 +-
 .../beam_LoadTests_Go_Combine_Flink_Batch.yml      |   18 +-
 .../beam_LoadTests_Go_GBK_Dataflow_Batch.yml       |    6 +-
 .../beam_LoadTests_Go_GBK_Flink_Batch.yml          |   18 +-
 .../beam_LoadTests_Go_ParDo_Dataflow_Batch.yml     |    6 +-
 .../beam_LoadTests_Go_ParDo_Flink_Batch.yml        |   10 +-
 .../beam_LoadTests_Go_SideInput_Dataflow_Batch.yml |    6 +-
 .../beam_LoadTests_Go_SideInput_Flink_Batch.yml    |   12 +-
 .../beam_LoadTests_Java_CoGBK_Dataflow_Batch.yml   |    6 +-
 ...eam_LoadTests_Java_CoGBK_Dataflow_Streaming.yml |    6 +-
 ...s_Java_CoGBK_Dataflow_V2_Batch_JavaVersions.yml |    6 +-
 ...va_CoGBK_Dataflow_V2_Streaming_JavaVersions.yml |    6 +-
 ...s_Java_CoGBK_SparkStructuredStreaming_Batch.yml |    6 +-
 .../beam_LoadTests_Java_Combine_Dataflow_Batch.yml |    6 +-
 ...m_LoadTests_Java_Combine_Dataflow_Streaming.yml |    6 +-
 ...Java_Combine_SparkStructuredStreaming_Batch.yml |    6 +-
 .../beam_LoadTests_Java_GBK_Dataflow_Batch.yml     |    6 +-
 .../beam_LoadTests_Java_GBK_Dataflow_Streaming.yml |    6 +-
 .../beam_LoadTests_Java_GBK_Dataflow_V2_Batch.yml  |    6 +-
 ...LoadTests_Java_GBK_Dataflow_V2_Batch_Java17.yml |    6 +-
 ...am_LoadTests_Java_GBK_Dataflow_V2_Streaming.yml |    6 +-
 ...Tests_Java_GBK_Dataflow_V2_Streaming_Java17.yml |    6 +-
 .../workflows/beam_LoadTests_Java_GBK_Smoke.yml    |    4 +-
 ...sts_Java_GBK_SparkStructuredStreaming_Batch.yml |    6 +-
 .../beam_LoadTests_Java_ParDo_Dataflow_Batch.yml   |    6 +-
 ...eam_LoadTests_Java_ParDo_Dataflow_Streaming.yml |    6 +-
 ...s_Java_ParDo_Dataflow_V2_Batch_JavaVersions.yml |    6 +-
 ...va_ParDo_Dataflow_V2_Streaming_JavaVersions.yml |    6 +-
 ...s_Java_ParDo_SparkStructuredStreaming_Batch.yml |    6 +-
 .github/workflows/beam_LoadTests_Java_PubsubIO.yml |    6 +-
 .../beam_LoadTests_Python_CoGBK_Dataflow_Batch.yml |    6 +-
 ...m_LoadTests_Python_CoGBK_Dataflow_Streaming.yml |    6 +-
 .../beam_LoadTests_Python_CoGBK_Flink_Batch.yml    |   10 +-
 ...eam_LoadTests_Python_Combine_Dataflow_Batch.yml |    6 +-
 ...LoadTests_Python_Combine_Dataflow_Streaming.yml |    6 +-
 .../beam_LoadTests_Python_Combine_Flink_Batch.yml  |   20 +-
 ...am_LoadTests_Python_Combine_Flink_Streaming.yml |   10 +-
 ...LoadTests_Python_FnApiRunner_Microbenchmark.yml |    6 +-
 .../beam_LoadTests_Python_GBK_Dataflow_Batch.yml   |    6 +-
 ...eam_LoadTests_Python_GBK_Dataflow_Streaming.yml |    6 +-
 .../beam_LoadTests_Python_GBK_Flink_Batch.yml      |   18 +-
 ...adTests_Python_GBK_reiterate_Dataflow_Batch.yml |    6 +-
 ...sts_Python_GBK_reiterate_Dataflow_Streaming.yml |    4 +-
 .../beam_LoadTests_Python_ParDo_Dataflow_Batch.yml |    6 +-
 ...m_LoadTests_Python_ParDo_Dataflow_Streaming.yml |    6 +-
 .../beam_LoadTests_Python_ParDo_Flink_Batch.yml    |   12 +-
 ...beam_LoadTests_Python_ParDo_Flink_Streaming.yml |   12 +-
 ...m_LoadTests_Python_SideInput_Dataflow_Batch.yml |    6 +-
 .github/workflows/beam_LoadTests_Python_Smoke.yml  |    6 +-
 .../workflows/beam_MetricsCredentialsRotation.yml  |    4 +-
 .github/workflows/beam_Metrics_Report.yml          |    4 +-
 .../workflows/beam_PerformanceTests_AvroIOIT.yml   |    4 +-
 .../beam_PerformanceTests_AvroIOIT_HDFS.yml        |    4 +-
 ...PerformanceTests_BigQueryIO_Batch_Java_Avro.yml |    6 +-
 ...PerformanceTests_BigQueryIO_Batch_Java_Json.yml |    6 +-
 ..._PerformanceTests_BigQueryIO_Streaming_Java.yml |    6 +-
 ...eam_PerformanceTests_BiqQueryIO_Read_Python.yml |    6 +-
 ...formanceTests_BiqQueryIO_Write_Python_Batch.yml |    6 +-
 .github/workflows/beam_PerformanceTests_Cdap.yml   |    4 +-
 .../beam_PerformanceTests_Compressed_TextIOIT.yml  |    6 +-
 ...m_PerformanceTests_Compressed_TextIOIT_HDFS.yml |    6 +-
 .../beam_PerformanceTests_HadoopFormat.yml         |    4 +-
 .github/workflows/beam_PerformanceTests_JDBC.yml   |    6 +-
 .../workflows/beam_PerformanceTests_Kafka_IO.yml   |    6 +-
 .../beam_PerformanceTests_ManyFiles_TextIOIT.yml   |    6 +-
 ...am_PerformanceTests_ManyFiles_TextIOIT_HDFS.yml |    6 +-
 .../beam_PerformanceTests_MongoDBIO_IT.yml         |    4 +-
 .../beam_PerformanceTests_ParquetIOIT.yml          |    6 +-
 .../beam_PerformanceTests_ParquetIOIT_HDFS.yml     |    6 +-
 ...erformanceTests_PubsubIOIT_Python_Streaming.yml |    6 +-
 ...m_PerformanceTests_SQLBigQueryIO_Batch_Java.yml |    6 +-
 .../beam_PerformanceTests_SingleStoreIO.yml        |    8 +-
 ..._PerformanceTests_SpannerIO_Read_2GB_Python.yml |    6 +-
 ...manceTests_SpannerIO_Write_2GB_Python_Batch.yml |    6 +-
 .../beam_PerformanceTests_SparkReceiver_IO.yml     |    4 +-
 .../beam_PerformanceTests_TFRecordIOIT.yml         |    6 +-
 .../beam_PerformanceTests_TFRecordIOIT_HDFS.yml    |    6 +-
 .../workflows/beam_PerformanceTests_TextIOIT.yml   |    6 +-
 .../beam_PerformanceTests_TextIOIT_HDFS.yml        |    6 +-
 .../beam_PerformanceTests_TextIOIT_Python.yml      |    6 +-
 ...PerformanceTests_WordCountIT_PythonVersions.yml |    6 +-
 .../workflows/beam_PerformanceTests_XmlIOIT.yml    |    6 +-
 .../beam_PerformanceTests_XmlIOIT_HDFS.yml         |    6 +-
 .../beam_PerformanceTests_xlang_KafkaIO_Python.yml |    5 +-
 .github/workflows/beam_Playground_CI_Nightly.yml   |    6 +-
 .github/workflows/beam_Playground_Precommit.yml    |    7 +-
 .github/workflows/beam_PostCommit_Go.yml           |    8 +-
 .../workflows/beam_PostCommit_Go_Dataflow_ARM.yml  |    6 +-
 .github/workflows/beam_PostCommit_Go_VR_Flink.yml  |    6 +-
 .github/workflows/beam_PostCommit_Go_VR_Spark.yml  |    6 +-
 .github/workflows/beam_PostCommit_Java.yml         |    6 +-
 .../beam_PostCommit_Java_Avro_Versions.yml         |    6 +-
 .../beam_PostCommit_Java_BigQueryEarlyRollout.yml  |    4 +-
 .../workflows/beam_PostCommit_Java_DataflowV1.yml  |    6 +-
 .../workflows/beam_PostCommit_Java_DataflowV2.yml  |    8 +-
 ... => beam_PostCommit_Java_Delta_IO_Dataflow.yml} |   30 +-
 .../beam_PostCommit_Java_Examples_Dataflow.yml     |    6 +-
 .../beam_PostCommit_Java_Examples_Dataflow_ARM.yml |    6 +-
 ...beam_PostCommit_Java_Examples_Dataflow_Java.yml |    6 +-
 .../beam_PostCommit_Java_Examples_Dataflow_V2.yml  |    4 +-
 ...m_PostCommit_Java_Examples_Dataflow_V2_Java.yml |    6 +-
 .../beam_PostCommit_Java_Examples_Direct.yml       |    6 +-
 .../beam_PostCommit_Java_Examples_Flink.yml        |    6 +-
 .../beam_PostCommit_Java_Examples_Spark.yml        |    6 +-
 .../beam_PostCommit_Java_Hadoop_Versions.yml       |    6 +-
 .../beam_PostCommit_Java_IO_Performance_Tests.yml  |    7 +-
 .../beam_PostCommit_Java_InfluxDbIO_IT.yml         |    4 +-
 .../beam_PostCommit_Java_Jpms_Dataflow.yml         |    7 +-
 ...beam_PostCommit_Java_Jpms_Dataflow_Versions.yml |    4 +-
 .../workflows/beam_PostCommit_Java_Jpms_Direct.yml |    6 +-
 .../beam_PostCommit_Java_Jpms_Direct_Versions.yml  |    4 +-
 .../beam_PostCommit_Java_Jpms_Flink_Java11.yml     |    6 +-
 .../beam_PostCommit_Java_Jpms_Spark_Java11.yml     |    6 +-
 .../beam_PostCommit_Java_Nexmark_Dataflow.yml      |    6 +-
 .../beam_PostCommit_Java_Nexmark_Dataflow_V2.yml   |    6 +-
 ...am_PostCommit_Java_Nexmark_Dataflow_V2_Java.yml |    6 +-
 .../beam_PostCommit_Java_Nexmark_Direct.yml        |    6 +-
 .../beam_PostCommit_Java_Nexmark_Flink.yml         |    9 +-
 .../beam_PostCommit_Java_Nexmark_Spark.yml         |    6 +-
 .../beam_PostCommit_Java_PVR_Flink_Batch.yml       |    6 +-
 .../beam_PostCommit_Java_PVR_Flink_Streaming.yml   |    8 +-
 .../beam_PostCommit_Java_PVR_Spark3_Streaming.yml  |    6 +-
 .../beam_PostCommit_Java_PVR_Spark4_Batch.yml      |    4 +-
 .../beam_PostCommit_Java_PVR_Spark4_Streaming.yml  |    6 +-
 .../beam_PostCommit_Java_PVR_Spark_Batch.yml       |    4 +-
 .../beam_PostCommit_Java_SingleStoreIO_IT.yml      |    6 +-
 .../beam_PostCommit_Java_Tpcds_Dataflow.yml        |    4 +-
 .../workflows/beam_PostCommit_Java_Tpcds_Flink.yml |    6 +-
 .../workflows/beam_PostCommit_Java_Tpcds_Spark.yml |    4 +-
 ...am_PostCommit_Java_ValidatesRunner_Dataflow.yml |    6 +-
 ..._Java_ValidatesRunner_Dataflow_JavaVersions.yml |    8 +-
 ...mit_Java_ValidatesRunner_Dataflow_Streaming.yml |    6 +-
 ...a_ValidatesRunner_Dataflow_Streaming_Engine.yml |    4 +-
 ...atesRunner_Dataflow_Streaming_TagEncodingV2.yml |    4 +-
 ...PostCommit_Java_ValidatesRunner_Dataflow_V2.yml |    6 +-
 ..._Java_ValidatesRunner_Dataflow_V2_Streaming.yml |    6 +-
 ...beam_PostCommit_Java_ValidatesRunner_Direct.yml |    6 +-
 ...it_Java_ValidatesRunner_Direct_JavaVersions.yml |    6 +-
 .../beam_PostCommit_Java_ValidatesRunner_Flink.yml |    8 +-
 .../beam_PostCommit_Java_ValidatesRunner_Spark.yml |    6 +-
 ...beam_PostCommit_Java_ValidatesRunner_Spark4.yml |    4 +-
 ...va_ValidatesRunner_SparkStructuredStreaming.yml |    6 +-
 ...am_PostCommit_Java_ValidatesRunner_Twister2.yml |    6 +-
 .../beam_PostCommit_Java_ValidatesRunner_ULR.yml   |    6 +-
 .github/workflows/beam_PostCommit_Javadoc.yml      |    6 +-
 .../beam_PostCommit_PortableJar_Flink.yml          |    4 +-
 .../beam_PostCommit_PortableJar_Spark.yml          |    6 +-
 .github/workflows/beam_PostCommit_Python.yml       |    4 +-
 .github/workflows/beam_PostCommit_Python_Arm.yml   |   24 +-
 .../beam_PostCommit_Python_Dependency.yml          |    4 +-
 .../beam_PostCommit_Python_Examples_Dataflow.yml   |    4 +-
 .../beam_PostCommit_Python_Examples_Direct.yml     |    6 +-
 .../beam_PostCommit_Python_Examples_Flink.yml      |    4 +-
 .../beam_PostCommit_Python_Examples_Spark.yml      |    6 +-
 .../beam_PostCommit_Python_MongoDBIO_IT.yml        |    4 +-
 .../beam_PostCommit_Python_Nexmark_Direct.yml      |    6 +-
 .../beam_PostCommit_Python_Portable_Flink.yml      |    4 +-
 ...stCommit_Python_ValidatesContainer_Dataflow.yml |    4 +-
 ..._Python_ValidatesContainer_Dataflow_With_RC.yml |    6 +-
 ..._PostCommit_Python_ValidatesRunner_Dataflow.yml |    6 +-
 ...eam_PostCommit_Python_ValidatesRunner_Flink.yml |    6 +-
 ...eam_PostCommit_Python_ValidatesRunner_Spark.yml |    6 +-
 .../workflows/beam_PostCommit_Python_Versions.yml  |    4 +-
 .../beam_PostCommit_Python_Xlang_Gcp_Dataflow.yml  |    7 +-
 .../beam_PostCommit_Python_Xlang_Gcp_Direct.yml    |    5 +-
 .../beam_PostCommit_Python_Xlang_IO_Dataflow.yml   |    7 +-
 .../beam_PostCommit_Python_Xlang_IO_Direct.yml     |   10 +-
 ...m_PostCommit_Python_Xlang_Messaging_Direct.yml} |   24 +-
 .github/workflows/beam_PostCommit_SQL.yml          |    4 +-
 .../beam_PostCommit_TransformService_Direct.yml    |   10 +-
 .github/workflows/beam_PostCommit_Website_Test.yml |    6 +-
 .github/workflows/beam_PostCommit_XVR_Direct.yml   |    7 +-
 .github/workflows/beam_PostCommit_XVR_Flink.yml    |    7 +-
 .../beam_PostCommit_XVR_GoUsingJava_Dataflow.yml   |    7 +-
 ...eam_PostCommit_XVR_JavaUsingPython_Dataflow.yml |    9 +-
 ..._PostCommit_XVR_PythonUsingJavaSQL_Dataflow.yml |    7 +-
 ...eam_PostCommit_XVR_PythonUsingJava_Dataflow.yml |    7 +-
 .github/workflows/beam_PostCommit_XVR_Spark3.yml   |    7 +-
 .../beam_PostCommit_Yaml_Xlang_Direct.yml          |   10 +-
 .../workflows/beam_PostRelease_NightlySnapshot.yml |   17 +-
 .../workflows/beam_PreCommit_CommunityMetrics.yml  |    8 +-
 .../workflows/beam_PreCommit_Flink_Container.yml   |   14 +-
 .github/workflows/beam_PreCommit_GHA.yml           |   32 +-
 .github/workflows/beam_PreCommit_Go.yml            |    6 +-
 .github/workflows/beam_PreCommit_GoPortable.yml    |    4 +-
 .github/workflows/beam_PreCommit_GoPrism.yml       |    6 +-
 .github/workflows/beam_PreCommit_ItFramework.yml   |    6 +-
 .github/workflows/beam_PreCommit_Java.yml          |    4 +-
 ...eCommit_Java_Amazon-Web-Services2_IO_Direct.yml |    6 +-
 .../beam_PreCommit_Java_Amqp_IO_Direct.yml         |    6 +-
 .../beam_PreCommit_Java_Azure_IO_Direct.yml        |    6 +-
 .../beam_PreCommit_Java_Cassandra_IO_Direct.yml    |    8 +-
 .../beam_PreCommit_Java_Cdap_IO_Direct.yml         |    6 +-
 .../beam_PreCommit_Java_Clickhouse_IO_Direct.yml   |    6 +-
 .../beam_PreCommit_Java_Csv_IO_Direct.yml          |    6 +-
 .../beam_PreCommit_Java_Datadog_IO_Direct.yml      |    6 +-
 .github/workflows/beam_PreCommit_Java_Dataflow.yml |    6 +-
 .../beam_PreCommit_Java_Debezium_IO_Direct.yml     |   16 +-
 ...yml => beam_PreCommit_Java_Delta_IO_Direct.yml} |   36 +-
 ...beam_PreCommit_Java_ElasticSearch_IO_Direct.yml |    4 +-
 .../beam_PreCommit_Java_Examples_Dataflow.yml      |    6 +-
 ...am_PreCommit_Java_Examples_Dataflow_Java11.yml} |   24 +-
 ...Commit_Java_File-schema-transform_IO_Direct.yml |    6 +-
 .../beam_PreCommit_Java_Flink_Versions.yml         |    8 +-
 .../beam_PreCommit_Java_GCP_IO_Direct.yml          |    6 +-
 .../beam_PreCommit_Java_Google-ads_IO_Direct.yml   |    6 +-
 .../beam_PreCommit_Java_HBase_IO_Direct.yml        |    8 +-
 .../beam_PreCommit_Java_HCatalog_IO_Direct.yml     |    8 +-
 .../beam_PreCommit_Java_Hadoop_IO_Direct.yml       |    6 +-
 .../workflows/beam_PreCommit_Java_IOs_Direct.yml   |   19 +-
 .../beam_PreCommit_Java_InfluxDb_IO_Direct.yml     |    6 +-
 .../beam_PreCommit_Java_JDBC_IO_Direct.yml         |    6 +-
 .../beam_PreCommit_Java_Jms_IO_Direct.yml          |    6 +-
 .../beam_PreCommit_Java_Kafka_IO_Direct.yml        |    4 +-
 .../beam_PreCommit_Java_Kudu_IO_Direct.yml         |    6 +-
 .../beam_PreCommit_Java_MongoDb_IO_Direct.yml      |    6 +-
 .../beam_PreCommit_Java_Mqtt_IO_Direct.yml         |    6 +-
 .../beam_PreCommit_Java_Neo4j_IO_Direct.yml        |    6 +-
 .../beam_PreCommit_Java_PVR_Flink_Batch.yml        |    6 +-
 .../beam_PreCommit_Java_PVR_Flink_Docker.yml       |    6 +-
 .../beam_PreCommit_Java_PVR_Prism_Loopback.yml     |    4 +-
 .../beam_PreCommit_Java_Parquet_IO_Direct.yml      |    6 +-
 .../beam_PreCommit_Java_Pulsar_IO_Direct.yml       |    4 +-
 .../beam_PreCommit_Java_RabbitMq_IO_Direct.yml     |    6 +-
 .../beam_PreCommit_Java_Redis_IO_Direct.yml        |    6 +-
 ...am_PreCommit_Java_RequestResponse_IO_Direct.yml |    6 +-
 .../beam_PreCommit_Java_SingleStore_IO_Direct.yml  |    6 +-
 .../beam_PreCommit_Java_Snowflake_IO_Direct.yml    |    6 +-
 .../beam_PreCommit_Java_Solace_IO_Direct.yml       |    6 +-
 .../beam_PreCommit_Java_Solr_IO_Direct.yml         |    6 +-
 .../beam_PreCommit_Java_Spark_Versions.yml         |   10 +-
 .../beam_PreCommit_Java_Splunk_IO_Direct.yml       |    6 +-
 .../beam_PreCommit_Java_Thrift_IO_Direct.yml       |    6 +-
 .../beam_PreCommit_Java_Tika_IO_Direct.yml         |    6 +-
 .../workflows/beam_PreCommit_Kotlin_Examples.yml   |   12 +-
 .../workflows/beam_PreCommit_Portable_Python.yml   |    6 +-
 .github/workflows/beam_PreCommit_Prism_Python.yml  |    6 +-
 .github/workflows/beam_PreCommit_Python.yml        |    4 +-
 .github/workflows/beam_PreCommit_PythonDocker.yml  |    8 +-
 .github/workflows/beam_PreCommit_PythonDocs.yml    |    4 +-
 .../workflows/beam_PreCommit_PythonFormatter.yml   |    4 +-
 .github/workflows/beam_PreCommit_PythonLint.yml    |    4 +-
 .../workflows/beam_PreCommit_Python_Coverage.yml   |    4 +-
 .../workflows/beam_PreCommit_Python_Dataframes.yml |    6 +-
 .github/workflows/beam_PreCommit_Python_Dill.yml   |    4 +-
 .../workflows/beam_PreCommit_Python_Examples.yml   |    4 +-
 .../beam_PreCommit_Python_Integration.yml          |    6 +-
 .github/workflows/beam_PreCommit_Python_ML.yml     |    4 +-
 .../workflows/beam_PreCommit_Python_PVR_Flink.yml  |    4 +-
 .../workflows/beam_PreCommit_Python_Runners.yml    |    6 +-
 .../workflows/beam_PreCommit_Python_Transforms.yml |    4 +-
 .github/workflows/beam_PreCommit_RAT.yml           |    4 +-
 .github/workflows/beam_PreCommit_SQL.yml           |    6 +-
 .github/workflows/beam_PreCommit_SQL_Java17.yml    |    6 +-
 .github/workflows/beam_PreCommit_Spotless.yml      |    6 +-
 .github/workflows/beam_PreCommit_Typescript.yml    |    6 +-
 .github/workflows/beam_PreCommit_Website.yml       |    6 +-
 .../workflows/beam_PreCommit_Website_Stage_GCS.yml |    4 +-
 .github/workflows/beam_PreCommit_Whitespace.yml    |    4 +-
 .../beam_PreCommit_Xlang_Generated_Transforms.yml  |    6 +-
 .../workflows/beam_PreCommit_Yaml_Xlang_Direct.yml |   10 +-
 .github/workflows/beam_Prober_CommunityMetrics.yml |    6 +-
 .github/workflows/beam_Publish_BeamMetrics.yml     |    4 +-
 .../workflows/beam_Publish_Beam_SDK_Snapshots.yml  |   30 +-
 .../workflows/beam_Publish_Docker_Snapshots.yml    |    8 +-
 ...Jobs.yml => beam_Publish_Python_VLLM_Image.yml} |   71 +-
 .github/workflows/beam_Publish_Website.yml         |    8 +-
 .../beam_Python_CostBenchmarks_Dataflow.yml        |    6 +-
 ...beam_Python_ValidatesContainer_Dataflow_ARM.yml |    8 +-
 .github/workflows/beam_Release_NightlySnapshot.yml |    4 +-
 .../beam_Release_Python_NightlySnapshot.yml        |    6 +-
 .../workflows/beam_StressTests_Java_BigQueryIO.yml |    6 +-
 .../workflows/beam_StressTests_Java_BigTableIO.yml |    6 +-
 .../workflows/beam_StressTests_Java_KafkaIO.yml    |    4 +-
 .../workflows/beam_StressTests_Java_PubSubIO.yml   |    6 +-
 .../workflows/beam_StressTests_Java_SpannerIO.yml  |    6 +-
 .github/workflows/beam_Upgrade_GCP_BOM.yml         |   98 +
 .github/workflows/build_release_candidate.yml      |  157 +-
 .github/workflows/build_runner_image.yml           |    9 +-
 .github/workflows/build_wheels.yml                 |   44 +-
 .github/workflows/cancel.yml                       |    2 +-
 .github/workflows/code_completion_plugin_tests.yml |    6 +-
 .github/workflows/codeql.yml                       |  197 +
 .github/workflows/cut_release_branch.yml           |   13 +-
 .github/workflows/dask_runner_tests.yml            |   12 +-
 .../workflows/deploy_release_candidate_pypi.yaml   |   14 +-
 .github/workflows/finalize_release.yml             |   41 +-
 .github/workflows/flaky_test_detection.yml         |    6 +-
 .github/workflows/git_tag_released_version.yml     |   18 +-
 .github/workflows/go_tests.yml                     |    3 +-
 .github/workflows/issue-tagger.yml                 |    4 +-
 .github/workflows/java_tests.yml                   |    4 +-
 ...s_Dataflow_MLTransform_Generate_Vocab_Batch.txt |   42 +
 ...aflow_MLTransform_Image_Embedding_CPU_Batch.txt |   42 +
 ...aflow_MLTransform_Image_Embedding_GPU_Batch.txt |   44 +
 ...ataflow_MLTransform_One_Hot_Encoding_Batch.txt} |   19 +-
 ...s_Dataflow_MLTransform_Text_Embedding_Batch.txt |   42 +
 ...nchmarks_Dataflow_Pytorch_Image_Captioning.txt} |   28 +-
 ...flow_Pytorch_Image_Classification_Rightfit.txt} |   30 +-
 ...ks_Dataflow_Pytorch_Image_Object_Detection.txt} |   24 +-
 ...rch_Sentiment_Batch_DistilBert_Base_Uncased.txt |    3 +
 ...Sentiment_Streaming_DistilBert_Base_Uncased.txt |    3 +
 .../go_CoGBK_Flink_Batch_MultipleKey.txt           |    4 +-
 .../go_CoGBK_Flink_Batch_Reiteration_10KB.txt      |    4 +-
 .../go_CoGBK_Flink_Batch_Reiteration_2MB.txt       |    4 +-
 .../go_GBK_Flink_Batch_100kb.txt                   |    2 +-
 .../go_GBK_Flink_Batch_Fanout_4.txt                |    2 +-
 .../go_GBK_Flink_Batch_Fanout_8.txt                |    2 +-
 .../go_GBK_Flink_Batch_Reiteration_10KB.txt        |    2 +-
 .github/workflows/local_env_tests.yml              |    8 +-
 .github/workflows/playground_frontend_test.yml     |    8 +-
 .github/workflows/pr-bot-new-prs.yml               |    6 +-
 .github/workflows/pr-bot-pr-updates.yml            |    5 +-
 .github/workflows/pr-bot-prs-needing-attention.yml |    8 +-
 .github/workflows/publish_github_release_notes.yml |   14 +-
 .github/workflows/python_dependency_tests.yml      |    7 +-
 .github/workflows/python_tests.yml                 |   16 +-
 .github/workflows/refresh_looker_metrics.yml       |    6 +-
 .github/workflows/reportGenerator.yml              |    6 +-
 .../republish_released_docker_containers.yml       |   18 +-
 .github/workflows/run_perf_alert_tool.yml          |    6 +-
 .../workflows/run_rc_validation_go_wordcount.yml   |   18 +-
 .../run_rc_validation_java_mobile_gaming.yml       |   23 +-
 .../run_rc_validation_java_quickstart.yml          |    5 +-
 .../run_rc_validation_python_mobile_gaming.yml     |   83 +-
 .../workflows/run_rc_validation_python_yaml.yml    |   39 +-
 .github/workflows/stale.yml                        |    2 +-
 .github/workflows/tour_of_beam_backend.yml         |   17 +-
 .../workflows/tour_of_beam_backend_integration.yml |    5 +-
 .github/workflows/tour_of_beam_frontend_test.yml   |    8 +-
 .github/workflows/typescript_tests.yml             |   30 +-
 .github/workflows/update_python_dependencies.yml   |   11 +-
 .github/zizmor.yml                                 |    5 +
 .test-infra/dataproc/flink_cluster.sh              |   92 +-
 .test-infra/dataproc/init-actions/docker.sh        |  103 -
 .test-infra/dataproc/init-actions/flink.sh         |  217 -
 .test-infra/metrics/build.gradle                   |    1 +
 .test-infra/metrics/influxdb/Dockerfile            |    9 +-
 .test-infra/metrics/influxdb/gsutil/.boto          |   24 -
 .test-infra/metrics/influxdb/gsutil/Dockerfile     |   25 -
 .../kubernetes/beam-influxdb-autobackup.yaml       |    7 +-
 .../metrics/src/test/groovy/ProberTests.groovy     |    8 +-
 .../github/github_runs_prefetcher/code/main.py     |   12 +-
 .test-infra/mock-apis/go.mod                       |   10 +-
 .test-infra/mock-apis/go.sum                       |   20 +-
 .test-infra/mock-apis/poetry.lock                  |   32 +-
 .test-infra/tools/refresh_looker_metrics.py        |   28 +-
 .test-infra/tools/stale_cleaner.py                 |   17 +-
 .test-infra/tools/test_stale_cleaner.py            |   61 +-
 .test-infra/validate-runner/build.gradle           |   13 +
 CHANGES.md                                         |  177 +-
 build.gradle.kts                                   |    8 +-
 buildSrc/build.gradle.kts                          |   19 +
 .../org/apache/beam/gradle/BeamDockerPlugin.groovy |   24 +-
 .../org/apache/beam/gradle/BeamModulePlugin.groovy |  141 +-
 .../org/apache/beam/gradle/Repositories.groovy     |   26 +-
 contributor-docs/README.md                         |    1 +
 contributor-docs/local-flink-python.md             |  204 +
 contributor-docs/release-guide.md                  |   12 +-
 examples/java/build.gradle                         |   11 +-
 examples/java/common.gradle                        |    3 +
 examples/java/iceberg/build.gradle                 |    4 +-
 .../apache/beam/examples/BatchElementsExample.java |   81 +
 .../beam/examples/complete/game/UserScore.java     |    2 +-
 .../beam/examples/subprocess/utils/FileUtils.java  |   30 +-
 .../examples/subprocess/utils/FileUtilsTest.java   |   80 +
 examples/kotlin/build.gradle                       |    7 +
 examples/multi-language/README.md                  |    6 +-
 .../beam-ml/automatic_model_refresh.ipynb          |    4 +-
 .../beam-ml/rag_usecase/opensearch_connector.py    |    8 +-
 .../beam-ml/rag_usecase/opensearch_enrichment.py   |    4 +-
 .../notebooks/beam-ml/run_inference_gemini.ipynb   |  609 ++
 .../get-started/learn_beam_basics_by_doing.ipynb   |    4 +-
 .../learn_beam_transforms_by_doing.ipynb           |    2 +-
 .../learn_beam_windowing_by_doing.ipynb            |    2 +-
 .../notebooks/get-started/try-apache-beam-go.ipynb |    4 +-
 .../get-started/try-apache-beam-java.ipynb         |    4 +-
 .../notebooks/get-started/try-apache-beam-py.ipynb |    4 +-
 gradle.properties                                  |    9 +-
 infra/enforcement/README.md                        |   12 +-
 infra/enforcement/account_keys.py                  |  110 +-
 infra/enforcement/sending.py                       |  169 +-
 infra/enforcement/test_sending.py                  |  103 +
 infra/iam/roles/roles_config.yaml                  |    2 +-
 infra/iam/users.yml                                |   17 +-
 it/build.gradle                                    |    6 +-
 it/clickhouse/build.gradle                         |    8 +-
 it/common/build.gradle                             |   29 +-
 .../apache/beam/it/common}/artifacts/Artifact.java |    2 +-
 .../beam/it/common}/artifacts/ArtifactClient.java  |    2 +-
 .../beam/it/common}/artifacts/GcsArtifact.java     |    2 +-
 .../beam/it/common}/artifacts/package-info.java    |    2 +-
 .../it/common}/artifacts/utils/ArtifactUtils.java  |    2 +-
 .../it/common}/artifacts/utils/AvroTestUtil.java   |    2 +-
 .../it/common}/artifacts/utils/JsonTestUtil.java   |    2 +-
 .../common}/artifacts/utils/ParquetTestUtil.java   |    2 +-
 .../it/common}/artifacts/utils/package-info.java   |    2 +-
 .../common}/bigquery/BigQueryResourceManager.java  |   44 +-
 .../bigquery/BigQueryResourceManagerException.java |    2 +-
 .../bigquery/BigQueryResourceManagerUtils.java     |    2 +-
 .../beam/it/common/bigquery}/package-info.java     |    4 +-
 .../common}/dataflow/AbstractPipelineLauncher.java |    2 +-
 .../common}/dataflow/DefaultPipelineLauncher.java  |   52 +-
 .../beam/it/common/dataflow}/package-info.java     |    4 +-
 .../it/common}/monitoring/MonitoringClient.java    |    2 +-
 .../beam/it/common}/monitoring/package-info.java   |    2 +-
 .../it/common}/storage/GcsResourceManager.java     |   10 +-
 .../beam/it/common}/storage/package-info.java      |    2 +-
 .../src/main/resources/test-artifact.json          |    0
 .../beam/it/common}/artifacts/GcsArtifactTest.java |    4 +-
 .../common}/artifacts/utils/ArtifactUtilsTest.java |    4 +-
 .../bigquery/BigQueryResourceManagerTest.java      |    2 +-
 .../bigquery/BigQueryResourceManagerUtilsTest.java |    4 +-
 .../dataflow/AbstractPipelineLauncherTest.java     |    2 +-
 .../dataflow/DefaultPipelineLauncherTest.java      |    8 +-
 .../beam/it/common/dataflow}/IOLoadTestBase.java   |   21 +-
 .../beam/it/common/dataflow}/IOStressTestBase.java |    2 +-
 .../beam/it/common/dataflow}/LoadTestBase.java     |    8 +-
 .../it/common}/storage/GcsResourceManagerTest.java |   19 +-
 it/google-cloud-platform/build.gradle              |   21 +-
 .../gcp/bigquery/conditions/BigQueryRowsCheck.java |    2 +-
 .../it/gcp/dataflow/ClassicTemplateClient.java     |    1 +
 .../beam/it/gcp/dataflow/FlexTemplateClient.java   |    1 +
 .../it/gcp/spanner/matchers/SpannerAsserts.java    |    2 +-
 .../java/org/apache/beam/it/gcp/WordCountIT.java   |    1 +
 .../apache/beam/it/gcp/bigquery/BigQueryIOLT.java  |    4 +-
 .../apache/beam/it/gcp/bigquery/BigQueryIOST.java  |    4 +-
 .../beam/it/gcp/bigquery/BigQueryStreamingLT.java  |    2 +-
 .../apache/beam/it/gcp/bigtable/BigTableIOLT.java  |    3 +-
 .../apache/beam/it/gcp/bigtable/BigTableIOST.java  |    3 +-
 .../it/gcp/dataflow/ClassicTemplateClientTest.java |    9 +-
 .../it/gcp/dataflow/FlexTemplateClientTest.java    |    8 +-
 .../beam/it/gcp/datagenerator/DataGenerator.java   |    2 +-
 .../org/apache/beam/it/gcp/pubsub/PubSubIOST.java  |    3 +-
 .../org/apache/beam/it/gcp/pubsub/PubsubIOLT.java  |    3 +-
 .../apache/beam/it/gcp/spanner/SpannerIOLT.java    |    3 +-
 .../apache/beam/it/gcp/spanner/SpannerIOST.java    |    3 +-
 .../apache/beam/it/gcp/storage/FileBasedIOLT.java  |    4 +-
 it/iceberg/build.gradle                            |  134 +
 .../org/apache/beam/it/iceberg/IcebergIOLT.java    |  363 +
 it/kafka/build.gradle                              |    3 +-
 .../java/org/apache/beam/it/kafka/KafkaIOLT.java   |    3 +-
 .../java/org/apache/beam/it/kafka/KafkaIOST.java   |    3 +-
 .../beam/it/kafka/KafkaResourceManagerTest.java    |    3 +-
 it/mongodb/build.gradle                            |    1 +
 .../beam/it/mongodb/MongoDBResourceManager.java    |   11 +-
 .../it/mongodb/MongoDBResourceManagerTest.java     |    5 +
 it/truthmatchers/build.gradle                      |    6 +-
 .../beam/it/truthmatchers}/ArtifactAsserts.java    |   15 +-
 .../beam/it/truthmatchers}/ArtifactsSubject.java   |   21 +-
 learning/tour-of-beam/backend/function.go          |    6 +-
 .../backend/integration_tests/client.go            |    8 +-
 .../backend/internal/fs_content/yaml.go            |   10 +-
 .../backend/internal/storage/datastore.go          |   11 +-
 .../tour-of-beam/backend/internal/storage/mock.go  |    8 +-
 .../beam/model/fn_execution/v1/beam_fn_api.proto   |  101 +-
 .../model/fn_execution/v1/beam_provision_api.proto |    2 +-
 .../beam/model/fnexecution/v1/standard_coders.yaml |   31 +
 .../interactive/v1/beam_interactive_api.proto      |    4 +
 .../job_management/v1/beam_artifact_api.proto      |   19 +-
 .../job_management/v1/beam_expansion_api.proto     |    9 +
 .../model/job_management/v1/beam_job_api.proto     |   32 +
 .../beam/model/pipeline/v1/beam_runner_api.proto   |   98 +-
 .../apache/beam/model/pipeline/v1/endpoints.proto  |    3 +
 .../model/pipeline/v1/external_transforms.proto    |   24 +
 .../apache/beam/model/pipeline/v1/metrics.proto    |   73 +-
 .../org/apache/beam/model/pipeline/v1/schema.proto |  107 +
 playground/backend/containers/scio/Dockerfile      |   22 +-
 playground/backend/go.mod                          |    7 +-
 playground/backend/go.sum                          |   10 +-
 .../life_cycle/life_cycle_setuper_test.go          |    3 +-
 playground/kafka-emulator/build.gradle             |   11 +
 release/build.gradle.kts                           |    4 +-
 release/src/main/groovy/TestScripts.groovy         |  156 +-
 .../main/groovy/mobilegaming-java-dataflow.groovy  |   87 +-
 .../groovy/mobilegaming-java-dataflowbom.groovy    |   12 +-
 .../main/groovy/mobilegaming-java-direct.groovy    |   73 +-
 .../main/groovy/quickstart-java-dataflow.groovy    |    8 +-
 .../main/groovy/quickstart-java-flinklocal.groovy  |    4 +-
 .../src/main/groovy/quickstart-java-spark.groovy   |   20 +-
 .../python_release_automation_utils.sh             |   30 +-
 .../run_release_candidate_python_quickstart.sh     |    6 +-
 runners/core-java/build.gradle                     |    1 +
 .../org/apache/beam/runners/core/DoFnRunner.java   |   12 +
 .../org/apache/beam/runners/core/DoFnRunners.java  |    4 +-
 .../core/GroupAlsoByWindowViaWindowSetNewDoFn.java |    2 +-
 .../apache/beam/runners/core/KeyedWorkItem.java    |    9 +
 .../runners/core/LateDataDroppingDoFnRunner.java   |   23 +-
 ...TimeBoundedSplittableProcessElementInvoker.java |  258 +-
 .../apache/beam/runners/core/ReduceFnRunner.java   |   70 +-
 .../apache/beam/runners/core/SimpleDoFnRunner.java |   83 +-
 .../core/SplittableParDoViaKeyedWorkItems.java     |   79 +-
 .../core/SplittableProcessElementInvoker.java      |   25 +-
 .../beam/runners/core/StatefulDoFnRunner.java      |    6 +
 .../org/apache/beam/runners/core/StepContext.java  |    6 +
 .../apache/beam/runners/core/WatermarkHold.java    |   17 +
 .../runners/core/construction/package-info.java    |    4 -
 .../beam/runners/core/metrics/package-info.java    |    4 -
 .../org/apache/beam/runners/core/package-info.java |    4 -
 .../core/triggers/FinishedTriggersBitSet.java      |   14 +
 .../core/triggers/TriggerStateMachineRunner.java   |   43 +-
 .../beam/runners/core/triggers/package-info.java   |    4 -
 .../core/LateDataDroppingDoFnRunnerTest.java       |    4 +-
 ...BoundedSplittableProcessElementInvokerTest.java |   47 +-
 .../beam/runners/core/ReduceFnRunnerTest.java      |   58 +
 .../apache/beam/runners/core/ReduceFnTester.java   |   33 +-
 .../SimplePushbackSideInputDoFnRunnerTest.java     |    4 +
 .../runners/core/SplittableParDoProcessFnTest.java |  405 +-
 .../triggers/TriggerStateMachineRunnerTest.java    |  117 +
 .../direct/ExecutorServiceParallelExecutor.java    |   28 +-
 .../direct/UnboundedReadEvaluatorFactory.java      |   13 +-
 .../runners/extensions/metrics/package-info.java   |    4 -
 runners/flink/1.17/build.gradle                    |   25 -
 runners/flink/1.18/build.gradle                    |   25 -
 runners/flink/1.19/build.gradle                    |    6 +
 runners/flink/1.19/job-server/build.gradle         |    6 +
 .../flink/streaming/MemoryStateBackendWrapper.java |   80 -
 .../runners/flink/streaming/StreamSources.java     |   61 -
 runners/flink/1.20/build.gradle                    |    6 +
 runners/flink/1.20/job-server/build.gradle         |    6 +
 runners/flink/2.0/build.gradle                     |    6 +
 runners/flink/2.0/job-server/build.gradle          |    6 +
 .../FlinkStreamingPortablePipelineTranslator.java  |    1 -
 .../functions/FlinkExecutableStageFunction.java    |   19 +-
 runners/flink/{2.0 => 2.1}/build.gradle            |    4 +-
 .../job-server-container/build.gradle              |    0
 .../flink/{1.17 => 2.1}/job-server/build.gradle    |    2 +-
 runners/flink/{2.0 => 2.2}/build.gradle            |   19 +-
 .../job-server-container/build.gradle              |    0
 .../flink/{1.18 => 2.2}/job-server/build.gradle    |    2 +-
 .../runners/flink/FlinkExecutionEnvironments.java  |  507 ++
 .../wrappers/streaming/DoFnOperator.java           | 1786 ++++
 .../runners/flink/FlinkPipelineOptionsTest.java    |  205 +
 runners/flink/flink_runner.gradle                  |   17 +
 runners/flink/job-server/flink_job_server.gradle   |   47 +-
 .../runners/flink/FlinkDetachedRunnerResult.java   |   60 +-
 .../beam/runners/flink/FlinkRunnerResult.java      |    6 +
 .../FlinkStreamingPortablePipelineTranslator.java  |    1 -
 .../flink/metrics/DoFnRunnerWithMetricsUpdate.java |    4 +
 .../functions/FlinkExecutableStageFunction.java    |   19 +-
 .../streaming/ExecutableStageDoFnOperator.java     |   21 +-
 .../wrappers/streaming/WindowDoFnOperator.java     |    2 +-
 .../streaming/stableinput/BufferingDoFnRunner.java |    3 +
 .../beam/runners/flink/FlinkRunnerResultTest.java  |   92 +
 .../beam/runners/flink/FlinkSubmissionTest.java    |   51 +-
 .../flink/streaming/MemoryStateBackendWrapper.java |   28 +-
 .../runners/flink/streaming/StreamSources.java     |    7 +
 .../wrappers/streaming/DoFnOperatorTest.java       |    3 +
 .../streaming/io/UnboundedSourceWrapperTest.java   |   14 +-
 .../google-cloud-dataflow-java/arm/build.gradle    |    2 +-
 runners/google-cloud-dataflow-java/build.gradle    |   22 +-
 .../beam/runners/dataflow/BatchViewOverrides.java  |   10 +-
 .../beam/runners/dataflow/DataflowPipelineJob.java |   68 +-
 .../beam/runners/dataflow/DataflowRunner.java      |   44 +-
 .../dataflow/RedistributeByKeyOverrideFactory.java |    1 +
 .../options/DataflowStreamingPipelineOptions.java  |   21 +-
 .../options/DataflowWorkerLoggingOptions.java      |   56 +-
 .../beam/runners/dataflow/util/MonitoringUtil.java |    6 +-
 .../beam/runners/dataflow/util/PackageUtil.java    |    3 +
 .../runners/dataflow/DataflowPipelineJobTest.java  |   46 +
 .../beam/runners/dataflow/DataflowRunnerTest.java  |  218 +-
 .../runners/dataflow/util/MonitoringUtilTest.java  |    4 +-
 .../runners/dataflow/util/PackageUtilTest.java     |    4 +
 .../google-cloud-dataflow-java/worker/build.gradle |   18 +
 .../worker/AssignWindowsParDoFnFactory.java        |    3 +
 .../worker/BatchModeUngroupingParDoFn.java         |    3 +
 .../CreateIsmShardKeyAndSortKeyDoFnFactory.java    |    3 +
 .../dataflow/worker/DataflowExecutionContext.java  |    4 +
 .../dataflow/worker/DataflowOutputCounter.java     |   68 +-
 .../dataflow/worker/DataflowProcessFnRunner.java   |    6 +
 .../dataflow/worker/DataflowWorkUnitClient.java    |   12 +-
 .../runners/dataflow/worker/ForwardingParDoFn.java |    6 +
 .../dataflow/worker/GroupAlsoByWindowFnRunner.java |    4 +
 .../dataflow/worker/GroupAlsoByWindowsParDoFn.java |    8 +-
 .../beam/runners/dataflow/worker/HotKeyLogger.java |   10 +-
 .../worker/IntrinsicMapTaskExecutorFactory.java    |   11 +-
 .../dataflow/worker/MultiKeyBundleOptions.java     |  138 +
 .../worker/PairWithConstantKeyDoFnFactory.java     |    3 +
 .../dataflow/worker/PartialGroupByKeyParDoFns.java |    6 +
 .../beam/runners/dataflow/worker/PubsubReader.java |   14 +-
 .../ReifyTimestampAndWindowsParDoFnFactory.java    |    3 +
 .../dataflow/worker/SimpleDoFnRunnerFactory.java   |   12 -
 .../runners/dataflow/worker/SimpleParDoFn.java     |  534 +-
 ...impleParDoFn.java => SimpleParDoFnHelpers.java} |  414 +-
 .../worker/SplittableProcessFnFactory.java         |   21 +-
 .../dataflow/worker/StreamingDataflowWorker.java   |   89 +-
 .../StreamingGroupAlsoByWindowViaWindowSetFn.java  |    2 +-
 .../StreamingKeyedWorkItemSideInputDoFnRunner.java |    6 +
 .../StreamingKeyedWorkItemSideInputParDoFn.java    |  249 +
 .../worker/StreamingModeExecutionContext.java      |  792 +-
 .../StreamingPCollectionViewWriterParDoFn.java     |    3 +
 .../worker/StreamingSideInputDoFnRunner.java       |   47 +-
 .../dataflow/worker/StreamingSideInputFetcher.java |   34 +-
 .../worker/StreamingSideInputProcessor.java        |  132 +
 .../worker/ToIsmRecordForMultimapDoFnFactory.java  |    3 +
 .../dataflow/worker/UngroupedWindmillReader.java   |   24 +-
 .../dataflow/worker/UserParDoFnFactory.java        |   56 +-
 .../runners/dataflow/worker/ValuesDoFnFactory.java |    3 +
 .../dataflow/worker/WindmillKeyedWorkItem.java     |   36 +-
 .../WindmillOpenTelemetryContextPropagator.java    |   22 +-
 .../worker/WindmillReaderIteratorBase.java         |   36 +-
 .../beam/runners/dataflow/worker/WindmillSink.java |   42 +-
 .../dataflow/worker/WindmillTimerInternals.java    |   23 +
 .../dataflow/worker/WindmillValueKindHelper.java   |   60 +
 .../dataflow/worker/WindowingWindmillReader.java   |  117 +-
 ...Exception.java => WorkCancellingException.java} |   30 +-
 .../worker/WorkItemCancelledException.java         |   29 +-
 .../WorkerCustomSourceOperationExecutor.java       |    3 +
 .../dataflow/worker/WorkerCustomSources.java       |    9 +-
 .../logging/DataflowWorkerLoggingHandler.java      |   35 +-
 .../logging/DataflowWorkerLoggingInitializer.java  |   16 +-
 .../worker/logging/DataflowWorkerLoggingMDC.java   |   15 +-
 .../dataflow/worker/streaming/ActiveWorkState.java |   51 +-
 .../streaming/BoundedQueueExecutorWorkHandle.java  |   17 +-
 .../worker/streaming/ComputationState.java         |   16 +-
 .../worker/streaming/ComputationWorkExecutor.java  |   44 +-
 .../dataflow/worker/streaming/ExecutableWork.java  |   63 +-
 .../worker/streaming/FailedWorkHandler.java        |    8 +-
 .../streaming/KeyCommitTooLargeException.java      |   20 +-
 .../runners/dataflow/worker/streaming/Work.java    |  152 +-
 .../FanOutStreamingEngineWorkerHarness.java        |    7 +-
 .../streaming/harness/MetricsDataProvider.java     |    4 +-
 .../harness/StreamingWorkerStatusReporter.java     |    2 +-
 .../dataflow/worker/util/BoundedQueueExecutor.java |  215 +-
 .../dataflow/worker/util/ExceptionUtils.java       |   36 +-
 .../dataflow/worker/util/KeyGroupWorkQueue.java    |  462 +
 .../util/common/worker/FlattenOperation.java       |    4 +
 .../worker/util/common/worker/MapTaskExecutor.java |   10 +-
 .../worker/util/common/worker/Operation.java       |    5 +
 .../worker/util/common/worker/ParDoFn.java         |    4 +
 .../worker/util/common/worker/ParDoOperation.java  |    9 +-
 .../worker/util/common/worker/ReadOperation.java   |    3 +
 .../worker/SimplePartialGroupByKeyParDoFn.java     |    5 +
 .../worker/util/common/worker/WorkExecutor.java    |    3 +
 .../worker/util/common/worker/WriteOperation.java  |    4 +
 .../worker/windmill/client/WindmillStream.java     |    5 +
 .../worker/windmill/client/commits/Commit.java     |   76 +-
 .../worker/windmill/client/commits/Commits.java    |    2 +-
 .../windmill/client/commits/CompleteCommit.java    |   30 +-
 .../commits/StreamingApplianceWorkCommitter.java   |   13 +-
 .../commits/StreamingEngineWorkCommitter.java      |  113 +-
 .../client/getdata/StreamGetDataClient.java        |    5 +-
 .../windmill/client/grpc/GrpcCommitWorkStream.java |  139 +-
 .../client/grpc/GrpcWindmillStreamFactory.java     |   16 +-
 .../windmill/state/WindmillStateInternals.java     |   17 +
 .../worker/windmill/state/WindmillStateReader.java |    9 +-
 .../windmill/state/WindmillWatermarkHold.java      |   94 +-
 .../processing/ComputationWorkExecutorFactory.java |   52 +-
 .../work/processing/StreamingCommitFinalizer.java  |    2 +-
 .../work/processing/StreamingWorkScheduler.java    |  507 +-
 .../processing/failures/WorkFailureProcessor.java  |  100 +-
 .../windmill/work/refresh/ActiveWorkRefresher.java |   21 +-
 .../worker/DataflowOperationContextTest.java       |    6 +-
 .../dataflow/worker/DataflowOutputCounterTest.java |  108 +
 .../worker/DataflowWorkUnitClientTest.java         |    8 +-
 .../dataflow/worker/FakeWindmillServer.java        |   67 +-
 .../IntrinsicMapTaskExecutorFactoryTest.java       |   14 +-
 .../worker/IntrinsicMapTaskExecutorTest.java       |   15 +
 .../worker/KeyTokenInvalidExceptionTest.java       |   39 -
 .../dataflow/worker/SimpleParDoFnHelpersTest.java  |  133 +
 .../runners/dataflow/worker/SimpleParDoFnTest.java |    6 +-
 .../worker/StreamingDataflowWorkerTest.java        |  627 +-
 .../worker/StreamingGroupAlsoByWindowFnsTest.java  |    9 +-
 ...eamingKeyedWorkItemSideInputDoFnRunnerTest.java |    3 +-
 ...StreamingKeyedWorkItemSideInputParDoFnTest.java |  488 ++
 .../worker/StreamingModeExecutionContextTest.java  |  591 +-
 .../worker/StreamingSideInputDoFnRunnerTest.java   |    4 +
 .../worker/StreamingSideInputProcessorTest.java    |  215 +
 .../dataflow/worker/UserParDoFnFactoryTest.java    |   27 +-
 .../worker/WindmillReaderIteratorBaseTest.java     |  154 +-
 .../worker/WindowingWindmillReaderTest.java        |  276 +
 .../dataflow/worker/WorkerCustomSourcesTest.java   |   88 +-
 .../logging/DataflowWorkerLoggingHandlerTest.java  |  107 +-
 .../DataflowWorkerLoggingInitializerTest.java      |   44 +-
 .../dataflow/worker/status/ThreadzServletTest.java |   15 +-
 .../worker/streaming/ActiveWorkStateTest.java      |   71 +-
 .../streaming/ComputationStateCacheTest.java       |    6 +-
 .../worker/streaming/ComputationStateTest.java     |  114 +
 .../dataflow/worker/streaming/WorkTest.java        |  134 +
 .../worker/testing/RestoreDataflowLoggingMDC.java  |    8 +-
 .../testing/RestoreDataflowLoggingMDCTest.java     |   10 +-
 .../worker/util/BoundedQueueExecutorTest.java      |  402 +-
 .../worker/util/KeyGroupWorkQueueTest.java         |  502 ++
 .../util/common/worker/ExecutorTestUtils.java      |    3 +
 .../util/common/worker/MapTaskExecutorTest.java    |   15 +
 .../util/common/worker/ParDoOperationTest.java     |    6 +
 .../StreamingApplianceWorkCommitterTest.java       |   15 +-
 .../commits/StreamingEngineWorkCommitterTest.java  |  323 +-
 .../client/grpc/GrpcCommitWorkStreamTest.java      |  355 +
 .../windmill/state/WindmillStateInternalsTest.java |  167 +-
 .../windmill/state/WindmillStateReaderTest.java    |   10 +-
 .../processing/StreamingCommitFinalizerTest.java   |    3 +-
 .../failures/WorkFailureProcessorTest.java         |  137 +-
 .../work/refresh/ActiveWorkRefresherTest.java      |   80 +-
 .../worker/windmill/src/main/proto/windmill.proto  |   32 +-
 .../control/BundleCheckpointHandlers.java          |   24 +-
 .../fnexecution/control/SdkHarnessClient.java      |    2 +-
 .../runners/fnexecution/data/FnDataService.java    |    3 +-
 .../runners/fnexecution/data/GrpcDataService.java  |   13 +-
 .../runners/fnexecution/ServerFactoryTest.java     |   19 +-
 .../control/BundleCheckpointHandlersTest.java      |  125 +
 .../fnexecution/data/GrpcDataServiceTest.java      |    2 +-
 .../beam/runners/jobsubmission/JobInvocation.java  |   15 +-
 .../runners/jobsubmission/JobInvocationTest.java   |   36 +
 .../portability/JobServicePipelineResult.java      |    7 +-
 .../runners/portability/PortableRunnerTest.java    |   21 +
 runners/portability/test_pipeline_jar.sh           |   11 +-
 runners/prism/java/build.gradle                    |    4 +
 runners/spark/3/build.gradle                       |    6 +
 runners/spark/3/job-server/build.gradle            |    8 +-
 runners/spark/4/build.gradle                       |    5 +
 runners/spark/4/job-server/build.gradle            |    6 +
 runners/spark/job-server/spark_job_server.gradle   |   31 +-
 runners/spark/spark_runner.gradle                  |   12 +-
 .../translation/batch/DoFnRunnerFactory.java       |    4 +
 .../translation/batch/DoFnRunnerWithMetrics.java   |    4 +
 .../spark/translation/DoFnRunnerWithMetrics.java   |    4 +
 .../SparkBatchPortablePipelineTranslator.java      |    8 +-
 .../translation/SparkExecutableStageFunction.java  |  172 +-
 .../SparkStreamingPortablePipelineTranslator.java  |    4 +-
 .../SparkExecutableStageFunctionTest.java          |  138 +-
 .../translation/SparkInputDataProcessorTest.java   |    4 +
 .../streaming/utils/EmbeddedKafkaCluster.java      |   17 +-
 scripts/beam-sql.sh                                |   75 +-
 scripts/ci/issue-report/package-lock.json          |   14 +-
 scripts/ci/issue-report/package.json               |    2 +-
 scripts/ci/pr-bot/processNewPrs.ts                 |   20 +
 scripts/ci/pr-bot/shared/githubUtils.ts            |   21 +
 scripts/tools/bomupgrader.py                       |  122 +-
 sdks/go.mod                                        |  186 +-
 sdks/go.sum                                        |  449 +-
 sdks/go/README.md                                  |    2 +-
 sdks/go/container/boot.go                          |   40 +-
 sdks/go/container/boot_test.go                     |   17 +-
 sdks/go/container/tools/buffered_logging.go        |   64 +-
 sdks/go/container/tools/buffered_logging_test.go   |  168 +-
 sdks/go/container/tools/pipeline_options.go        |  218 +
 sdks/go/container/tools/pipeline_options_test.go   |  243 +
 sdks/go/examples/wasm/README.md                    |    6 +-
 sdks/go/pkg/beam/artifact/options.go               |   48 -
 sdks/go/pkg/beam/artifact/options_test.go          |   78 -
 sdks/go/pkg/beam/coder.go                          |   17 +
 sdks/go/pkg/beam/core/core.go                      |    2 +-
 sdks/go/pkg/beam/core/graph/coder/coder.go         |  116 +
 sdks/go/pkg/beam/core/graph/coder/coder_test.go    |   66 +
 sdks/go/pkg/beam/core/graph/coder/registry.go      |   60 +-
 .../pkg/beam/core/graph/coder/sharded_key_test.go  |   81 +
 sdks/go/pkg/beam/core/runtime/exec/coder.go        |   63 +
 sdks/go/pkg/beam/core/runtime/exec/coder_test.go   |   84 +
 .../pkg/beam/core/runtime/exec/datasampler_test.go |   13 +-
 sdks/go/pkg/beam/core/runtime/graphx/coder.go      |   21 +
 sdks/go/pkg/beam/core/runtime/symbols.go           |   20 +
 sdks/go/pkg/beam/core/typex/class.go               |    4 +-
 sdks/go/pkg/beam/core/typex/fulltype.go            |   23 +
 sdks/go/pkg/beam/core/typex/special.go             |   18 +-
 sdks/go/pkg/beam/core/util/reflectx/call.go        |   31 +
 sdks/go/pkg/beam/io/filesystem/gcs/gcs.go          |   88 +-
 sdks/go/pkg/beam/io/filesystem/gcs/gcs_test.go     |  146 +
 sdks/go/pkg/beam/pcollection.go                    |   16 +
 sdks/go/pkg/beam/runners/dataflow/dataflow.go      |    9 +-
 sdks/go/pkg/beam/runners/dataflow/dataflow_test.go |   43 +-
 .../beam/runners/dataflow/dataflowlib/execute.go   |    1 +
 .../pkg/beam/runners/dataflow/dataflowlib/job.go   |    4 +
 .../beam/runners/dataflow/dataflowlib/job_test.go  |   16 +
 sdks/go/pkg/beam/runners/prism/internal/coders.go  |   11 +
 .../pkg/beam/runners/prism/internal/coders_test.go |   16 +
 .../prism/internal/engine/elementmanager.go        |    8 +-
 .../beam/runners/prism/internal/environments.go    |   93 +-
 .../beam/runners/prism/internal/execute_test.go    |   21 +
 .../beam/runners/prism/internal/handlerunner.go    |    4 +-
 .../runners/prism/internal/handlerunner_test.go    |   78 +
 sdks/go/pkg/beam/runners/prism/internal/stage.go   |   17 +-
 .../beam/runners/prism/internal/testdofns_test.go  |   34 +
 .../beam/runners/prism/internal/worker/bundle.go   |   58 +-
 .../runners/prism/internal/worker/worker_test.go   |   67 +
 sdks/go/pkg/beam/transforms/batch/batch.go         |  677 ++
 .../pkg/beam/transforms/batch/batch_prism_test.go  |  222 +
 .../beam/transforms/batch/batch_test.go}           |   53 +-
 sdks/go/pkg/beam/transforms/batch/doc.go           |   58 +
 sdks/go/pkg/beam/transforms/batch/size.go          |   88 +
 sdks/go/pkg/beam/transforms/batch/size_test.go     |   91 +
 sdks/go/test/build.gradle                          |    3 +
 .../test/integration/io/xlang/debezium/debezium.go |    2 +-
 .../integration/io/xlang/debezium/debezium_test.go |    2 +-
 .../go/test/integration/io/xlang/jdbc/jdbc_test.go |    6 +-
 sdks/go/test/run_validatesrunner_tests.sh          |    4 +
 sdks/java/container/boot.go                        |   19 +-
 sdks/java/container/common.gradle                  |    1 +
 sdks/java/container/java17/option-arrow.json       |    9 +
 sdks/java/container/java21/option-arrow.json       |    9 +
 sdks/java/container/java25/option-arrow.json       |    9 +
 .../container/license_scripts/dep_urls_java.yaml   |   12 +-
 .../org/apache/beam/sdk/jmh/util/package-info.java |    4 -
 .../java/org/apache/beam/sdk/PipelineResult.java   |   19 +
 .../apache/beam/sdk/annotations/package-info.java  |    4 -
 .../org/apache/beam/sdk/coders/package-info.java   |    4 -
 .../apache/beam/sdk/expansion/package-info.java    |    4 -
 .../sdk/fn/data/BeamFnDataGrpcMultiplexer.java     |    2 +-
 .../sdk/fn/data/BeamFnDataOutboundAggregator.java  |  131 +-
 .../beam/sdk/fn/splittabledofn/package-info.java   |    4 -
 .../org/apache/beam/sdk/harness/package-info.java  |    3 -
 .../org/apache/beam/sdk/io/CountingSource.java     |   40 +-
 .../org/apache/beam/sdk/io/fs/package-info.java    |    4 -
 .../java/org/apache/beam/sdk/io/package-info.java  |    4 -
 .../org/apache/beam/sdk/io/range/package-info.java |    4 -
 .../org/apache/beam/sdk/metrics/package-info.java  |    4 -
 .../beam/sdk/options/PipelineOptionsFactory.java   |   26 +-
 .../apache/beam/sdk/options/SdkHarnessOptions.java |    7 +
 .../java/org/apache/beam/sdk/package-info.java     |    4 -
 .../org/apache/beam/sdk/runners/package-info.java  |    3 -
 .../sdk/schemas/FieldValueTypeInformation.java     |   13 +-
 .../beam/sdk/schemas/FromRowUsingCreator.java      |    6 +-
 .../sdk/schemas/GetterBasedSchemaProvider.java     |   98 +-
 .../apache/beam/sdk/schemas/SchemaTranslation.java |    2 +
 .../org/apache/beam/sdk/schemas/SchemaUtils.java   |    7 +
 .../beam/sdk/schemas/annotations/package-info.java |    4 -
 .../apache/beam/sdk/schemas/io/package-info.java   |    4 -
 .../beam/sdk/schemas/io/payloads/package-info.java |    4 -
 .../beam/sdk/schemas/logicaltypes/Timestamp.java   |    9 +-
 .../sdk/schemas/logicaltypes/package-info.java     |    4 -
 .../org/apache/beam/sdk/schemas/package-info.java  |    4 -
 .../sdk/schemas/parser/generated/package-info.java |    4 -
 .../beam/sdk/schemas/parser/package-info.java      |    4 -
 .../beam/sdk/schemas/transforms/package-info.java  |    4 -
 .../schemas/transforms/providers/package-info.java |    4 -
 .../beam/sdk/schemas/utils/AutoValueUtils.java     |    9 +-
 .../beam/sdk/schemas/utils/ByteBuddyUtils.java     |  128 +
 .../sdk/schemas/utils/StaticSchemaInference.java   |   19 +-
 .../beam/sdk/schemas/utils/package-info.java       |    4 -
 .../apache/beam/sdk/state/WatermarkHoldState.java  |   10 +
 .../org/apache/beam/sdk/state/package-info.java    |    4 -
 .../UsesSideInputsInTimer.java}                    |   15 +-
 .../org/apache/beam/sdk/testing/package-info.java  |    4 -
 .../apache/beam/sdk/transforms/AsyncWrapper.java   |  781 ++
 .../apache/beam/sdk/transforms/BatchElements.java  |    9 +-
 .../java/org/apache/beam/sdk/transforms/DoFn.java  |   34 +
 .../org/apache/beam/sdk/transforms/DoFnTester.java |   25 +-
 .../org/apache/beam/sdk/transforms/JsonToRow.java  |    8 +-
 .../apache/beam/sdk/transforms/LogElements.java    |  234 +
 .../apache/beam/sdk/transforms/Redistribute.java   |   22 +-
 .../java/org/apache/beam/sdk/transforms/Reify.java |    7 +-
 .../org/apache/beam/sdk/transforms/Reshuffle.java  |    1 +
 .../org/apache/beam/sdk/transforms/WithKeys.java   |   27 +-
 .../beam/sdk/transforms/display/DisplayData.java   |    7 +-
 .../beam/sdk/transforms/display/package-info.java  |    4 -
 .../sdk/transforms/errorhandling/package-info.java |    4 -
 .../beam/sdk/transforms/join/package-info.java     |    4 -
 .../apache/beam/sdk/transforms/package-info.java   |    4 -
 .../reflect/ByteBuddyDoFnInvokerFactory.java       |  106 +-
 .../beam/sdk/transforms/reflect/DoFnInvoker.java   |   32 +
 .../beam/sdk/transforms/reflect/DoFnSignature.java |   35 +
 .../sdk/transforms/reflect/DoFnSignatures.java     |   70 +-
 .../beam/sdk/transforms/reflect/package-info.java  |    3 -
 .../transforms/splittabledofn/package-info.java    |    4 -
 .../sdk/transforms/windowing/package-info.java     |    4 -
 .../java/org/apache/beam/sdk/util/ApiSurface.java  |    7 +-
 .../java/org/apache/beam/sdk/util/RowJson.java     |   65 +-
 .../beam/sdk/util/RowJsonValueExtractors.java      |    9 +-
 .../beam/sdk/util/common/ReflectHelpers.java       |   22 +-
 .../beam/sdk/util/construction/Environments.java   |    1 +
 .../construction/SplittableParDoNaiveBounded.java  |   17 +
 .../sdk/util/construction/graph/package-info.java  |    4 -
 .../beam/sdk/util/construction/package-info.java   |    4 -
 .../sdk/values/OpenTelemetryContextPropagator.java |    8 +-
 .../main/java/org/apache/beam/sdk/values/Row.java  |    6 +-
 .../org/apache/beam/sdk/values/RowWithGetters.java |   13 +-
 .../org/apache/beam/sdk/values/WindowedValues.java |   11 +-
 .../org/apache/beam/sdk/values/package-info.java   |    4 -
 .../fn/data/BeamFnDataOutboundAggregatorTest.java  |  135 +-
 .../org/apache/beam/sdk/io/CountingSourceTest.java |   36 +
 .../java/org/apache/beam/sdk/io/FileIOTest.java    |   76 +-
 .../sdk/options/PipelineOptionsFactoryTest.java    |   56 +-
 .../beam/sdk/schemas/AutoValueSchemaTest.java      |  426 +-
 .../beam/sdk/schemas/JavaBeanSchemaTest.java       |  250 +-
 .../beam/sdk/schemas/JavaFieldSchemaTest.java      |  318 +-
 .../beam/sdk/schemas/SchemaTranslationTest.java    |    7 +
 .../apache/beam/sdk/schemas/SchemaUtilsTest.java   |   43 +
 .../beam/sdk/schemas/transforms/ConvertTest.java   |   99 +
 .../providers/JavaFilterTransformProviderTest.java |    8 +-
 .../beam/sdk/schemas/utils/JavaBeanUtilsTest.java  |   10 +
 .../beam/sdk/schemas/utils/POJOUtilsTest.java      |   32 +
 .../beam/sdk/schemas/utils/TestJavaBeans.java      |  124 +
 .../apache/beam/sdk/schemas/utils/TestPOJOs.java   |  143 +
 .../beam/sdk/transforms/AsyncWrapperTest.java      |  941 ++
 .../apache/beam/sdk/transforms/JsonToRowTest.java  |   26 +
 .../beam/sdk/transforms/LogElementsTest.java       |  108 +
 .../org/apache/beam/sdk/transforms/ParDoTest.java  |  153 +
 .../apache/beam/sdk/transforms/ValueKindTest.java  |  623 ++
 .../sdk/transforms/reflect/DoFnInvokersTest.java   |   47 +
 .../sdk/transforms/reflect/DoFnSignaturesTest.java |  125 +-
 .../java/org/apache/beam/sdk/util/RowJsonTest.java |   26 +-
 .../apache/beam/sdk/util/WindowedValueTest.java    |   33 +
 .../beam/sdk/util/common/ReflectHelpersTest.java   |   15 +-
 sdks/java/expansion-service/container/Dockerfile   |    2 +-
 .../service/WindowIntoTransformProvider.java       |    1 +
 sdks/java/extensions/arrow/build.gradle            |    7 +-
 .../beam/sdk/extensions/arrow/ArrowConversion.java |   65 +-
 .../sdk/extensions/arrow/ArrowConversionTest.java  |   33 +
 sdks/java/extensions/avro/build.gradle             |    7 +
 .../sdk/extensions/avro/coders/package-info.java   |    4 -
 .../beam/sdk/extensions/avro/io/package-info.java  |    4 -
 .../beam/sdk/extensions/avro/package-info.java     |    4 -
 .../avro/schemas/io/payloads/package-info.java     |    4 -
 .../sdk/extensions/avro/schemas/package-info.java  |    4 -
 .../extensions/avro/schemas/utils/AvroUtils.java   |   89 +-
 .../avro/schemas/utils/package-info.java           |    4 -
 .../avro/schemas/utils/AvroUtilsTest.java          |    4 +-
 .../google-cloud-platform-core/build.gradle        |    1 +
 .../sdk/extensions/gcp/options/GcsOptions.java     |   85 +-
 .../beam/sdk/extensions/gcp/util/GcsUtilV1.java    |   26 +-
 .../sdk/extensions/gcp/GcpCoreApiSurfaceTest.java  |    2 +
 .../sdk/extensions/gcp/options/GcsOptionsTest.java |   47 +
 .../beam/sdk/extensions/gcp/util/GcsUtilTest.java  |   13 +
 .../opentelemetry-gcp-auth-extension/build.gradle  |   51 +
 .../opentelemetry/gcp/auth/ConfigurableOption.java |  164 +
 ...GcpAuthAutoConfigurationCustomizerProvider.java |  292 +
 .../gcp/auth/GoogleAuthException.java              |   69 +
 .../opentelemetry/gcp/auth}/package-info.java      |    4 +-
 ...uthAutoConfigurationCustomizerProviderTest.java | 1310 +++
 .../extensions/protobuf/ProtoByteBuddyUtils.java   |   49 +-
 sdks/java/extensions/sql/iceberg/build.gradle      |    4 +-
 .../meta/provider/iceberg/IcebergMetastore.java    |   21 +-
 .../sql/meta/provider/iceberg/IcebergTable.java    |   19 +-
 .../provider/iceberg/BeamSqlCliIcebergTest.java    |    6 +-
 .../provider/iceberg/IcebergMetastoreTest.java     |   13 +
 .../meta/provider/iceberg/IcebergReadWriteIT.java  |    3 -
 .../meta/provider/iceberg/PubsubToIcebergIT.java   |    4 +-
 .../sql/src/main/codegen/includes/parserImpls.ftl  |    4 +-
 .../beam/sdk/extensions/sql/impl/BeamSqlEnv.java   |   13 +
 .../extensions/sql/impl/CalciteQueryPlanner.java   |  103 +-
 .../sdk/extensions/sql/impl/JdbcConnection.java    |   22 +
 .../beam/sdk/extensions/sql/impl/UdfImpl.java      |   19 +-
 .../extensions/sql/impl/parser/SqlDdlNodes.java    |    7 +-
 .../sdk/extensions/sql/impl/rel/BeamCalcRel.java   |   40 +
 .../sdk/extensions/sql/impl/rel/BeamSortRel.java   |   58 +-
 .../sdk/extensions/sql/impl/rel/package-info.java  |    4 -
 .../sdk/extensions/sql/impl/rule/package-info.java |    4 -
 .../impl/transform/BeamBuiltinAggregations.java    |    2 +
 .../sql/impl/transform/agg/VarianceFn.java         |   38 +-
 .../sql/impl/transform/agg/package-info.java       |    4 -
 .../extensions/sql/impl/utils/CalciteUtils.java    |   39 +-
 .../sql/meta/catalog/InMemoryCatalog.java          |   12 +-
 .../sql/meta/provider/mongodb/package-info.java    |    4 -
 .../sql/meta/provider/pubsub/package-info.java     |    4 -
 .../sdk/extensions/sql/BeamComplexTypeTest.java    |   32 +
 .../{BeamSqlAliasTest => BeamSqlAliasTest.java}    |   12 +-
 .../sdk/extensions/sql/BeamSqlCliDatabaseTest.java |   52 +
 .../sql/BeamSqlDslAggregationVarianceTest.java     |   43 +-
 .../extensions/sql/BeamSqlDslParametersTest.java   |   78 +
 .../sql/impl/BeamSqlEnvRegisterOperatorTest.java   |  100 +
 .../sql/impl/LazyAggregateCombineFnTest.java       |    5 +-
 .../beam/sdk/extensions/sql/impl/UdfImplTest.java  |   65 +
 .../extensions/sql/impl/rel/BeamCalcRelTest.java   |   20 +
 .../extensions/sql/impl/rel/BeamSortRelTest.java   |   43 +
 .../sql/impl/transform/agg/VarianceFnTest.java     |   25 +-
 .../sql/impl/utils/CalciteUtilsTest.java           |   63 +
 .../apache/beam/fn/harness/FnApiDoFnRunner.java    |   72 +-
 .../java/org/apache/beam/fn/harness/FnHarness.java |   66 +-
 .../fn/harness/control/ProcessBundleHandler.java   |   51 +-
 .../beam/fn/harness/data/BeamFnDataClient.java     |   37 +-
 .../beam/fn/harness/data/BeamFnDataGrpcClient.java |  103 +-
 .../beam/fn/harness/BeamFnDataWriteRunnerTest.java |   44 +-
 .../beam/fn/harness/FnApiDoFnRunnerTest.java       |    2 +-
 .../PTransformRunnerFactoryTestContext.java        |   12 +-
 .../harness/control/ProcessBundleHandlerTest.java  |   47 +-
 .../fn/harness/data/BeamFnDataGrpcClientTest.java  |   41 +-
 .../sdk/io/aws2/schemas/AwsSchemaProvider.java     |    7 +-
 .../apache/beam/sdk/io/aws2/schemas/AwsTypes.java  |   16 +
 .../sdk/io/aws2/sqs/SqsIOWriteBatchesTest.java     |   97 +-
 .../remote => io/arrow-flight}/build.gradle        |   39 +-
 .../beam/sdk/io/arrowflight/ArrowFlightIO.java     |  840 ++
 .../beam/sdk/io/arrowflight}/package-info.java     |   16 +-
 .../beam/sdk/io/arrowflight/ArrowFlightIOTest.java |  330 +
 sdks/java/io/cassandra/build.gradle                |   21 +
 sdks/java/io/cdap/build.gradle                     |   10 +
 .../beam/sdk/io/clickhouse/ClickHouseIO.java       |    3 +
 .../beam/sdk/io/clickhouse/ClickHouseWriter.java   |   40 +
 .../apache/beam/sdk/io/clickhouse/TableSchema.java |   53 +-
 .../clickhouse/src/main/javacc/ColumnTypeParser.jj |   27 +
 .../beam/sdk/io/clickhouse/ClickHouseIOIT.java     |  119 +
 .../sdk/io/clickhouse/ClickHouseWriterTest.java    |  141 +
 .../beam/sdk/io/clickhouse/TableSchemaTest.java    |  102 +
 .../components/throttling/AdaptiveThrottler.java   |  105 +
 .../components/throttling/ReactiveThrottler.java   |   78 +
 .../throttling/AdaptiveThrottlerTest.java          |  107 +
 .../org/apache/beam/sdk/io/csv/CsvIOParseTest.java |    6 +-
 .../apache/beam/sdk/io/datadog/DatadogEvent.java   |    6 +
 .../DatadogWriteSchemaTransformConfiguration.java  |  114 +
 .../DatadogWriteSchemaTransformProvider.java       |  294 +
 .../DatadogWriteSchemaTransformProviderTest.java   |  535 ++
 sdks/java/io/debezium/build.gradle                 |   17 +-
 .../io/debezium/expansion-service/build.gradle     |    6 +-
 sdks/java/io/debezium/src/README.md                |    8 +-
 .../org/apache/beam/io/debezium/DebeziumIO.java    |   15 +-
 .../DebeziumReadSchemaTransformProvider.java       |   96 +-
 .../beam/io/debezium/KafkaSourceConsumerFn.java    |    5 +
 .../io/debezium/DebeziumIOMySqlConnectorIT.java    |    6 +-
 .../debezium/DebeziumIOPostgresSqlConnectorIT.java |    4 +-
 .../apache/beam/io/debezium/DebeziumIOTest.java    |    3 +-
 .../DebeziumReadSchemaTransformProviderTest.java   |  139 +
 .../debezium/DebeziumReadSchemaTransformTest.java  |   29 +-
 sdks/java/io/delta/build.gradle                    |  141 +
 .../org/apache/beam/sdk/io/delta/BeamEngine.java   |   55 +
 .../beam/sdk/io/delta/BeamParquetHandler.java      |  376 +
 .../beam/sdk/io/delta/CreateCDCReadTasksDoFn.java  |  296 +
 .../beam/sdk/io/delta/CreateReadTasksDoFn.java     |  140 +
 .../apache/beam/sdk/io/delta/DeltaCDCReadTask.java |  125 +
 .../beam/sdk/io/delta/DeltaCDCSourceDoFn.java      |  367 +
 .../delta/DeltaCdcReadSchemaTransformProvider.java |  178 +
 .../java/org/apache/beam/sdk/io/delta/DeltaIO.java |  374 +
 .../delta/DeltaReadSchemaTransformProvider.java}   |   94 +-
 .../apache/beam/sdk/io/delta/DeltaReadTask.java    |   88 +
 .../beam/sdk/io/delta/DeltaReadTaskTracker.java    |   56 +
 .../apache/beam/sdk/io/delta/DeltaSourceDoFn.java  |  455 +
 .../apache/beam/sdk/io/delta/SerializableRow.java  |  544 ++
 .../beam/sdk/io/delta/SerializableStructType.java  |   69 +
 .../apache/beam/sdk/io/delta}/package-info.java    |    8 +-
 .../org/apache/beam/sdk/io/delta/DeltaIOIT.java    |  428 +
 .../org/apache/beam/sdk/io/delta/DeltaIOS3IT.java  |  278 +
 .../org/apache/beam/sdk/io/delta/DeltaIOTest.java  | 1718 ++++
 .../DeltaReadSchemaTransformProviderTest.java      |  127 +
 .../beam/sdk/io/delta/DeltaWriteTestUtils.java     |  371 +
 .../beam/sdk/io/delta/SerializableRowTest.java     |  498 ++
 sdks/java/io/expansion-service/build.gradle        |    9 +-
 sdks/java/io/google-cloud-platform/build.gradle    |   12 +-
 .../beam/sdk/io/gcp/bigquery/BigQueryIO.java       |   17 +-
 .../beam/sdk/io/gcp/bigquery/BigQueryUtils.java    |    6 +
 .../bigquery/StorageApiWriteUnshardedRecords.java  |   12 +-
 .../bigquery/StorageApiWritesShardedRecords.java   |   20 +-
 ...ueryStorageWriteApiSchemaTransformProvider.java |    9 +-
 .../providers/BigQueryWriteConfiguration.java      |    9 +-
 .../sdk/io/gcp/datastore/AdaptiveThrottler.java    |  111 -
 .../beam/sdk/io/gcp/datastore/DatastoreV1.java     |    1 +
 .../apache/beam/sdk/io/gcp/healthcare/FhirIO.java  |    6 +-
 .../apache/beam/sdk/io/gcp/healthcare/HL7v2IO.java |   17 +-
 .../sdk/io/gcp/pubsub/AddTimestampAttribute.java   |   13 +-
 .../beam/sdk/io/gcp/pubsub/ExternalWrite.java      |   21 +-
 .../beam/sdk/io/gcp/pubsub/NestedRowToMessage.java |    8 +-
 .../io/gcp/pubsub/PubSubPayloadTranslation.java    |   46 +-
 .../beam/sdk/io/gcp/pubsub/PubsubClient.java       |   59 +-
 .../beam/sdk/io/gcp/pubsub/PubsubGrpcClient.java   |   47 +-
 .../apache/beam/sdk/io/gcp/pubsub/PubsubIO.java    |  306 +-
 .../beam/sdk/io/gcp/pubsub/PubsubJsonClient.java   |  100 +-
 .../beam/sdk/io/gcp/pubsub/PubsubMessage.java      |   23 +-
 .../beam/sdk/io/gcp/pubsub/PubsubMessageToRow.java |   42 +-
 ...hAttributesAndMessageIdAndOrderingKeyCoder.java |   19 +-
 ...bsubMessageWithAttributesAndMessageIdCoder.java |   14 +-
 .../pubsub/PubsubMessageWithAttributesCoder.java   |    9 +-
 .../pubsub/PubsubMessageWithMessageIdCoder.java    |    9 +-
 .../pubsub/PubsubReadSchemaTransformProvider.java  |   34 +-
 .../beam/sdk/io/gcp/pubsub/PubsubRowToMessage.java |   58 +-
 .../sdk/io/gcp/pubsub/PubsubSchemaIOProvider.java  |   51 +-
 .../beam/sdk/io/gcp/pubsub/PubsubTestClient.java   |  118 +-
 .../sdk/io/gcp/pubsub/PubsubUnboundedSink.java     |   62 +-
 .../sdk/io/gcp/pubsub/PubsubUnboundedSource.java   |  190 +-
 .../pubsub/PubsubWriteSchemaTransformProvider.java |   59 +-
 .../apache/beam/sdk/io/gcp/pubsub/TestPubsub.java  |   92 +-
 .../beam/sdk/io/gcp/pubsub/TestPubsubSignal.java   |   79 +-
 .../beam/sdk/io/gcp/spanner/BatchSpannerRead.java  |   13 +-
 .../sdk/io/gcp/spanner/CreateTransactionFn.java    |    8 +-
 .../beam/sdk/io/gcp/spanner/NaiveSpannerRead.java  |    8 +-
 .../beam/sdk/io/gcp/spanner/ReadSpannerSchema.java |    8 +-
 .../beam/sdk/io/gcp/spanner/SpannerAccessor.java   |   50 +-
 .../beam/sdk/io/gcp/spanner/SpannerConfig.java     |   93 +
 .../apache/beam/sdk/io/gcp/spanner/SpannerIO.java  |  280 +-
 .../io/gcp/spanner/SpannerTransformRegistrar.java  |   30 +
 .../MetadataSpannerConfigFactory.java              |    6 +
 .../changestreams/action/ActionFactory.java        |    7 +-
 .../action/HeartbeatRecordAction.java              |    2 +-
 .../action/QueryChangeStreamAction.java            |   44 +-
 .../gcp/spanner/changestreams/dao/DaoFactory.java  |   16 +-
 .../dofn/CleanUpReadChangeStreamDoFn.java          |    7 +
 .../dofn/DetectNewPartitionsDoFn.java              |    5 +-
 .../spanner/changestreams/dofn/InitializeDoFn.java |    7 +
 .../dofn/ReadChangeStreamPartitionDoFn.java        |    9 +-
 .../sdk/io/gcp/testing/FakeDatasetService.java     |   18 +
 .../apache/beam/sdk/io/gcp/GcpApiSurfaceTest.java  |    1 +
 .../gcp/bigquery/BeamRowToStorageApiProtoTest.java |   11 +-
 .../sdk/io/gcp/bigquery/BigQueryIOWriteTest.java   |   63 +
 .../sdk/io/gcp/bigquery/BigQueryUtilsTest.java     |   42 +-
 ....java => StorageApiSinkSchemaUpdateITBase.java} |   61 +-
 ...torageApiSinkSchemaUpdateWithInputSchemaIT.java |   50 +
 ...ageApiSinkSchemaUpdateWithoutInputSchemaIT.java |   51 +
 ...StorageWriteApiSchemaTransformProviderTest.java |    5 +-
 .../io/gcp/datastore/AdaptiveThrottlerTest.java    |  114 -
 .../sdk/io/gcp/spanner/SpannerAccessorTest.java    |   51 +
 .../gcp/spanner/SpannerIOReadChangeStreamTest.java |   16 +
 .../sdk/io/gcp/spanner/SpannerIOWriteTest.java     |   14 +-
 .../beam/sdk/io/gcp/spanner/SpannerReadIT.java     |   33 +-
 .../beam/sdk/io/gcp/spanner/SpannerTestHelper.java |  110 +
 .../beam/sdk/io/gcp/spanner/SpannerWriteIT.java    |  110 +-
 .../SpannerChangeStreamErrorTest.java              |   41 +-
 .../action/HeartbeatRecordActionTest.java          |    4 +-
 .../action/QueryChangeStreamActionTest.java        |   10 +-
 .../dofn/ReadChangeStreamPartitionDoFnTest.java    |    6 +-
 .../changestreams/it/IntegrationTestEnv.java       |   17 +-
 .../changestreams/it/SpannerChangeStreamIT.java    |   19 +-
 ...StreamOrderedByTimestampAndTransactionIdIT.java |   10 +-
 ...nnerChangeStreamOrderedWithinKeyGloballyIT.java |   10 +-
 .../it/SpannerChangeStreamOrderedWithinKeyIT.java  |   10 +-
 .../it/SpannerChangeStreamPlacementTableIT.java    |   19 +-
 ...pannerChangeStreamPlacementTablePostgresIT.java |   10 +-
 .../it/SpannerChangeStreamPostgresIT.java          |   10 +-
 ...SpannerChangeStreamTransactionBoundariesIT.java |   10 +-
 .../it/SpannerChangeStreamsSchemaTransformIT.java  |    4 +
 .../apache/beam/sdk/io/hdfs/HadoopFileSystem.java  |   17 +
 .../beam/sdk/io/hdfs/HadoopFileSystemTest.java     |   20 +
 sdks/java/io/hadoop-format/build.gradle            |   22 +
 .../beam/sdk/io/hadoop/format/HadoopFormatIO.java  |   25 +
 sdks/java/io/hcatalog/build.gradle                 |    4 +
 sdks/java/io/iceberg/build.gradle                  |   10 +-
 .../org/apache/beam/sdk/io/iceberg/AddFiles.java   |   49 +-
 .../beam/sdk/io/iceberg/AppendFilesToTables.java   |    5 +-
 .../iceberg/AssignDestinationsAndPartitions.java   |   53 +-
 .../beam/sdk/io/iceberg/CreateReadTasksDoFn.java   |    7 +-
 .../beam/sdk/io/iceberg/DynamicDestinations.java   |   12 +-
 .../beam/sdk/io/iceberg/FileWriteResult.java       |    4 +-
 .../beam/sdk/io/iceberg/IcebergCatalogConfig.java  |   18 +-
 .../IcebergCdcReadSchemaTransformProvider.java     |   34 +-
 .../org/apache/beam/sdk/io/iceberg/IcebergIO.java  |  119 +-
 .../IcebergReadSchemaTransformProvider.java        |    3 +-
 .../beam/sdk/io/iceberg/IcebergScanConfig.java     |  244 +-
 .../apache/beam/sdk/io/iceberg/IcebergUtils.java   |  374 +-
 .../IcebergWriteSchemaTransformProvider.java       |   12 +
 .../beam/sdk/io/iceberg/IncrementalScanSource.java |  102 -
 .../io/iceberg/OneTableDynamicDestinations.java    |   39 +-
 .../apache/beam/sdk/io/iceberg/PartitionUtils.java |  100 +-
 .../io/iceberg/PortableIcebergDestinations.java    |    3 +-
 .../apache/beam/sdk/io/iceberg/ReadFromTasks.java  |  103 -
 .../org/apache/beam/sdk/io/iceberg/ReadUtils.java  |  142 +-
 .../apache/beam/sdk/io/iceberg/RecordWriter.java   |   24 +-
 .../beam/sdk/io/iceberg/RecordWriterManager.java   |  166 +-
 .../org/apache/beam/sdk/io/iceberg/ScanSource.java |    7 +-
 .../apache/beam/sdk/io/iceberg/ScanTaskReader.java |    9 +-
 .../beam/sdk/io/iceberg/SerializableDataFile.java  |  192 +-
 .../sdk/io/iceberg/SerializableDeleteFile.java     |  335 +
 .../apache/beam/sdk/io/iceberg/SnapshotInfo.java   |    3 +-
 .../org/apache/beam/sdk/io/iceberg/TableCache.java |  206 +-
 .../beam/sdk/io/iceberg/WatchForSnapshots.java     |  195 -
 .../sdk/io/iceberg/WriteDirectRowsToFiles.java     |   25 +-
 .../sdk/io/iceberg/WriteGroupedRowsToFiles.java    |   27 +-
 .../io/iceberg/WritePartitionedRowsToFiles.java    |  195 +-
 .../beam/sdk/io/iceberg/WriteToDestinations.java   |   29 +-
 .../beam/sdk/io/iceberg/WriteToPartitions.java     |    9 +-
 .../sdk/io/iceberg/WriteUngroupedRowsToFiles.java  |   26 +-
 .../sdk/io/iceberg/cdc/ApplyWatermarkColumn.java   |   99 +
 .../beam/sdk/io/iceberg/cdc/CdcOutputUtils.java    |  195 +
 .../beam/sdk/io/iceberg/cdc/CdcReadUtils.java      |  700 ++
 .../beam/sdk/io/iceberg/cdc/CdcResolver.java       |  191 +
 .../beam/sdk/io/iceberg/cdc/CdcRowDescriptor.java  |   89 +
 .../sdk/io/iceberg/cdc/ChangelogDescriptor.java    |  106 +
 .../beam/sdk/io/iceberg/cdc/ChangelogScanner.java  | 1014 +++
 .../beam/sdk/io/iceberg/cdc/DeleteReader.java      |  309 +
 .../io/iceberg/cdc/IcebergCdcMetadataColumns.java  |   93 +
 .../io/iceberg/cdc/IncrementalChangelogSource.java |  211 +
 .../beam/sdk/io/iceberg/cdc/LocalResolveDoFn.java  |  249 +
 .../beam/sdk/io/iceberg/cdc/OverlapRange.java      |  102 +
 .../sdk/io/iceberg/cdc/ReadFromChangelogs.java     |  499 ++
 .../beam/sdk/io/iceberg/cdc/ResolveChanges.java    |  168 +
 .../io/iceberg/cdc/SerializableChangelogTask.java  |  279 +
 .../beam/sdk/io/iceberg/cdc/SnapshotWindowFn.java  |   87 +
 .../sdk/io/iceberg/cdc/WatchForSnapshotsSdf.java   |  317 +
 .../beam/sdk/io/iceberg/cdc}/package-info.java     |    4 +-
 .../iceberg/BeamBaseIncrementalChangelogScan.java  | 1014 +++
 .../java/org/apache/iceberg}/package-info.java     |    4 +-
 .../org/apache/beam/sdk/io/iceberg/AddFilesIT.java |   18 +-
 .../beam/sdk/io/iceberg/BeamRowWrapperTest.java    |    9 +-
 .../IcebergCdcReadSchemaTransformProviderTest.java |  126 +-
 .../beam/sdk/io/iceberg/IcebergIOReadTest.java     |  134 +-
 .../beam/sdk/io/iceberg/IcebergIOWriteTest.java    |   57 +-
 .../beam/sdk/io/iceberg/IcebergScanConfigTest.java |  270 +
 .../IcebergSchemaTransformTranslationTest.java     |   34 +-
 .../beam/sdk/io/iceberg/IcebergUtilsTest.java      |  117 +-
 .../IcebergWriteSchemaTransformProviderTest.java   |   76 +-
 .../beam/sdk/io/iceberg/PartitionUtilsTest.java    |   68 +
 .../apache/beam/sdk/io/iceberg/ReadUtilsTest.java  |   59 +-
 .../sdk/io/iceberg/RecordWriterManagerTest.java    |  212 +-
 .../sdk/io/iceberg/SerializableDataFileTest.java   |  222 +
 .../sdk/io/iceberg/SerializableDeleteFileTest.java |  222 +
 .../apache/beam/sdk/io/iceberg/TableCacheTest.java |  126 +
 .../beam/sdk/io/iceberg/TestDataWarehouse.java     |    2 +-
 .../catalog/BigQueryMetastoreCatalogIT.java        |    1 +
 .../io/iceberg/catalog/IcebergCatalogBaseIT.java   |  299 +-
 .../sdk/io/iceberg/catalog/RESTCatalogBLMSIT.java  |    1 -
 .../io/iceberg/cdc/ApplyWatermarkColumnTest.java   |  158 +
 .../beam/sdk/io/iceberg/cdc/CdcReadUtilsTest.java  |  381 +
 .../beam/sdk/io/iceberg/cdc/CdcResolverTest.java   |  156 +
 .../sdk/io/iceberg/cdc/ChangelogScannerTest.java   |  465 +
 .../beam/sdk/io/iceberg/cdc/DeleteReaderTest.java  |  419 +
 .../cdc/IncrementalChangelogSourceTest.java        |  514 ++
 .../sdk/io/iceberg/cdc/LocalResolveDoFnTest.java   |  340 +
 .../beam/sdk/io/iceberg/cdc/OverlapRangeTest.java  |  161 +
 .../sdk/io/iceberg/cdc/ReadFromChangelogsTest.java |  366 +
 .../sdk/io/iceberg/cdc/ResolveChangesTest.java     |  222 +
 .../iceberg/cdc/SerializableChangelogTaskTest.java |  256 +
 .../sdk/io/iceberg/cdc/SnapshotWindowFnTest.java   |   93 +
 .../io/iceberg/cdc/WatchForSnapshotsSdfTest.java   |  304 +
 sdks/java/io/influxdb/build.gradle                 |   10 +-
 .../apache/beam/sdk/io/influxdb/InfluxDbIO.java    |   22 +
 .../beam/sdk/io/influxdb/InfluxDbIOTest.java       |  104 +-
 sdks/java/io/jms/build.gradle                      |    6 +
 .../io/jms/BeamGenericJmsConnectionFactory.java    |   42 +
 .../beam/sdk/io/jms/ConnectionConfiguration.java   |  252 +
 .../apache/beam/sdk/io/jms/JmsCheckpointMark.java  |  113 +-
 .../java/org/apache/beam/sdk/io/jms/JmsIO.java     |  224 +-
 .../sdk/io/jms/JmsReadSchemaTransformProvider.java |  227 +
 .../io/jms/JmsWriteSchemaTransformProvider.java    |  191 +
 .../java/org/apache/beam/sdk/io/jms/CommonJms.java |   27 +-
 .../sdk/io/jms/ConnectionConfigurationTest.java    |  159 +
 .../java/org/apache/beam/sdk/io/jms/JmsIOIT.java   |  209 +-
 .../java/org/apache/beam/sdk/io/jms/JmsIOTest.java |  222 +-
 .../org/apache/beam/sdk/io/jms/JmsLocalTest.java   |  245 +
 .../sdk/io/jms/JmsSchemaTransformProviderTest.java |  253 +
 sdks/java/io/kafka/build.gradle                    |   14 +-
 sdks/java/io/kafka/kafka-201/build.gradle          |   24 -
 sdks/java/io/kafka/kafka-241/build.gradle          |   24 -
 sdks/java/io/kafka/kafka-251/build.gradle          |   24 -
 sdks/java/io/kafka/kafka-282/build.gradle          |   24 -
 sdks/java/io/kafka/kafka-312/build.gradle          |   24 -
 sdks/java/io/kafka/kafka-390/build.gradle          |   24 -
 .../io/kafka/{kafka-231 => kafka-392}/build.gradle |    6 +-
 .../beam/sdk/io/kafka/KafkaCommitOffset.java       |    6 +
 .../java/org/apache/beam/sdk/io/kafka/KafkaIO.java |  163 +-
 .../KafkaIOReadImplementationCompatibility.java    |    6 +
 .../beam/sdk/io/kafka/KafkaUnboundedReader.java    |   28 +-
 .../beam/sdk/io/kafka/ReadFromKafkaDoFn.java       |  114 +-
 .../beam/sdk/io/kafka/KafkaCommitOffsetTest.java   |    4 +-
 ...KafkaIOReadImplementationCompatibilityTest.java |   22 +
 .../io/messaging-expansion-service/build.gradle    |   55 +
 sdks/java/io/mongodb/build.gradle                  |    2 +
 .../beam/sdk/io/mongodb/AggregationQuery.java      |    9 +-
 .../beam/sdk/io/mongodb/MongoDbGridFSIO.java       |    8 +-
 .../org/apache/beam/sdk/io/mongodb/MongoDbIO.java  |   14 +-
 .../MongoDbReadSchemaTransformConfiguration.java   |   89 +
 .../MongoDbReadSchemaTransformProvider.java        |  151 +
 .../apache/beam/sdk/io/mongodb/MongoDbUtils.java   |  205 +
 .../MongoDbWriteSchemaTransformConfiguration.java  |   85 +
 .../MongoDbWriteSchemaTransformProvider.java       |  168 +
 .../MongoDbReadSchemaTransformProviderTest.java    |  221 +
 .../beam/sdk/io/mongodb/MongoDbUtilsTest.java      |  258 +
 .../MongoDbWriteSchemaTransformProviderTest.java   |  166 +
 sdks/java/io/mqtt/build.gradle                     |    3 +-
 .../java/org/apache/beam/sdk/io/mqtt/MqttIO.java   |  189 +-
 .../io/mqtt/MqttReadSchemaTransformProvider.java   |  150 +
 .../io/mqtt/MqttWriteSchemaTransformProvider.java  |  132 +
 .../org/apache/beam/sdk/io/mqtt/MqttIOTest.java    |  204 +-
 .../io/mqtt/MqttSchemaTransformProviderTest.java   |  252 +
 .../org/apache/beam/sdk/io/parquet/ParquetIO.java  |   47 +
 .../org/apache/beam/io/requestresponse/Call.java   |   12 +-
 .../apache/beam/io/requestresponse/CallTest.java   |  104 +-
 .../singlestore/SingleStoreDefaultRowMapper.java   |    3 +-
 .../io/snowflake/crosslanguage/package-info.java   |    4 -
 .../beam/sdk/io/solace/broker/BrokerResponse.java  |   13 +-
 .../sdk/io/solace/broker/BrokerResponseTest.java   |   66 +
 sdks/java/io/sparkreceiver/3/build.gradle          |    6 +
 .../java/org/apache/beam/sdk/managed/Managed.java  |    9 +
 .../ml/inference/{remote => gemini}/build.gradle   |   31 +-
 .../ml/inference/gemini/GeminiImageResponse.java   |   17 +-
 .../inference/gemini/GeminiInferenceFunctions.java |   67 +
 .../ml/inference/gemini/GeminiModelHandler.java    |   86 +
 .../ml/inference/gemini/GeminiModelParameters.java |   60 +
 .../ml/inference/gemini/GeminiRequestFunction.java |   35 +-
 .../sdk/ml/inference/gemini/GeminiStringInput.java |   17 +-
 .../ml/inference/gemini/GeminiStringResponse.java  |   17 +-
 .../sdk/ml/inference/gemini}/package-info.java     |    4 +-
 .../ml/inference/gemini/GeminiModelHandlerIT.java  |  163 +
 .../inference/gemini/GeminiModelHandlerTest.java   |  133 +
 sdks/java/ml/inference/remote/build.gradle         |    1 +
 .../ml/inference/remote/PredictionResultCoder.java |   74 +
 .../sdk/ml/inference/remote/RemoteInference.java   |  167 +-
 .../beam/sdk/ml/inference/remote/RetryHandler.java |   27 +-
 .../ml/inference/remote/RemoteInferenceTest.java   |  144 +-
 .../sdk/ml/inference/remote/RetryHandlerTest.java  |  145 +
 sdks/java/testing/kafka-service/build.gradle       |    3 +-
 .../apache/beam/sdk/testing/kafka/LocalKafka.java  |   13 +-
 sdks/java/testing/load-tests/build.gradle          |   30 +
 sdks/java/testing/nexmark/build.gradle             |   30 +
 .../apache/beam/sdk/nexmark/NexmarkLauncher.java   |    2 +
 sdks/java/testing/tpcds/build.gradle               |   30 +
 sdks/java/testing/watermarks/build.gradle          |   25 +-
 .../controller-container/Dockerfile                |    2 +-
 sdks/python/.yapfignore                            |    1 +
 sdks/python/apache_beam/coders/row_coder_test.py   |   54 +
 sdks/python/apache_beam/dataframe/frames.py        |    7 +-
 sdks/python/apache_beam/dataframe/frames_test.py   |   37 +-
 sdks/python/apache_beam/dataframe/io.py            |   38 +-
 sdks/python/apache_beam/dataframe/io_test.py       |   66 +-
 .../apache_beam/examples/complete/autocomplete.py  |    5 +-
 .../examples/complete/game/game_stats.py           |   14 +-
 .../examples/complete/game/hourly_team_score.py    |   19 +-
 .../examples/complete/game/leader_board.py         |   24 +-
 .../examples/complete/game/user_score.py           |    9 +-
 .../examples/complete/top_wikipedia_sessions.py    |    5 +-
 .../examples/cookbook/bigtableio_it_test.py        |    4 +-
 .../cookbook/ordered_window_elements/streaming.py  |    2 +-
 .../anomaly_detection_pipeline/setup.py            |    2 +-
 .../examples/inference/gemini_image_generation.py  |    4 +-
 .../inference/gemini_text_classification.py        |    4 +-
 .../online_clustering/clustering_pipeline/setup.py |    6 +-
 .../examples/inference/pytorch_image_captioning.py |  641 ++
 ...ytorch_image_classification_with_side_inputs.py |    2 +-
 .../inference/pytorch_image_object_detection.py    |  517 ++
 .../inference/pytorch_imagenet_rightfit.py         |  537 ++
 .../examples/inference/pytorch_sentiment.py        |   80 +-
 .../examples/inference/vllm_text_completion.py     |   24 +-
 .../apache_beam/examples/kafkataxi/README.md       |    2 +-
 .../apache_beam/examples/matrix_power_test.py      |   13 +-
 .../kfp/components/preprocessing/requirements.txt  |    2 +-
 .../kfp/components/train/requirements.txt          |    4 +-
 .../apache_beam/examples/ml_transform/README.md    |  116 +
 .../ml_transform/mltransform_generate_vocab.py     |  265 +
 .../mltransform_generate_vocab_requirements.txt}   |   14 +-
 .../mltransform_generate_vocab_test.py             |  191 +
 .../ml_transform/mltransform_image_embedding.py    |  263 +
 .../mltransform_image_embedding_test.py            |  105 +
 .../ml_transform/mltransform_one_hot_encoding.py   |  266 +
 .../mltransform_one_hot_encoding_requirements.txt} |   15 +-
 .../mltransform_one_hot_encoding_test.py           |  257 +
 .../ml_transform/mltransform_text_embedding.py     |  186 +
 .../mltransform_text_embedding_test.py             |   94 +
 .../transforms/elementwise/enrichment_test.py      |    2 +-
 .../apache_beam/examples/wordcount_debugging.py    |    4 +-
 .../apache_beam/examples/wordcount_rust/README.md  |    4 +-
 .../examples/wordcount_rust/requirements.txt       |    2 +-
 .../wordcount_rust/word_processing/Cargo.lock      |   65 +-
 .../wordcount_rust/word_processing/Cargo.toml      |    2 +-
 .../apache_beam/examples/wordcount_with_metrics.py |    4 +-
 .../apache_beam/io/components/rate_limit.proto     |   64 +
 .../apache_beam/io/components/rate_limit_pb2.py    |   63 +
 .../apache_beam/io/components/rate_limiter.py      |   39 +-
 .../apache_beam/io/components/rate_limiter_test.py |   62 +-
 sdks/python/apache_beam/io/external/gcp/pubsub.py  |   11 +-
 .../io/external/xlang_bigqueryio_it_test.py        |    2 +-
 .../io/external/xlang_debeziumio_it_test.py        |    2 +-
 .../apache_beam/io/external/xlang_jmsio_it_test.py |  383 +
 .../io/external/xlang_mqttio_it_test.py            |  262 +
 sdks/python/apache_beam/io/fileio.py               |    3 +-
 sdks/python/apache_beam/io/filesystemio.py         |    5 +-
 sdks/python/apache_beam/io/gcp/__init__.py         |   20 -
 sdks/python/apache_beam/io/gcp/bigquery.py         |   54 +-
 .../apache_beam/io/gcp/bigquery_change_history.py  |   10 +-
 .../apache_beam/io/gcp/bigquery_file_loads.py      |   85 +-
 .../apache_beam/io/gcp/bigquery_file_loads_test.py |  151 +
 .../apache_beam/io/gcp/bigquery_schema_tools.py    |   63 +-
 .../io/gcp/bigquery_schema_tools_test.py           |  183 +-
 sdks/python/apache_beam/io/gcp/bigquery_test.py    |   53 +-
 sdks/python/apache_beam/io/gcp/bigquery_tools.py   |   21 +-
 .../apache_beam/io/gcp/bigquery_tools_test.py      |  161 +
 .../apache_beam/io/gcp/bigquery_write_it_test.py   |   71 +-
 .../apache_beam/io/gcp/bigtableio_it_test.py       |   10 +-
 .../apache_beam/io/gcp/gcsfilesystem_test.py       |    4 +-
 sdks/python/apache_beam/io/gcp/gcsio.py            |   61 +-
 .../apache_beam/io/gcp/gcsio_integration_test.py   |   21 +-
 sdks/python/apache_beam/io/gcp/gcsio_test.py       |  132 +
 .../apache_beam/io/gcp/healthcare/dicomclient.py   |    5 +-
 .../clients/bigquery/bigquery_v2_client.py         |  799 +-
 .../clients/bigquery/bigquery_v2_messages.py       | 9263 ++++++++++++++++----
 sdks/python/apache_beam/io/gcp/pubsub.py           |   48 +-
 .../apache_beam/io/gcp/pubsub_integration_test.py  |  100 +
 sdks/python/apache_beam/io/gcp/pubsub_test.py      |  193 +-
 sdks/python/apache_beam/io/iobase.py               |   14 +-
 sdks/python/apache_beam/io/iobase_test.py          |   49 +-
 sdks/python/apache_beam/io/requestresponse.py      |    5 +-
 sdks/python/apache_beam/io/requestresponse_test.py |   10 +
 sdks/python/apache_beam/io/textio_test.py          |    7 +-
 sdks/python/apache_beam/io/unbounded_source.py     | 1136 +++
 .../python/apache_beam/io/unbounded_source_test.py | 1252 +++
 sdks/python/apache_beam/io/watch.py                |  901 ++
 sdks/python/apache_beam/io/watch_test.py           |  858 ++
 sdks/python/apache_beam/metrics/metric.py          |    2 +-
 sdks/python/apache_beam/ml/anomaly/specifiable.py  |    2 +-
 .../apache_beam/ml/anomaly/univariate/mean_test.py |    8 +-
 .../apache_beam/ml/anomaly/univariate/perf_test.py |   51 +-
 .../ml/anomaly/univariate/quantile_test.py         |    8 +-
 .../ml/anomaly/univariate/stdev_test.py            |    8 +-
 .../ml/inference/agent_development_kit.py          |   36 +-
 .../ml/inference/agent_development_kit_test.py     |   38 +-
 sdks/python/apache_beam/ml/inference/base.py       |   11 +-
 sdks/python/apache_beam/ml/inference/base_test.py  |    2 +-
 .../pytorch_image_captioning_requirements.txt}     |   13 +-
 ...ytorch_image_object_detection_requirements.txt} |   10 +-
 .../ml/inference/pytorch_inference_test.py         |    5 +-
 .../inference/pytorch_rightfit_requirements.txt}   |   12 +-
 .../ml/inference/sklearn_inference_test.py         |    4 +-
 .../ml/inference/tensorflow_inference_test.py      |    4 +-
 .../ml/inference/tensorrt_inference_test.py        |   11 +-
 .../inference/test_resources/vllm.dockerfile.old   |   14 +-
 .../apache_beam/ml/inference/vllm_inference.py     |  331 +-
 .../ml/inference/vllm_inference_test.py            |  157 +
 .../ml/rag/enrichment/milvus_search_it_test.py     |    1 +
 .../apache_beam/ml/rag/ingestion/milvus_search.py  |   13 +-
 .../ml/rag/ingestion/milvus_search_it_test.py      |    1 +
 .../ml/rag/ingestion/milvus_search_test.py         |   35 +
 sdks/python/apache_beam/ml/rag/ingestion/qdrant.py |  328 +
 .../apache_beam/ml/rag/ingestion/qdrant_it_test.py |  326 +
 .../apache_beam/ml/rag/ingestion/qdrant_test.py    |  480 +
 sdks/python/apache_beam/ml/rag/test_utils.py       |    1 +
 .../mltransform_embedding_tests_requirements.txt}  |    5 +-
 .../python/apache_beam/options/pipeline_options.py |  107 +-
 .../apache_beam/options/pipeline_options_test.py   |   62 +
 .../options/pipeline_options_validator.py          |   12 +-
 sdks/python/apache_beam/portability/common_urns.py |    1 +
 sdks/python/apache_beam/runners/common.pxd         |    1 +
 sdks/python/apache_beam/runners/common.py          |   35 +
 sdks/python/apache_beam/runners/common_test.py     |   58 +
 .../runners/dataflow/dataflow_metrics.py           |   65 +-
 .../runners/dataflow/dataflow_metrics_test.py      |  314 +-
 .../runners/dataflow/dataflow_runner.py            |   72 +-
 .../runners/dataflow/dataflow_runner_test.py       |   59 +-
 .../runners/dataflow/internal/apiclient.py         |  537 +-
 .../runners/dataflow/internal/apiclient_test.py    |  354 +-
 .../runners/dataflow/internal/clients/README.txt   |   11 -
 .../dataflow/internal/clients/dataflow/__init__.py |   34 -
 .../clients/dataflow/dataflow_v1b3_client.py       | 1318 ---
 .../clients/dataflow/dataflow_v1b3_messages.py     | 8072 -----------------
 .../internal/clients/dataflow/message_matchers.py  |  118 -
 .../clients/dataflow/message_matchers_test.py      |   74 -
 .../apache_beam/runners/dataflow/internal/names.py |    2 +-
 .../runners/direct/transform_evaluator.py          |  118 +-
 .../dataproc/dataproc_cluster_manager.py           |    3 +-
 .../runners/interactive/interactive_beam_test.py   |    8 +-
 .../runners/interactive/interactive_environment.py |   10 +-
 .../runners/interactive/recording_manager.py       |  199 +-
 .../runners/interactive/recording_manager_test.py  |  191 +-
 .../runners/interactive/user_pipeline_tracker.py   |  109 +-
 .../interactive/user_pipeline_tracker_test.py      |   48 +
 .../runners/portability/beam_plugins_it_test.py    |   70 +
 .../portability/fn_api_runner/fn_runner_test.py    |   29 +
 .../apache_beam/runners/portability/job_server.py  |    9 +-
 .../runners/portability/job_server_test.py         |    2 +-
 .../runners/portability/portable_runner_test.py    |   17 +
 .../runners/portability/prism_runner.py            |   15 +-
 .../runners/portability/prism_runner_test.py       |   11 +
 .../runners/portability/sdk_container_builder.py   |    2 +-
 .../portability/sdk_container_builder_test.py      |   41 +
 .../runners/portability/spark_runner.py            |   11 +
 .../runners/portability/spark_runner_test.py       |   20 +-
 .../apache_beam/runners/portability/stager.py      |  149 +-
 .../apache_beam/runners/worker/bundle_processor.py |   37 +-
 .../apache_beam/runners/worker/data_plane.py       |   60 +-
 .../apache_beam/runners/worker/opcounters.py       |    2 +
 .../apache_beam/runners/worker/opcounters_test.py  |   16 +
 .../apache_beam/runners/worker/sdk_worker.py       |   12 +-
 .../apache_beam/runners/worker/sdk_worker_main.py  |   20 +-
 .../runners/worker/sdk_worker_main_test.py         |   49 +
 .../apache_beam/runners/worker/sdk_worker_test.py  |   38 +
 .../runners/worker/statesampler_fast.pyx           |   19 +-
 .../runners/worker/statesampler_test.py            |   43 +
 .../testing/benchmarks/chicago_taxi/run_chicago.sh |    6 +-
 .../benchmarks/cloudml/criteo_tft/criteo.py        |   30 +-
 .../benchmarks/cloudml/criteo_tft/criteo_test.py   |   93 +
 .../testing/benchmarks/cloudml/requirements.txt    |   16 +-
 .../mltransform_generate_vocab_benchmark.py        |   56 +
 .../mltransform_image_embedding_benchmark.py       |  128 +
 .../mltransform_one_hot_encoding_benchmark.py      |  150 +
 .../mltransform_text_embedding_benchmark.py        |  117 +
 .../pytorch_image_captioning_benchmarks.py         |   42 +
 .../pytorch_image_object_detection_benchmarks.py   |   42 +
 .../pytorch_imagenet_rightfit_benchmarks.py        |   42 +
 .../inference/table_row_inference_benchmark.py     |   48 +-
 .../testing/load_tests/dataflow_cost_benchmark.py  |  189 +-
 sdks/python/apache_beam/testing/test_pipeline.py   |    4 +-
 sdks/python/apache_beam/testing/util.py            |   43 +
 sdks/python/apache_beam/testing/util_test.py       |   44 +
 sdks/python/apache_beam/transforms/async_dofn.py   |    4 +-
 .../apache_beam/transforms/async_dofn_test.py      |   62 +-
 sdks/python/apache_beam/transforms/combiners.py    |   58 +
 .../apache_beam/transforms/combiners_test.py       |   54 +
 sdks/python/apache_beam/transforms/core.py         |    4 +-
 sdks/python/apache_beam/transforms/external.py     |    1 +
 sdks/python/apache_beam/transforms/managed.py      |    4 +-
 .../transforms/managed_iceberg_it_test.py          |    7 +-
 sdks/python/apache_beam/transforms/stats_test.py   |   20 +-
 .../apache_beam/transforms/userstate_test.py       |   69 +
 sdks/python/apache_beam/transforms/util_test.py    |   49 +
 sdks/python/apache_beam/typehints/opcodes.py       |    8 +
 sdks/python/apache_beam/typehints/row_type_test.py |   21 +
 sdks/python/apache_beam/typehints/schemas.py       |  152 +-
 sdks/python/apache_beam/typehints/schemas_test.py  |  120 +
 .../apache_beam/typehints/trivial_inference.py     |   26 +
 .../typehints/trivial_inference_test.py            |   27 +
 sdks/python/apache_beam/typehints/typehints.py     |   14 +-
 .../python/apache_beam/typehints/typehints_test.py |   18 +
 sdks/python/apache_beam/utils/subprocess_server.py |  147 +-
 .../apache_beam/utils/subprocess_server_test.py    |  171 +-
 sdks/python/apache_beam/utils/timestamp.py         |  332 +-
 sdks/python/apache_beam/utils/timestamp_test.py    |  204 +-
 sdks/python/apache_beam/version.py                 |    2 +-
 sdks/python/apache_beam/yaml/examples/README.md    |    4 +
 .../yaml/extended_tests/databases/debezium.yaml    |   50 +
 .../yaml/extended_tests/databases/iceberg.yaml     |   57 +-
 sdks/python/apache_beam/yaml/integration_tests.py  |  259 +-
 sdks/python/apache_beam/yaml/main.py               |    7 +-
 sdks/python/apache_beam/yaml/readme_test.py        |   16 +-
 sdks/python/apache_beam/yaml/standard_io.yaml      |   79 +-
 .../apache_beam/yaml/standard_providers.yaml       |    1 +
 .../clients => yaml/test_utils}/__init__.py        |    2 +
 .../yaml/test_utils/datadog_test_utils.py          |  131 +
 sdks/python/apache_beam/yaml/tests/datadog.yaml    |   63 +
 .../Cargo.toml => yaml/tests/delta.yaml}           |   28 +-
 .../apache_beam/yaml/tests/iceberg_add_files.yaml  |    3 +
 sdks/python/apache_beam/yaml/tests/match_all.yaml  |   96 +
 sdks/python/apache_beam/yaml/tests/mongodb.yaml    |   74 +
 .../yaml/tests/runinference_huggingface.yaml       |   62 +
 ...uninference.yaml => runinference_vertexai.yaml} |    0
 sdks/python/apache_beam/yaml/yaml_io.py            |  333 +-
 sdks/python/apache_beam/yaml/yaml_io_test.py       |  440 +
 sdks/python/apache_beam/yaml/yaml_mapping.py       |  171 +-
 sdks/python/apache_beam/yaml/yaml_ml.py            |   95 +-
 sdks/python/apache_beam/yaml/yaml_provider.py      |   11 +-
 sdks/python/apache_beam/yaml/yaml_transform.py     |   56 +-
 .../python/apache_beam/yaml/yaml_transform_test.py |  121 +-
 sdks/python/apache_beam/yaml/yaml_udf_test.py      |  119 +-
 sdks/python/build.gradle                           |   33 +-
 sdks/python/container/Dockerfile                   |   20 +-
 .../container/base_image_requirements_manual.txt   |    5 +-
 sdks/python/container/boot.go                      |  179 +-
 .../license_scripts/upgrade_bundled_pip.py         |   76 +
 .../container/ml/py310/base_image_requirements.txt |  194 +-
 .../container/ml/py310/gpu_image_requirements.txt  |  290 +-
 .../container/ml/py311/base_image_requirements.txt |  200 +-
 .../container/ml/py311/gpu_image_requirements.txt  |  292 +-
 .../container/ml/py312/base_image_requirements.txt |  199 +-
 .../container/ml/py312/gpu_image_requirements.txt  |  292 +-
 .../container/ml/py313/base_image_requirements.txt |  197 +-
 sdks/python/container/piputil.go                   |   39 +-
 sdks/python/container/profiler.go                  |  624 ++
 sdks/python/container/profiler_test.go             |  134 +
 .../container/py310/base_image_requirements.txt    |  181 +-
 .../container/py311/base_image_requirements.txt    |  187 +-
 .../container/py312/base_image_requirements.txt    |  182 +-
 .../container/py313/base_image_requirements.txt    |  184 +-
 .../container/py314/base_image_requirements.txt    |  186 +-
 sdks/python/container/run_generate_requirements.sh |    4 +-
 sdks/python/pyproject.toml                         |  118 +-
 sdks/python/pytest.ini                             |    1 +
 sdks/python/scripts/run_integration_test.sh        |   22 +-
 sdks/python/scripts/run_job_server.sh              |    8 +-
 sdks/python/scripts/run_lint.sh                    |   38 +-
 sdks/python/scripts/run_snapshot_publish.sh        |    4 +-
 sdks/python/setup.py                               |   71 +-
 sdks/python/test-suites/dataflow/common.gradle     |   54 +-
 sdks/python/test-suites/direct/build.gradle        |    7 +
 sdks/python/test-suites/direct/common.gradle       |   21 +-
 sdks/python/test-suites/tox/py310/build.gradle     |   27 +-
 sdks/python/test-suites/xlang/build.gradle         |   11 +-
 sdks/python/tox.ini                                |   11 +-
 sdks/standard_expansion_services.yaml              |   21 +
 sdks/standard_external_transforms.yaml             |  257 +-
 sdks/typescript/container/boot.go                  |    4 +-
 sdks/typescript/package-lock.json                  |  226 +-
 sdks/typescript/package.json                       |    6 +-
 sdks/typescript/src/apache_beam/runners/flink.ts   |    2 +-
 settings.gradle.kts                                |   39 +-
 website/Dockerfile                                 |    2 +-
 website/build.gradle                               |    5 +-
 website/www/site/assets/js/bootstrap.js            |   46 +-
 website/www/site/assets/js/bootstrap/alert.js      |    7 +-
 website/www/site/assets/js/bootstrap/carousel.js   |    9 +-
 website/www/site/assets/js/bootstrap/collapse.js   |    6 +-
 website/www/site/assets/js/bootstrap/dropdown.js   |    7 +-
 website/www/site/assets/js/bootstrap/modal.js      |    8 +-
 website/www/site/assets/js/bootstrap/tooltip.js    |    9 +-
 website/www/site/assets/js/page-nav.js             |    2 +-
 .../www/site/assets/scss/_capability-matrix.scss   |  123 +-
 .../www/site/assets/scss/capability-matrix.scss    |  121 +-
 website/www/site/config.toml                       |    2 +-
 .../content/en/blog/apache-hop-with-dataflow.md    |   12 +-
 website/www/site/content/en/blog/beam-2.68.0.md    |    2 +-
 website/www/site/content/en/blog/beam-2.74.0.md    |   70 +
 website/www/site/content/en/blog/beam-2.75.0.md    |   72 +
 .../content/en/blog/beam-sql-with-notebooks.md     |    4 +-
 .../beam-summit-2026-interview-with-raj-katakam.md |   88 +
 .../site/content/en/documentation/dsls/sql/ddl.md  |  344 +
 .../dsls/sql/extensions/create-external-table.md   |   52 -
 .../content/en/documentation/dsls/sql/overview.md  |    2 +
 .../documentation/io/built-in/google-bigquery.md   |    2 +-
 .../en/documentation/io/built-in/iceberg.md        |    2 +-
 .../site/content/en/documentation/io/connectors.md |   56 +-
 .../site/content/en/documentation/io/managed-io.md |  771 +-
 .../site/content/en/documentation/runners/flink.md |   37 +-
 .../site/content/en/documentation/runners/spark.md |    4 +-
 .../sdks/python-multi-language-pipelines.md        |    2 +-
 .../sdks/python-pipeline-dependencies.md           |   41 +-
 .../site/content/en/documentation/sdks/python.md   |    6 +-
 .../www/site/content/en/documentation/sdks/yaml.md |   34 +
 .../transforms/java/aggregation/batchelements.md   |   31 +
 .../en/documentation/transforms/java/overview.md   |    1 +
 .../www/site/content/en/get-started/downloads.md   |   30 +-
 .../site/content/en/get-started/quickstart-java.md |    4 +-
 .../site/content/en/get-started/quickstart-py.md   |   19 +-
 website/www/site/content/en/performance/_index.md  |   18 +
 .../content/en/performance/icebergio/_index.md     |   50 +
 .../mltransform-image-embedding-cpu/_index.md      |   48 +
 .../mltransform-image-embedding-gpu/_index.md      |   48 +
 .../mltransform-text-embedding/_index.md           |   44 +
 .../en/performance/mltransformonehot/_index.md     |   42 +
 .../en/performance/mltransformvocab/_index.md      |   37 +
 .../pytorchimagecaptioningbatchcpu/_index.md       |   41 +
 .../pytorchimagecaptioningbatchgpu/_index.md       |   41 +
 .../pytorchimagecaptioningstreamingcpu/_index.md   |   41 +
 .../pytorchimagecaptioningstreaminggpu/_index.md   |   41 +
 .../pytorchimagenetrightfitcpu/_index.md           |   43 +
 .../pytorchimagenetrightfitgpu/_index.md           |   43 +
 .../pytorchimagenetrightfitoncecpu/_index.md       |   44 +
 .../pytorchimagenetrightfitoncegpu/_index.md       |   44 +
 .../pytorchimageobjectdetectionbatchcpu/_index.md  |   41 +
 .../pytorchimageobjectdetectionbatchgpu/_index.md  |   41 +
 .../_index.md                                      |   41 +
 .../_index.md                                      |   41 +
 website/www/site/data/authors.yml                  |    3 +
 website/www/site/data/performance.yaml             |  303 +
 .../partials/section-menu/en/documentation.html    |    1 +
 .../layouts/partials/section-menu/en/sdks.html     |    1 +
 .../documentation/capability-matrix-big.html       |   36 +-
 .../documentation/capability-matrix-single.html    |   37 +-
 website/www/site/layouts/shortcodes/tab.html       |   22 +
 website/www/site/static/.htaccess                  |   14 -
 website/www/yarn.lock                              |    6 +-
 1700 files changed, 97187 insertions(+), 25381 deletions(-)


Reply via email to