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 c8bacb4dabc Update activemq to 5.19.5 (#39593)
add 44d4089a67e Potential fix for environment variable built from
user-controlled sources (#38942)
add 8ded79b7278 Part 1: Log systemName in DataflowWorkUnitClient, Commit,
and core worker states (#39561)
add 8b326b96561 Create span in spanner CDC to start new trace when otel is
enabled. (#39567)
add f2c622b3c03 Fix OpenTelemetry dependencies in published POMs (#39608)
add 24eb1e60d37 Bump github.com/aws/aws-sdk-go-v2/config in /sdks (#39606)
add 945fcfc3f8f Bump github.com/aws/aws-sdk-go-v2/service/s3 in /sdks
(#39609)
add f1125f7c6ca Bump com.gradle.common-custom-user-data-gradle-plugin
(#39602)
add 3a0985a07da [Go SDK] Add GroupIntoBatches transform (#19868) (#38220)
add 77a11347746 Bump github.com/aws/aws-sdk-go-v2/feature/s3/manager in
/sdks (#39605)
add 14c61d13821 Bump zizmorcore/zizmor-action from 0.6.1 to 0.6.2 (#39604)
No new revisions were added by this update.
Summary of changes:
.github/workflows/beam_PreCommit_GHA.yml | 2 +-
.../workflows/beam_Publish_Beam_SDK_Snapshots.yml | 17 +-
CHANGES.md | 10 +
.../dataflow/worker/DataflowWorkUnitClient.java | 12 +-
.../worker/StreamingModeExecutionContext.java | 6 +-
.../beam/runners/dataflow/worker/WindmillSink.java | 14 +-
.../logging/DataflowWorkerLoggingHandler.java | 5 +-
.../worker/logging/DataflowWorkerLoggingMDC.java | 15 +-
.../worker/streaming/ComputationState.java | 4 +
.../streaming/harness/MetricsDataProvider.java | 4 +-
.../harness/StreamingWorkerStatusReporter.java | 2 +-
.../worker/windmill/client/commits/Commit.java | 8 +-
.../work/processing/StreamingWorkScheduler.java | 6 +-
.../worker/DataflowWorkUnitClientTest.java | 8 +-
.../logging/DataflowWorkerLoggingHandlerTest.java | 4 +-
.../worker/testing/RestoreDataflowLoggingMDC.java | 8 +-
.../testing/RestoreDataflowLoggingMDCTest.java | 10 +-
sdks/go.mod | 38 +-
sdks/go.sum | 76 +--
sdks/go/pkg/beam/coder.go | 17 +
sdks/go/pkg/beam/core/graph/coder/coder.go | 116 ++++
sdks/go/pkg/beam/core/graph/coder/coder_test.go | 66 ++
sdks/go/pkg/beam/core/graph/coder/registry.go | 60 +-
.../pkg/beam/core/graph/coder/sharded_key_test.go | 81 +++
sdks/go/pkg/beam/core/runtime/exec/coder.go | 63 ++
sdks/go/pkg/beam/core/runtime/exec/coder_test.go | 84 +++
sdks/go/pkg/beam/core/runtime/graphx/coder.go | 21 +
sdks/go/pkg/beam/core/runtime/symbols.go | 20 +
sdks/go/pkg/beam/core/typex/class.go | 4 +-
sdks/go/pkg/beam/core/typex/fulltype.go | 23 +
sdks/go/pkg/beam/core/typex/special.go | 18 +-
sdks/go/pkg/beam/core/util/reflectx/call.go | 31 +
sdks/go/pkg/beam/pcollection.go | 16 +
sdks/go/pkg/beam/transforms/batch/batch.go | 677 +++++++++++++++++++++
.../pkg/beam/transforms/batch/batch_prism_test.go | 222 +++++++
.../batch/batch_test.go} | 49 +-
sdks/go/pkg/beam/transforms/batch/doc.go | 58 ++
sdks/go/pkg/beam/transforms/batch/size.go | 88 +++
sdks/go/pkg/beam/transforms/batch/size_test.go | 91 +++
sdks/java/io/google-cloud-platform/build.gradle | 1 +
.../changestreams/action/ActionFactory.java | 7 +-
.../action/QueryChangeStreamAction.java | 44 +-
.../dofn/ReadChangeStreamPartitionDoFn.java | 4 +-
.../action/QueryChangeStreamActionTest.java | 10 +-
.../dofn/ReadChangeStreamPartitionDoFnTest.java | 3 +-
sdks/java/io/kafka/build.gradle | 1 +
settings.gradle.kts | 2 +-
47 files changed, 1976 insertions(+), 150 deletions(-)
create mode 100644 sdks/go/pkg/beam/core/graph/coder/sharded_key_test.go
create mode 100644 sdks/go/pkg/beam/transforms/batch/batch.go
create mode 100644 sdks/go/pkg/beam/transforms/batch/batch_prism_test.go
copy sdks/go/pkg/beam/{core/runtime/xlangx/payload_test.go =>
transforms/batch/batch_test.go} (52%)
create mode 100644 sdks/go/pkg/beam/transforms/batch/doc.go
create mode 100644 sdks/go/pkg/beam/transforms/batch/size.go
create mode 100644 sdks/go/pkg/beam/transforms/batch/size_test.go