This is an automated email from the ASF dual-hosted git repository.
github-bot pushed a change to branch nightly-refs/heads/master
in repository https://gitbox.apache.org/repos/asf/beam.git
from 3cdb9fe2b2d Bump Python FnAPI beam-master container. (#28704)
add 218bda98624 added jobs to GitHub Actions (#28679)
add b74a0dc8e65 Add Load Tests Combine Dataflow Batch Java workflow
(#28676)
add e53b68b368d Add GitHub Workflow Replacements for Jenkins
job_PerformanceTests_SpannerIO and
job_PerformanceTests_SQLBigQueryIO_Batch_Java (#28555)
add 941e77d295a beam_PostCommit_Java_Tpcds (#28495)
add f824adcae1f Add GitHub Workflow Replacements for Jenkins
job_PerformanceTests_JDBC (#28602)
add 072848f069d Add Load Tests Combine/ParDo SparkStructuredStreaming
Batch Java workflows (#28714)
add c3cfccae1e1 Add Load Tests ParDo Dataflow Java workflows (#28713)
add 052d2644c17 Add Load Tests Combine Dataflow Streaming Java workflow
(#28677)
add 50f2b116596 Add GitHub Workflow Replacements for Jenkins
job_PerformanceTests_Kafka_IO (#28712)
add 1c0f1beede0 Github Workflow Replacement for Jenkins Jobs,
beam_PerformanceTests_XmlIOIT* and beam_PerformanceTests_TFRecordIOIT (#28584)
add e328ab51f1b Github Workflow Replacement for Jenkins Jobs,
beam_PerformanceTests_ManyFiles_TextIOIT* (#28581)
add dda0eb9d642 Create RequestResponseIO gradle project (#28706)
add b9ac87ecbe2 add write to checks due to publish unit test result
(#28720)
add b371ab9e0b3 Bump get-func-name from 2.0.0 to 2.0.2 in /sdks/typescript
(#28707)
add 5f6817e8213 Run other arm tests on Dataflow Java (#28666)
add 2b7df3b64b3 Bump github.com/aws/aws-sdk-go-v2/service/s3 in /sdks
(#28674)
add e6dcd372123 Update example in Tour of Beam (#28718)
add 56b1259cac9 Revert "Replace StorageV1 client with GCS client (#28079)"
(#28721)
add 170310ee83a Break out nested classes from StreamingDataflowWorker.
(#28537)
add 5515f18b915 Bump github.com/aws/aws-sdk-go-v2/feature/s3/manager in
/sdks (#28673)
add a4ee8548424 assign highmem runner to beam_PostCommit_Python and to
beam_PreCommit_Java_GCP_IO_Direct (#28719)
add 2d0a52f9346 Remove repeated test scenarios. (#28669)
add dc9bec8f7c9 Github Workflow Replacement for Jenkins Jobs,
beam_PerformanceTests_Compressed_TextIOIT* (#28606)
add a6de3033072 Github Workflow Replacement for Jenkins Jobs,
beam_PerformanceTests_ParquetIOIT* (#28582)
add a9516ba18a7 Add arm tests to postcommit suite (#28644)
add a46bc12a256 Ensure configuration schema used to decode configuration.
(#28727)
No new revisions were added by this update.
Summary of changes:
.github/workflows/README.md | 8 +-
...beam_LoadTests_Java_Combine_Dataflow_Batch.yml} | 71 +-
...m_LoadTests_Java_Combine_Dataflow_Streaming.yml | 107 +
...Java_Combine_SparkStructuredStreaming_Batch.yml | 107 +
.../beam_LoadTests_Java_ParDo_Dataflow_Batch.yml | 117 +
...eam_LoadTests_Java_ParDo_Dataflow_Streaming.yml | 117 +
...s_Java_ParDo_SparkStructuredStreaming_Batch.yml | 117 +
...eam_PerformanceTests_BiqQueryIO_Read_Python.yml | 35 +-
...formanceTests_BiqQueryIO_Write_Python_Batch.yml | 33 +-
...d_Python.yml => beam_PerformanceTests_Cdap.yml} | 69 +-
... beam_PerformanceTests_Compressed_TextIOIT.yml} | 52 +-
..._PerformanceTests_Compressed_TextIOIT_HDFS.yml} | 70 +-
....yml => beam_PerformanceTests_HadoopFormat.yml} | 69 +-
...d_Python.yml => beam_PerformanceTests_JDBC.yml} | 68 +-
.../workflows/beam_PerformanceTests_Kafka_IO.yml | 122 +
...> beam_PerformanceTests_ManyFiles_TextIOIT.yml} | 52 +-
...m_PerformanceTests_ManyFiles_TextIOIT_HDFS.yml} | 70 +-
....yml => beam_PerformanceTests_MongoDBIO_IT.yml} | 69 +-
...n.yml => beam_PerformanceTests_ParquetIOIT.yml} | 52 +-
... => beam_PerformanceTests_ParquetIOIT_HDFS.yml} | 70 +-
...rformanceTests_PubsubIOIT_Python_Streaming.yml} | 51 +-
..._PerformanceTests_SQLBigQueryIO_Batch_Java.yml} | 65 +-
...PerformanceTests_SpannerIO_Read_2GB_Python.yml} | 51 +-
...anceTests_SpannerIO_Write_2GB_Python_Batch.yml} | 51 +-
... => beam_PerformanceTests_SparkReceiver_IO.yml} | 69 +-
....yml => beam_PerformanceTests_TFRecordIOIT.yml} | 45 +-
.../beam_PerformanceTests_TFRecordIOIT_HDFS.yml | 112 +
...erformanceTests_WordCountIT_PythonVersions.yml} | 72 +-
..._Test.yml => beam_PerformanceTests_XmlIOIT.yml} | 45 +-
....yml => beam_PerformanceTests_XmlIOIT_HDFS.yml} | 70 +-
.github/workflows/beam_PostCommit_Java.yml | 4 +-
.../beam_PostCommit_Java_Examples_Dataflow_ARM.yml | 1 +
...yml => beam_PostCommit_Java_Tpcds_Dataflow.yml} | 49 +-
...st.yml => beam_PostCommit_Java_Tpcds_Flink.yml} | 46 +-
...st.yml => beam_PostCommit_Java_Tpcds_Spark.yml} | 45 +-
.github/workflows/beam_PostCommit_Python.yml | 2 +-
...t_Python.yml => beam_PostCommit_Python_Arm.yml} | 22 +-
.github/workflows/beam_PostCommit_Website_Test.yml | 2 +-
.../beam_PreCommit_Java_GCP_IO_Direct.yml | 2 +-
.../config_Combine_Java_Dataflow_Batch_10b.txt} | 29 +-
...onfig_Combine_Java_Dataflow_Batch_Fanout_4.txt} | 29 +-
...onfig_Combine_Java_Dataflow_Batch_Fanout_8.txt} | 29 +-
.../java_Combine_Dataflow_Streaming_10b.txt | 31 +
.../java_Combine_Dataflow_Streaming_Fanout_4.txt | 31 +
.../java_Combine_Dataflow_Streaming_Fanout_8.txt | 31 +
..._Combine_SparkStructuredStreaming_Batch_10b.txt | 27 +
...ine_SparkStructuredStreaming_Batch_Fanout_4.txt | 27 +
...ine_SparkStructuredStreaming_Batch_Fanout_8.txt | 27 +
.../java_ParDo_Dataflow_Batch_100_counters.txt | 29 +
.../java_ParDo_Dataflow_Batch_10_counters.txt | 29 +
.../java_ParDo_Dataflow_Batch_10_times.txt | 29 +
.../java_ParDo_Dataflow_Batch_200_times.txt | 30 +
.../java_ParDo_Dataflow_Streaming_100_counters.txt | 30 +
.../java_ParDo_Dataflow_Streaming_10_counters.txt | 30 +
.../java_ParDo_Dataflow_Streaming_10_times.txt | 30 +
.../java_ParDo_Dataflow_Streaming_200_times.txt | 30 +
...SparkStructuredStreaming_Batch_100_counters.txt | 26 +
..._SparkStructuredStreaming_Batch_10_counters.txt | 26 +
...rDo_SparkStructuredStreaming_Batch_10_times.txt | 26 +
...Do_SparkStructuredStreaming_Batch_200_times.txt | 26 +
.../performance-tests-job-configs/JDBC.txt | 29 +
.../SQLBigQueryIO_Batch_Java.txt | 24 +
.../TFRecordIOIT_HDFS.txt | 23 +
..._Read_Python.txt => biqQueryIO_Read_Python.txt} | 11 +-
...ython.txt => biqQueryIO_Write_Python_Batch.txt} | 11 +-
.../performance-tests-job-configs/cdap.txt | 29 +
.../config_Compressed_TextIOIT.txt | 27 +
.../config_Compressed_TextIOIT_HDFS.txt | 27 +
.../config_ManyFiles_TextIOIT.txt | 29 +
.../config_ManyFiles_TextIOIT_HDFS.txt | 29 +
.../config_ParquetIOIT.txt | 26 +
.../config_ParquetIOIT_HDFS.txt | 26 +
.../config_TFRecordIOIT.txt | 26 +
.../config_XmlIOIT.txt | 27 +
.../config_XmlIOIT_HDFS.txt | 27 +
.../performance-tests-job-configs/hadoopFormat.txt | 29 +
.../kafka_IO_Batch.txt | 26 +
.../kafka_IO_Streaming.txt | 27 +
.../performance-tests-job-configs/mongoDBIO_IT.txt | 28 +
..._Python.txt => pubsubIOIT_Python_Streaming.txt} | 23 +-
...ad_Python.txt => spannerIO_Read_2GB_Python.txt} | 23 +-
...d_Python.txt => spannerIO_Write_2GB_Python.txt} | 23 +-
.../sparkReceiver_IO.txt | 26 +
...ryIO_Read_Python.txt => wordCountIT_Python.txt} | 26 +-
CHANGES.md | 6 -
.../setting-pipeline/python-example/task.py | 17 +-
.../assets/symbols/python.g.yaml | 8 +
.../google-cloud-dataflow-java/arm/build.gradle | 5 +-
.../dataflow/worker/BatchDataflowWorker.java | 95 +-
.../dataflow/worker/DataflowWorkUnitClient.java | 18 +-
.../dataflow/worker/StreamingDataflowWorker.java | 1004 ++------
.../runners/dataflow/worker/WorkUnitClient.java | 6 +-
.../dataflow/worker/counters/NameContext.java | 7 +-
.../dataflow/worker/streaming/ActiveWorkState.java | 292 +++
.../runners/dataflow/worker/streaming/Commit.java | 43 +
.../worker/streaming/ComputationState.java | 139 +
.../dataflow/worker/streaming/ExecutionState.java | 54 +
.../streaming/KeyCommitTooLargeException.java | 50 +
.../dataflow/worker/streaming/ShardedKey.java | 38 +
.../dataflow/worker/streaming/StageInfo.java | 114 +
.../worker/streaming/WeightedBoundedQueue.java | 101 +
.../runners/dataflow/worker/streaming/Work.java | 173 ++
.../dataflow/worker/BatchDataflowWorkerTest.java | 17 +-
.../worker/DataflowWorkUnitClientTest.java | 31 +-
.../worker/StreamingDataflowWorkerTest.java | 1303 +++++-----
.../worker/streaming/ActiveWorkStateTest.java | 296 +++
.../worker/streaming/WeightBoundedQueueTest.java | 194 ++
sdks/go.mod | 4 +-
sdks/go.sum | 8 +-
.../ExpansionServiceSchemaTransformProvider.java | 3 +-
...xpansionServiceSchemaTransformProviderTest.java | 111 +-
sdks/java/io/rrio/build.gradle | 36 +
.../examples/complete/game/user_score.py | 1 -
.../apache_beam/examples/wordcount_it_test.py | 10 +-
sdks/python/apache_beam/internal/gcp/auth.py | 7 +-
sdks/python/apache_beam/io/gcp/bigquery_test.py | 4 +-
sdks/python/apache_beam/io/gcp/bigquery_tools.py | 4 -
.../apache_beam/io/gcp/bigquery_tools_test.py | 5 +-
sdks/python/apache_beam/io/gcp/gcsfilesystem.py | 37 +-
.../io/gcp/gcsfilesystem_integration_test.py | 9 +-
.../apache_beam/io/gcp/gcsfilesystem_test.py | 74 +-
sdks/python/apache_beam/io/gcp/gcsio.py | 641 +++--
.../apache_beam/io/gcp/gcsio_integration_test.py | 188 +-
sdks/python/apache_beam/io/gcp/gcsio_overrides.py | 55 +
sdks/python/apache_beam/io/gcp/gcsio_test.py | 886 +++++--
.../io/gcp/internal/clients/storage/__init__.py | 33 +
.../internal/clients/storage/storage_v1_client.py | 1517 +++++++++++
.../clients/storage/storage_v1_messages.py | 2714 ++++++++++++++++++++
.../options/pipeline_options_validator_test.py | 1 -
.../dataflow_exercise_metrics_pipeline_test.py | 10 -
.../runners/dataflow/internal/apiclient.py | 60 +-
.../apache_beam/runners/interactive/utils.py | 26 +-
.../apache_beam/runners/interactive/utils_test.py | 41 +-
.../runners/portability/sdk_container_builder.py | 41 +-
sdks/python/mypy.ini | 3 +
sdks/python/scripts/run_integration_test.sh | 10 +
sdks/python/setup.py | 1 -
sdks/python/test-suites/dataflow/common.gradle | 20 +
sdks/typescript/package-lock.json | 20 +-
settings.gradle.kts | 1 +
140 files changed, 11061 insertions(+), 3014 deletions(-)
copy .github/workflows/{beam_PerformanceTests_BiqQueryIO_Read_Python.yml =>
beam_LoadTests_Java_Combine_Dataflow_Batch.yml} (51%)
create mode 100644
.github/workflows/beam_LoadTests_Java_Combine_Dataflow_Streaming.yml
create mode 100644
.github/workflows/beam_LoadTests_Java_Combine_SparkStructuredStreaming_Batch.yml
create mode 100644
.github/workflows/beam_LoadTests_Java_ParDo_Dataflow_Batch.yml
create mode 100644
.github/workflows/beam_LoadTests_Java_ParDo_Dataflow_Streaming.yml
create mode 100644
.github/workflows/beam_LoadTests_Java_ParDo_SparkStructuredStreaming_Batch.yml
copy .github/workflows/{beam_PerformanceTests_BiqQueryIO_Read_Python.yml =>
beam_PerformanceTests_Cdap.yml} (55%)
copy .github/workflows/{beam_PerformanceTests_BiqQueryIO_Read_Python.yml =>
beam_PerformanceTests_Compressed_TextIOIT.yml} (60%)
copy .github/workflows/{beam_PerformanceTests_BiqQueryIO_Read_Python.yml =>
beam_PerformanceTests_Compressed_TextIOIT_HDFS.yml} (50%)
copy
.github/workflows/{beam_PerformanceTests_BiqQueryIO_Write_Python_Batch.yml =>
beam_PerformanceTests_HadoopFormat.yml} (54%)
copy .github/workflows/{beam_PerformanceTests_BiqQueryIO_Read_Python.yml =>
beam_PerformanceTests_JDBC.yml} (55%)
create mode 100644 .github/workflows/beam_PerformanceTests_Kafka_IO.yml
copy .github/workflows/{beam_PerformanceTests_BiqQueryIO_Read_Python.yml =>
beam_PerformanceTests_ManyFiles_TextIOIT.yml} (60%)
copy .github/workflows/{beam_PerformanceTests_BiqQueryIO_Read_Python.yml =>
beam_PerformanceTests_ManyFiles_TextIOIT_HDFS.yml} (51%)
copy
.github/workflows/{beam_PerformanceTests_BiqQueryIO_Write_Python_Batch.yml =>
beam_PerformanceTests_MongoDBIO_IT.yml} (54%)
copy .github/workflows/{beam_PerformanceTests_BiqQueryIO_Read_Python.yml =>
beam_PerformanceTests_ParquetIOIT.yml} (60%)
copy
.github/workflows/{beam_PerformanceTests_BiqQueryIO_Write_Python_Batch.yml =>
beam_PerformanceTests_ParquetIOIT_HDFS.yml} (51%)
copy .github/workflows/{beam_PerformanceTests_BiqQueryIO_Read_Python.yml =>
beam_PerformanceTests_PubsubIOIT_Python_Streaming.yml} (62%)
copy .github/workflows/{beam_PerformanceTests_BiqQueryIO_Read_Python.yml =>
beam_PerformanceTests_SQLBigQueryIO_Batch_Java.yml} (57%)
copy
.github/workflows/{beam_PerformanceTests_BiqQueryIO_Write_Python_Batch.yml =>
beam_PerformanceTests_SpannerIO_Read_2GB_Python.yml} (63%)
copy
.github/workflows/{beam_PerformanceTests_BiqQueryIO_Write_Python_Batch.yml =>
beam_PerformanceTests_SpannerIO_Write_2GB_Python_Batch.yml} (63%)
copy .github/workflows/{beam_PerformanceTests_BiqQueryIO_Read_Python.yml =>
beam_PerformanceTests_SparkReceiver_IO.yml} (53%)
copy .github/workflows/{beam_PostCommit_Website_Test.yml =>
beam_PerformanceTests_TFRecordIOIT.yml} (60%)
create mode 100644
.github/workflows/beam_PerformanceTests_TFRecordIOIT_HDFS.yml
copy .github/workflows/{beam_PerformanceTests_BiqQueryIO_Read_Python.yml =>
beam_PerformanceTests_WordCountIT_PythonVersions.yml} (50%)
copy .github/workflows/{beam_PostCommit_Website_Test.yml =>
beam_PerformanceTests_XmlIOIT.yml} (60%)
copy .github/workflows/{beam_PerformanceTests_BiqQueryIO_Read_Python.yml =>
beam_PerformanceTests_XmlIOIT_HDFS.yml} (52%)
copy .github/workflows/{beam_PostCommit_Website_Test.yml =>
beam_PostCommit_Java_Tpcds_Dataflow.yml} (60%)
copy .github/workflows/{beam_PostCommit_Website_Test.yml =>
beam_PostCommit_Java_Tpcds_Flink.yml} (62%)
copy .github/workflows/{beam_PostCommit_Website_Test.yml =>
beam_PostCommit_Java_Tpcds_Spark.yml} (63%)
copy .github/workflows/{beam_PostCommit_Python.yml =>
beam_PostCommit_Python_Arm.yml} (88%)
copy
.github/workflows/{performance-tests-job-configs/config_PerformanceTests_BiqQueryIO_Read_Python.txt
=> load-tests-job-configs/config_Combine_Java_Dataflow_Batch_10b.txt} (51%)
copy
.github/workflows/{performance-tests-job-configs/config_PerformanceTests_BiqQueryIO_Read_Python.txt
=> load-tests-job-configs/config_Combine_Java_Dataflow_Batch_Fanout_4.txt}
(51%)
copy
.github/workflows/{performance-tests-job-configs/config_PerformanceTests_BiqQueryIO_Read_Python.txt
=> load-tests-job-configs/config_Combine_Java_Dataflow_Batch_Fanout_8.txt}
(51%)
create mode 100644
.github/workflows/load-tests-job-configs/java_Combine_Dataflow_Streaming_10b.txt
create mode 100644
.github/workflows/load-tests-job-configs/java_Combine_Dataflow_Streaming_Fanout_4.txt
create mode 100644
.github/workflows/load-tests-job-configs/java_Combine_Dataflow_Streaming_Fanout_8.txt
create mode 100644
.github/workflows/load-tests-job-configs/java_Combine_SparkStructuredStreaming_Batch_10b.txt
create mode 100644
.github/workflows/load-tests-job-configs/java_Combine_SparkStructuredStreaming_Batch_Fanout_4.txt
create mode 100644
.github/workflows/load-tests-job-configs/java_Combine_SparkStructuredStreaming_Batch_Fanout_8.txt
create mode 100644
.github/workflows/load-tests-job-configs/java_ParDo_Dataflow_Batch_100_counters.txt
create mode 100644
.github/workflows/load-tests-job-configs/java_ParDo_Dataflow_Batch_10_counters.txt
create mode 100644
.github/workflows/load-tests-job-configs/java_ParDo_Dataflow_Batch_10_times.txt
create mode 100644
.github/workflows/load-tests-job-configs/java_ParDo_Dataflow_Batch_200_times.txt
create mode 100644
.github/workflows/load-tests-job-configs/java_ParDo_Dataflow_Streaming_100_counters.txt
create mode 100644
.github/workflows/load-tests-job-configs/java_ParDo_Dataflow_Streaming_10_counters.txt
create mode 100644
.github/workflows/load-tests-job-configs/java_ParDo_Dataflow_Streaming_10_times.txt
create mode 100644
.github/workflows/load-tests-job-configs/java_ParDo_Dataflow_Streaming_200_times.txt
create mode 100644
.github/workflows/load-tests-job-configs/java_ParDo_SparkStructuredStreaming_Batch_100_counters.txt
create mode 100644
.github/workflows/load-tests-job-configs/java_ParDo_SparkStructuredStreaming_Batch_10_counters.txt
create mode 100644
.github/workflows/load-tests-job-configs/java_ParDo_SparkStructuredStreaming_Batch_10_times.txt
create mode 100644
.github/workflows/load-tests-job-configs/java_ParDo_SparkStructuredStreaming_Batch_200_times.txt
create mode 100644 .github/workflows/performance-tests-job-configs/JDBC.txt
create mode 100644
.github/workflows/performance-tests-job-configs/SQLBigQueryIO_Batch_Java.txt
create mode 100644
.github/workflows/performance-tests-job-configs/TFRecordIOIT_HDFS.txt
copy
.github/workflows/performance-tests-job-configs/{config_PerformanceTests_BiqQueryIO_Read_Python.txt
=> biqQueryIO_Read_Python.txt} (75%)
rename
.github/workflows/performance-tests-job-configs/{config_PerformanceTests_BiqQueryIO_Write_Python.txt
=> biqQueryIO_Write_Python_Batch.txt} (74%)
create mode 100644 .github/workflows/performance-tests-job-configs/cdap.txt
create mode 100644
.github/workflows/performance-tests-job-configs/config_Compressed_TextIOIT.txt
create mode 100644
.github/workflows/performance-tests-job-configs/config_Compressed_TextIOIT_HDFS.txt
create mode 100644
.github/workflows/performance-tests-job-configs/config_ManyFiles_TextIOIT.txt
create mode 100644
.github/workflows/performance-tests-job-configs/config_ManyFiles_TextIOIT_HDFS.txt
create mode 100644
.github/workflows/performance-tests-job-configs/config_ParquetIOIT.txt
create mode 100644
.github/workflows/performance-tests-job-configs/config_ParquetIOIT_HDFS.txt
create mode 100644
.github/workflows/performance-tests-job-configs/config_TFRecordIOIT.txt
create mode 100644
.github/workflows/performance-tests-job-configs/config_XmlIOIT.txt
create mode 100644
.github/workflows/performance-tests-job-configs/config_XmlIOIT_HDFS.txt
create mode 100644
.github/workflows/performance-tests-job-configs/hadoopFormat.txt
create mode 100644
.github/workflows/performance-tests-job-configs/kafka_IO_Batch.txt
create mode 100644
.github/workflows/performance-tests-job-configs/kafka_IO_Streaming.txt
create mode 100644
.github/workflows/performance-tests-job-configs/mongoDBIO_IT.txt
copy
.github/workflows/performance-tests-job-configs/{config_PerformanceTests_BiqQueryIO_Read_Python.txt
=> pubsubIOIT_Python_Streaming.txt} (58%)
copy
.github/workflows/performance-tests-job-configs/{config_PerformanceTests_BiqQueryIO_Read_Python.txt
=> spannerIO_Read_2GB_Python.txt} (57%)
copy
.github/workflows/performance-tests-job-configs/{config_PerformanceTests_BiqQueryIO_Read_Python.txt
=> spannerIO_Write_2GB_Python.txt} (57%)
create mode 100644
.github/workflows/performance-tests-job-configs/sparkReceiver_IO.txt
rename
.github/workflows/performance-tests-job-configs/{config_PerformanceTests_BiqQueryIO_Read_Python.txt
=> wordCountIT_Python.txt} (52%)
create mode 100644
runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/ActiveWorkState.java
create mode 100644
runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/Commit.java
create mode 100644
runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/ComputationState.java
create mode 100644
runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/ExecutionState.java
create mode 100644
runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/KeyCommitTooLargeException.java
create mode 100644
runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/ShardedKey.java
create mode 100644
runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/StageInfo.java
create mode 100644
runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/WeightedBoundedQueue.java
create mode 100644
runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/Work.java
create mode 100644
runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/streaming/ActiveWorkStateTest.java
create mode 100644
runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/streaming/WeightBoundedQueueTest.java
create mode 100644 sdks/java/io/rrio/build.gradle
create mode 100644 sdks/python/apache_beam/io/gcp/gcsio_overrides.py
create mode 100644
sdks/python/apache_beam/io/gcp/internal/clients/storage/__init__.py
create mode 100644
sdks/python/apache_beam/io/gcp/internal/clients/storage/storage_v1_client.py
create mode 100644
sdks/python/apache_beam/io/gcp/internal/clients/storage/storage_v1_messages.py