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 d3ccb3a9a75 upgrade grpc to 1.82.1 for test infra and playground
(#39915)
add fc26dacad45 Bump transformers (#39963)
add b0a14f61c2c Bump google.golang.org/api from 0.294.0 to 0.295.0 in
/sdks (#39967)
add 42129fd3154 Upgrade GCP Libraries BOM to 26.87.0 (#39887)
add 8abe4a39a4a Accept integral JSON values that a double represents
exactly (#39744)
add 8b431c8ed2b Update Dataflow Python container (#39958)
add 7151bed2bbd Bump github.com/apache/thrift from 0.23.0 to 0.24.0 in
/sdks (#39974)
add 10408793b66 AddFiles: read side of the schema pre-pass
(ReadFooterSchema, FileSchemas.canonical, CollectDistinctSchemas) (#39933)
add 3d265d2b079 Bump golang.org/x/crypto from 0.50.0 to 0.52.0 in
/playground/backend (#39966)
add c979b5cfbe2 [Flink] Select bounded-source split assignment by
estimated size (#39874)
add 32d5037d336 [Go] Add parsing for cpu_count resource hint (#39977)
add 1d652b557df Bump google.golang.org/grpc from 1.49.0 to 1.82.1 in
/learning/tour-of-beam/backend (#39480)
add da6355fbd18 Bump golang.org/x/crypto in /learning/tour-of-beam/backend
(#39983)
add 4391d8d7f38 Fix GCS glob matching for multi-byte characters (#39972)
add ce18612133a Bump github.com/cloudevents/sdk-go/v2 in
/learning/tour-of-beam/backend (#39985)
add b3a5df21e1c Bump golang.org/x/net in /learning/tour-of-beam/backend
(#39984)
No new revisions were added by this update.
Summary of changes:
.../workflows/tour_of_beam_backend_integration.yml | 2 +-
CHANGES.md | 2 +
.../org/apache/beam/gradle/BeamModulePlugin.groovy | 14 +-
.../beam/it/gcp/pubsub/PubsubResourceManager.java | 6 +-
learning/tour-of-beam/backend/function.go | 2 +-
learning/tour-of-beam/backend/go.mod | 76 ++-
learning/tour-of-beam/backend/go.sum | 411 +++++++---------
.../backend/integration_tests/client_pg.go | 6 +-
playground/backend/go.mod | 10 +-
playground/backend/go.sum | 24 +-
.../beam/runners/flink/FlinkPipelineOptions.java | 11 +
.../beam/runners/flink/FlinkPipelineOptions.java | 11 +
.../wrappers/streaming/io/source/FlinkSource.java | 64 +--
.../io/source/FlinkSourceEnumeratorState.java} | 32 +-
.../FlinkSourceEnumeratorStateSerializer.java | 103 ++++
.../io/source/FlinkSourceSplitEnumerator.java | 141 +++---
.../streaming/io/source/FlinkSourceSplitUtils.java | 84 ++++
.../io/source/LazyFlinkSourceSplitEnumerator.java | 184 ++++----
.../SizeBasedFlinkSourceSplitEnumerator.java | 205 ++++++++
.../io/source/FlinkSourceSplitEnumeratorTest.java | 522 +++++++++++++++++----
sdks/go.mod | 4 +-
sdks/go.sum | 8 +-
sdks/go/pkg/beam/io/filesystem/gcs/gcs.go | 22 +-
sdks/go/pkg/beam/io/filesystem/gcs/gcs_test.go | 8 +
sdks/go/pkg/beam/options/jobopts/options.go | 2 +
sdks/go/pkg/beam/options/jobopts/options_test.go | 4 +-
sdks/go/pkg/beam/options/resource/hint.go | 17 +-
sdks/go/pkg/beam/options/resource/hint_test.go | 63 +++
.../container/license_scripts/dep_urls_java.yaml | 2 +-
.../beam/sdk/util/RowJsonValueExtractors.java | 22 +-
.../java/org/apache/beam/sdk/util/RowJsonTest.java | 44 ++
.../changestreams/dao/MetadataTableDao.java | 26 +-
.../changestreams/dao/MetadataTableDaoTest.java | 2 +
.../changestreams/dofn/InitializeDoFnTest.java | 30 +-
.../changestreams/it/BigtableChangeStreamIT.java | 33 +-
.../sdk/io/iceberg/CollectDistinctSchemas.java | 103 ++++
.../apache/beam/sdk/io/iceberg/FileSchemas.java | 98 ++++
.../beam/sdk/io/iceberg/ReadFooterSchema.java | 181 +++++++
.../sdk/io/iceberg/CollectDistinctSchemasTest.java | 131 ++++++
.../beam/sdk/io/iceberg/FileSchemasTest.java | 131 ++++++
.../beam/sdk/io/iceberg/ReadFooterSchemaTest.java | 279 +++++++++++
.../inference/runinference_metrics/setup.py | 2 +-
.../apache_beam/runners/dataflow/internal/names.py | 2 +-
.../shortcodes/flink_java_pipeline_options.html | 5 +
.../shortcodes/flink_python_pipeline_options.html | 5 +
45 files changed, 2514 insertions(+), 620 deletions(-)
copy
runners/flink/src/{test/java/org/apache/beam/runners/flink/translation/wrappers/streaming/io/source/TestSource.java
=>
main/java/org/apache/beam/runners/flink/translation/wrappers/streaming/io/source/FlinkSourceEnumeratorState.java}
(50%)
create mode 100644
runners/flink/src/main/java/org/apache/beam/runners/flink/translation/wrappers/streaming/io/source/FlinkSourceEnumeratorStateSerializer.java
create mode 100644
runners/flink/src/main/java/org/apache/beam/runners/flink/translation/wrappers/streaming/io/source/FlinkSourceSplitUtils.java
create mode 100644
runners/flink/src/main/java/org/apache/beam/runners/flink/translation/wrappers/streaming/io/source/SizeBasedFlinkSourceSplitEnumerator.java
create mode 100644
sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/CollectDistinctSchemas.java
create mode 100644
sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/FileSchemas.java
create mode 100644
sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/ReadFooterSchema.java
create mode 100644
sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/CollectDistinctSchemasTest.java
create mode 100644
sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/FileSchemasTest.java
create mode 100644
sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/ReadFooterSchemaTest.java