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 51966762894 update containers (#39575)
add e4779cf79f2 Fix flaky FileIOTest.testMatchWatchForNewFiles test under
CI filesystems (#38047)
add 7f96ee4f2bf Bump google.golang.org/grpc from 1.82.1 to 1.83.0 in /sdks
(#39585)
add 93af6786bbb Bump github/codeql-action from 4.37.3 to 4.37.4 (#39586)
add 8b4c6751963 Support core dump analysis with pystack and gdb. (#39484)
add a72451d8072 Bump github.com/nats-io/nats-server/v2 from 2.14.3 to
2.14.4 in /sdks (#39584)
add 3a02af8146c [Docs] Update Flink version references on the Flink runner
page (#39212)
add 539b048eb8d Fix dataframe CSV tests on Windows (#39563)
add 4d3e1f017e5 Support array-valued schema options in Python (#39583)
add f0734081aca Feat: new cleaning rule to orphaned subscriptions (#39538)
add db57e4ae351 [Docs] Add CHANGES entries for Python UnboundedSource and
Watch (#39579)
add a6b3399e41d Fix internal test failure after #39487 (#39591)
add 52f6e46bede Add query_output_schema to ReadFromBigQuery for BEAM_ROW +
query support (#39160)
add f3e12fee1b2 [DebeziumIO] Upgrade to Debezium 3.5.2.Final (#39569)
add 82c6ee4e976 [Python] Bound Watch state with a timestamp cursor (#39090)
add 789e1d1480c Redistribute - trace propagation (#39590)
add 83821eb2fd1 update containers (#39596)
add a89b9f7120e remove gsutil usage (#39448)
add c4b4bda2a5d [Iceberg CDC] Add Changelog readers and update resolver
(#38837)
add c8bacb4dabc Update activemq to 5.19.5 (#39593)
No new revisions were added by this update.
Summary of changes:
.github/workflows/beam_PreCommit_GHA.yml | 16 +
.github/workflows/build_wheels.yml | 12 +-
.github/workflows/codeql.yml | 4 +-
.../workflows/run_rc_validation_go_wordcount.yml | 12 +-
.../run_rc_validation_python_mobile_gaming.yml | 6 +-
.../workflows/run_rc_validation_python_yaml.yml | 6 +-
.test-infra/dataproc/flink_cluster.sh | 2 +-
.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 +-
.test-infra/tools/stale_cleaner.py | 15 +-
.test-infra/tools/test_stale_cleaner.py | 61 ++-
CHANGES.md | 6 +-
.../org/apache/beam/gradle/BeamModulePlugin.groovy | 2 +-
.../beam/examples/complete/game/UserScore.java | 2 +-
examples/multi-language/README.md | 6 +-
.../beam-ml/automatic_model_refresh.ipynb | 4 +-
.../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 +-
.../main/groovy/mobilegaming-java-dataflow.groovy | 12 +-
.../groovy/mobilegaming-java-dataflowbom.groovy | 12 +-
.../main/groovy/quickstart-java-dataflow.groovy | 8 +-
.../python_release_automation_utils.sh | 6 +-
.../run_release_candidate_python_quickstart.sh | 6 +-
.../dataflow/worker/WindmillKeyedWorkItem.java | 12 +-
sdks/go.mod | 14 +-
sdks/go.sum | 28 +-
sdks/go/README.md | 2 +-
.../test/integration/io/xlang/debezium/debezium.go | 2 +-
.../integration/io/xlang/debezium/debezium_test.go | 2 +-
.../apache/beam/sdk/transforms/Redistribute.java | 23 +-
.../java/org/apache/beam/sdk/transforms/Reify.java | 4 +-
.../java/org/apache/beam/sdk/io/FileIOTest.java | 22 +-
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 | 31 +-
.../beam/io/debezium/KafkaSourceConsumerFn.java | 5 +
.../io/debezium/DebeziumIOMySqlConnectorIT.java | 6 +-
.../debezium/DebeziumIOPostgresSqlConnectorIT.java | 4 +-
.../apache/beam/io/debezium/DebeziumIOTest.java | 3 +-
.../debezium/DebeziumReadSchemaTransformTest.java | 29 +-
.../beam/sdk/io/iceberg/IcebergScanConfig.java | 107 ++++-
.../apache/beam/sdk/io/iceberg/IcebergUtils.java | 244 +++++-----
.../beam/sdk/io/iceberg/cdc/CdcOutputUtils.java | 6 +-
.../beam/sdk/io/iceberg/cdc/CdcReadUtils.java | 2 +
.../beam/sdk/io/iceberg/cdc/CdcResolver.java | 180 ++++++++
...ngelogDescriptor.java => CdcRowDescriptor.java} | 57 +--
.../beam/sdk/io/iceberg/cdc/ChangelogScanner.java | 13 +-
.../beam/sdk/io/iceberg/cdc/LocalResolveDoFn.java | 245 ++++++++++
.../beam/sdk/io/iceberg/cdc/OverlapRange.java | 102 +++++
.../sdk/io/iceberg/cdc/ReadFromChangelogs.java | 494 +++++++++++++++++++++
.../beam/sdk/io/iceberg/IcebergUtilsTest.java | 10 +-
.../beam/sdk/io/iceberg/cdc/CdcResolverTest.java | 156 +++++++
.../sdk/io/iceberg/cdc/ChangelogScannerTest.java | 20 +
.../sdk/io/iceberg/cdc/LocalResolveDoFnTest.java | 340 ++++++++++++++
.../beam/sdk/io/iceberg/cdc/OverlapRangeTest.java | 161 +++++++
.../sdk/io/iceberg/cdc/ReadFromChangelogsTest.java | 366 +++++++++++++++
sdks/python/apache_beam/dataframe/io.py | 2 +-
sdks/python/apache_beam/dataframe/io_test.py | 15 +-
.../io/external/xlang_debeziumio_it_test.py | 2 +-
sdks/python/apache_beam/io/gcp/bigquery.py | 26 +-
.../io/gcp/bigquery_schema_tools_test.py | 74 ++-
sdks/python/apache_beam/io/gcp/bigquery_test.py | 49 ++
sdks/python/apache_beam/io/watch.py | 264 ++++++++---
sdks/python/apache_beam/io/watch_test.py | 387 +++++++++++++++-
.../python/apache_beam/options/pipeline_options.py | 4 +
.../apache_beam/options/pipeline_options_test.py | 9 +
.../apache_beam/runners/dataflow/internal/names.py | 2 +-
.../testing/benchmarks/chicago_taxi/run_chicago.sh | 6 +-
sdks/python/apache_beam/typehints/schemas.py | 52 ++-
sdks/python/apache_beam/typehints/schemas_test.py | 16 +
sdks/python/apache_beam/yaml/yaml_io.py | 14 +-
sdks/python/apache_beam/yaml/yaml_io_test.py | 43 ++
.../container/base_image_requirements_manual.txt | 1 +
sdks/python/container/boot.go | 55 ++-
.../container/ml/py310/base_image_requirements.txt | 1 +
.../container/ml/py310/gpu_image_requirements.txt | 1 +
.../container/ml/py311/base_image_requirements.txt | 1 +
.../container/ml/py311/gpu_image_requirements.txt | 1 +
.../container/ml/py312/base_image_requirements.txt | 1 +
.../container/ml/py312/gpu_image_requirements.txt | 1 +
.../container/ml/py313/base_image_requirements.txt | 1 +
sdks/python/container/profiler.go | 303 +++++++++++--
sdks/python/container/profiler_test.go | 125 ++++++
.../container/py310/base_image_requirements.txt | 1 +
.../container/py311/base_image_requirements.txt | 1 +
.../container/py312/base_image_requirements.txt | 1 +
.../container/py313/base_image_requirements.txt | 1 +
.../container/py314/base_image_requirements.txt | 1 +
sdks/python/scripts/run_snapshot_publish.sh | 4 +-
website/Dockerfile | 2 +-
website/build.gradle | 5 +-
.../content/en/blog/apache-hop-with-dataflow.md | 12 +-
.../content/en/blog/beam-sql-with-notebooks.md | 4 +-
.../site/content/en/documentation/runners/flink.md | 23 +-
.../site/content/en/documentation/runners/spark.md | 4 +-
.../sdks/python-multi-language-pipelines.md | 2 +-
.../site/content/en/get-started/quickstart-java.md | 4 +-
105 files changed, 3981 insertions(+), 545 deletions(-)
delete mode 100644 .test-infra/metrics/influxdb/gsutil/.boto
delete mode 100644 .test-infra/metrics/influxdb/gsutil/Dockerfile
create mode 100644
sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/cdc/CdcResolver.java
copy
sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/cdc/{ChangelogDescriptor.java
=> CdcRowDescriptor.java} (59%)
create mode 100644
sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/cdc/LocalResolveDoFn.java
create mode 100644
sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/cdc/OverlapRange.java
create mode 100644
sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/cdc/ReadFromChangelogs.java
create mode 100644
sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/cdc/CdcResolverTest.java
create mode 100644
sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/cdc/LocalResolveDoFnTest.java
create mode 100644
sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/cdc/OverlapRangeTest.java
create mode 100644
sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/cdc/ReadFromChangelogsTest.java
create mode 100644 sdks/python/container/profiler_test.go