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 720326829a7 Fill gaps between GcsUtilV2 and GcsUtilV1 (#40244)
add 1e7255931d8 Fix stale ExecutorClassLoader leak and local JAR staging
in portable Spark runner
add 13a3083968f Merge pull request #40270 from Abacn/fix-spark-staging
add e3ff1977173 Tag the Kafka write error output with the schema it
actually emits (#39760)
add 4fbc6a12113 Fix nullness in extensions/avro (#39976)
add b71ebdc33e8 Bump github.com/aws/aws-sdk-go-v2/service/s3 in /sdks
(#40276)
add 5d6918f68da Bump google.golang.org/api from 0.298.0 to 0.299.0 in
/sdks (#40278)
add ea74e85486e Bump github.com/aws/aws-sdk-go-v2/feature/s3/manager in
/sdks (#40277)
add 8ac487a04d8 Add test case verifying metric capture in Setup()
lifecycle method (#39965)
add c9669152cb3 fix flaky python tests (#40228)
add 6317b1942f2 Reduce severity for state requests cancelled by runner
(#40254)
add b6c8937bbf1 Merge pull request #40274 from
reuvenlax/fix_flaky_schema_test
add dd89aa1e4e3 Merge pull request #40269 from
reuvenlax/schema_update_fixups
add 011e1b33e24 Fix pubsubio import of removed genproto package (fixes
#40018) (#40059)
add de8d83c107d Merge pull request #40259 from
reuvenlax/improve_gcp_presubmit_speed
No new revisions were added by this update.
Summary of changes:
...Commit_Java_PVR_Spark4_StructuredStreaming.json | 3 +-
.../beam_PreCommit_Java_GCP_IO_Direct.yml | 48 ++---
CHANGES.md | 2 +
.../beam/model/fn_execution/v1/beam_fn_api.proto | 13 ++
.../beam/runners/jobsubmission/JobInvoker.java | 15 +-
.../runners/spark/SparkCommonPipelineOptions.java | 7 +-
.../beam/runners/spark/SparkPipelineRunner.java | 8 +
.../spark/translation/SparkContextFactory.java | 12 +-
.../runners/spark/SparkPipelineOptionsTest.java | 17 ++
sdks/go.mod | 22 +--
sdks/go.sum | 36 ++--
.../beam/core/runtime/exec/setup_metrics_test.go | 215 +++++++++++++++++++++
sdks/go/pkg/beam/io/pubsubio/pubsubio.go | 2 +-
.../beam/sdk/extensions/avro/coders/AvroCoder.java | 45 +++--
.../sdk/extensions/avro/io/AvroDatumFactory.java | 12 +-
.../apache/beam/sdk/extensions/avro/io/AvroIO.java | 133 +++++++------
.../extensions/avro/io/AvroSchemaIOProvider.java | 18 +-
.../beam/sdk/extensions/avro/io/AvroSink.java | 26 +--
.../beam/sdk/extensions/avro/io/AvroSource.java | 87 ++++++---
.../avro/io/ConstantAvroDestination.java | 5 +-
.../avro/io/SerializableAvroCodecFactory.java | 24 +--
.../avro/schemas/utils/AvroByteBuddyUtils.java | 5 +-
.../extensions/avro/schemas/utils/AvroUtils.java | 37 ++--
.../avro/schemas/utils/AvroUtilsTest.java | 83 +++++++-
sdks/java/io/google-cloud-platform/build.gradle | 128 +++++++-----
.../beam/sdk/io/gcp/bigquery/BigQueryOptions.java | 2 +-
.../sdk/io/gcp/bigquery/CreateTableHelpers.java | 16 +-
.../io/gcp/bigquery/StorageApiWritePayload.java | 3 +-
.../bigquery/StorageApiWriteUnshardedRecords.java | 82 +++++---
.../bigquery/StorageApiWritesShardedRecords.java | 21 +-
.../sdk/io/gcp/testing/FakeDatasetService.java | 38 +++-
.../sdk/io/gcp/bigquery/AppendRowsPacketTest.java | 6 +-
.../BigQueryIOWriteStorageApiBatchTest.java | 22 +--
.../BigQueryIOWriteStorageApiStreamTest.java} | 29 ++-
.../sdk/io/gcp/bigquery/BigQueryIOWriteTest.java | 67 ++++++-
.../SchemaChangeDetectorHelperBufferingTest.java | 17 +-
.../bigquery/SchemaChangeDetectorHelperTest.java | 46 ++---
.../StorageApiSchemaMismatchDrainTest.java | 8 +-
.../bigquery/StorageApiSinkSchemaUpdateITBase.java | 4 +-
.../dofn/ReadChangeStreamPartitionDoFnTest.java | 14 +-
.../beam/sdk/io/gcp/spanner/SpannerReadIT.java | 52 ++---
.../beam/sdk/io/gcp/spanner/SpannerWriteIT.java | 57 ++++--
.../changestreams/it/IntegrationTestEnv.java | 80 ++++----
.../kafka/KafkaWriteSchemaTransformProvider.java | 3 +-
.../KafkaWriteSchemaTransformProviderTest.java | 24 +++
sdks/python/apache_beam/runners/common.py | 7 +
.../apache_beam/runners/worker/sdk_worker.py | 12 ++
.../apache_beam/runners/worker/sdk_worker_test.py | 28 +++
.../python/apache_beam/yaml/yaml_transform_test.py | 37 ++--
49 files changed, 1185 insertions(+), 493 deletions(-)
create mode 100644 sdks/go/pkg/beam/core/runtime/exec/setup_metrics_test.go
copy
runners/spark/src/test/java/org/apache/beam/runners/spark/TestSparkPipelineOptionsRegistrar.java
=>
sdks/java/io/google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/bigquery/BigQueryIOWriteStorageApiBatchTest.java
(58%)
copy
sdks/java/{extensions/avro/src/test/java/org/apache/beam/sdk/extensions/avro/AvroVersionVerificationTest.java
=>
io/google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/bigquery/BigQueryIOWriteStorageApiStreamTest.java}
(56%)