This is an automated email from the ASF dual-hosted git repository.
github-actions[bot] pushed a change to branch nightly-refs/heads/master
in repository https://gitbox.apache.org/repos/asf/beam.git
from 67018be2650 [runners-spark] Support stateful ParDo in the Structured
Streaming batch runner (#39793)
add 2d94b11acaf normalize KinesisIO
add 02060be7b7a sync Kinesis
add 661862a4ba2 Add LocalStack Kinesis YAML integration test
add 2ff4f0eda94 Merge pull request #39763 from
aIbrahiim/yaml-normalize-kinesis
add 7420ad22aea Update google-cloud-bigtable
add 11e9f8b5598 Merge pull request #39837 from apache/update-bigtable
add e626f54690b Implement Vertex AI Model Monitoring v2 (#39738)
add 3b9e647dec1 [Prism] Schedule consumers of a self checkpointing source
(#39572)
add 88b3ee7b488 JmsIO yaml (#39818)
add 4e5ae91554e Support lakehouse PCNT format in BigQueryIO storage read.
(#39597)
add 6434f7411e7 [Docs] Document UnboundedSource in the Python I/O
connector guide (#39529)
add 61ed38fd5fd Support Secret Manager in JdbcIO for Java, Python and YAML
(#39834)
add 7de4d3ab73b Fix KafkaIO wrtie SchemaTransform parallelism (#39844)
add 015e823a1a9 [GSoC-273] Feat: Integrate TestPubsubContext to prevent
Pub/Sub resource leaks and expand stale cleaner scope (#39826)
add c1017953ac3 Adds support for reading at a given Delta Lake version or
timestamp (#39758)
No new revisions were added by this update.
Summary of changes:
.../IO_Iceberg_Integration_Tests.json | 2 +-
.../beam_PostCommit_Java_Delta_IO_Dataflow.json | 2 +-
.github/trigger_files/beam_PostCommit_Python.json | 2 +-
.../beam_PostCommit_Python_Xlang_Gcp_Direct.json | 2 +-
...m_PostCommit_Python_Xlang_Messaging_Direct.json | 2 +-
.../beam_PostCommit_Yaml_Xlang_Direct.json | 2 +-
CHANGES.md | 1 +
runners/google-cloud-dataflow-java/build.gradle | 2 +
.../prism/internal/engine/elementmanager.go | 91 ++-
.../engine/elementmanager_continuation_test.go | 299 ++++++++++
.../schemaio-expansion-service/build.gradle | 6 +
.../KinesisReadSchemaTransformProvider.java | 315 ++++++++++
.../KinesisWriteSchemaTransformProvider.java | 280 +++++++++
.../KinesisSchemaTransformProviderTest.java | 221 +++++++
.../beam/sdk/io/delta/CreateReadTasksDoFn.java | 23 +-
.../java/org/apache/beam/sdk/io/delta/DeltaIO.java | 26 +-
.../io/delta/DeltaReadSchemaTransformProvider.java | 6 +-
.../org/apache/beam/sdk/io/delta/DeltaIOIT.java | 141 +++--
.../org/apache/beam/sdk/io/delta/DeltaIOTest.java | 102 ++++
.../DeltaReadSchemaTransformProviderTest.java | 51 ++
.../beam/sdk/io/delta/DeltaWriteTestUtils.java | 30 +
sdks/java/io/google-cloud-platform/build.gradle | 33 ++
.../beam/sdk/io/gcp/bigquery/BigQueryHelpers.java | 100 +++-
.../beam/sdk/io/gcp/bigquery/BigQueryIO.java | 13 +
.../io/gcp/bigquery/BigQueryStorageSourceBase.java | 7 +-
.../gcp/bigquery/BigQueryStorageTableSource.java | 3 +-
.../sdk/io/gcp/bigquery/BigQueryTableSource.java | 5 +
.../sdk/io/gcp/bigquery/BigQueryHelpersTest.java | 348 ++++++++++-
.../bigquery/BigQueryIOIcebergManagedTableIT.java | 408 +++++++++++++
.../io/gcp/bigquery/BigQueryIOStorageReadTest.java | 63 ++
sdks/java/io/iceberg/build.gradle | 11 +
.../org/apache/beam/sdk/io/iceberg/AddFilesIT.java | 71 ++-
.../iceberg/BigQueryManagedTableCrossEngineIT.java | 183 ++++++
.../catalog/BigQueryMetastoreCatalogIT.java | 10 +
.../io/iceberg/catalog/IcebergCatalogBaseIT.java | 259 ++++++++
.../sdk/io/iceberg/catalog/RESTCatalogBLMSIT.java | 21 +-
.../java/org/apache/beam/sdk/io/jdbc/JdbcIO.java | 70 +++
.../io/jdbc/JdbcReadSchemaTransformProvider.java | 34 +-
.../beam/sdk/io/jdbc/JdbcSchemaIOProvider.java | 6 +
.../io/jdbc/JdbcWriteSchemaTransformProvider.java | 34 +-
.../org/apache/beam/sdk/io/jdbc/JdbcIOTest.java | 45 ++
.../kafka/KafkaWriteSchemaTransformProvider.java | 41 +-
.../KafkaWriteSchemaTransformProviderTest.java | 27 +-
.../streaming_wordcount_debugging_it_test.py | 12 +-
.../examples/streaming_wordcount_it_test.py | 12 +-
.../io/external/xlang_jdbcio_it_test.py | 113 ++++
.../apache_beam/io/external/xlang_jmsio_it_test.py | 13 +-
.../apache_beam/io/gcp/pubsub_integration_test.py | 17 +-
sdks/python/apache_beam/io/jdbc.py | 48 +-
sdks/python/apache_beam/ml/inference/base.py | 28 +
.../ml/inference/vertex_ai_model_monitoring_v2.py | 458 ++++++++++++++
.../vertex_ai_model_monitoring_v2_it_test.py | 485 +++++++++++++++
.../vertex_ai_model_monitoring_v2_test.py | 655 +++++++++++++++++++++
sdks/python/apache_beam/testing/README.md | 109 ++++
.../{jdbc.yaml => jdbc_secret_manager.yaml} | 22 +-
.../messaging/{kafka.yaml => kinesis.yaml} | 52 +-
sdks/python/apache_beam/yaml/integration_tests.py | 269 +++++++++
sdks/python/apache_beam/yaml/standard_io.yaml | 104 ++++
.../databases/jdbc.yaml => tests/ibm_mq.yaml} | 47 +-
.../databases/jdbc.yaml => tests/jms.yaml} | 41 +-
sdks/python/apache_beam/yaml/yaml_provider.py | 6 +-
sdks/python/build.gradle | 2 +
.../container/ml/py310/base_image_requirements.txt | 2 +-
.../container/ml/py310/gpu_image_requirements.txt | 2 +-
.../container/ml/py311/base_image_requirements.txt | 2 +-
.../container/ml/py311/gpu_image_requirements.txt | 2 +-
.../container/ml/py312/base_image_requirements.txt | 2 +-
.../container/ml/py312/gpu_image_requirements.txt | 2 +-
.../container/ml/py313/base_image_requirements.txt | 2 +-
.../container/py310/base_image_requirements.txt | 2 +-
.../container/py311/base_image_requirements.txt | 2 +-
.../container/py312/base_image_requirements.txt | 2 +-
.../container/py313/base_image_requirements.txt | 2 +-
.../container/py314/base_image_requirements.txt | 2 +-
sdks/standard_external_transforms.yaml | 134 ++++-
.../en/documentation/io/developing-io-python.md | 239 +++++++-
.../site/content/en/documentation/io/managed-io.md | 4 +-
77 files changed, 6027 insertions(+), 265 deletions(-)
create mode 100644
sdks/go/pkg/beam/runners/prism/internal/engine/elementmanager_continuation_test.go
create mode 100644
sdks/java/io/amazon-web-services2/src/main/java/org/apache/beam/sdk/io/aws2/kinesis/KinesisReadSchemaTransformProvider.java
create mode 100644
sdks/java/io/amazon-web-services2/src/main/java/org/apache/beam/sdk/io/aws2/kinesis/KinesisWriteSchemaTransformProvider.java
create mode 100644
sdks/java/io/amazon-web-services2/src/test/java/org/apache/beam/sdk/io/aws2/kinesis/KinesisSchemaTransformProviderTest.java
create mode 100644
sdks/java/io/google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/bigquery/BigQueryIOIcebergManagedTableIT.java
create mode 100644
sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/BigQueryManagedTableCrossEngineIT.java
create mode 100644
sdks/python/apache_beam/ml/inference/vertex_ai_model_monitoring_v2.py
create mode 100644
sdks/python/apache_beam/ml/inference/vertex_ai_model_monitoring_v2_it_test.py
create mode 100644
sdks/python/apache_beam/ml/inference/vertex_ai_model_monitoring_v2_test.py
create mode 100644 sdks/python/apache_beam/testing/README.md
copy sdks/python/apache_beam/yaml/extended_tests/databases/{jdbc.yaml =>
jdbc_secret_manager.yaml} (69%)
copy sdks/python/apache_beam/yaml/extended_tests/messaging/{kafka.yaml =>
kinesis.yaml} (59%)
copy sdks/python/apache_beam/yaml/{extended_tests/databases/jdbc.yaml =>
tests/ibm_mq.yaml} (51%)
copy sdks/python/apache_beam/yaml/{extended_tests/databases/jdbc.yaml =>
tests/jms.yaml} (57%)