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(-)
