This is an automated email from the ASF dual-hosted git repository.
github-bot pushed a change to branch nightly-refs/heads/master
in repository https://gitbox.apache.org/repos/asf/beam.git
from 428ec97e30c improve error message for mismatched pipelines (#24834)
add 7ad44c84585 Handle schema updates in Storage API writes.
add f5020e7ac2b Merge pull request #24145: Handle updates to table schema
when using Storage API writes.
add 9cecec3f74c Bump cloud.google.com/go/storage from 1.28.1 to 1.29.0 in
/sdks (#25095)
add 5e1ebee8b47 Allow to set timeout for finishing a remote bundle in
Samza portable runner (#25031)
add e379c23c885 Fix truncate copy job when WRITE_TRUNCATE in BigQuery
batch load (#25101)
add 482401411b7 fix(sec): upgrade torch to 1.13.1 (#24933)
add 4dad3c696d4 [#24515] Delete the JRH (#24967)
add cd20288318d Support DoFn metrics in portable Samza Runner (#25068)
No new revisions were added by this update.
Summary of changes:
...mit_Java_ValidatesRunner_Dataflow_Java11.groovy | 2 +-
...mit_Java_ValidatesRunner_Dataflow_Java17.groovy | 2 +-
CHANGES.md | 4 +-
build.gradle.kts | 3 +-
.../org/apache/beam/gradle/BeamModulePlugin.groovy | 6 +-
runners/google-cloud-dataflow-java/build.gradle | 13 +-
.../examples-streaming/build.gradle | 10 +-
.../examples/build.gradle | 17 +-
.../google-cloud-dataflow-java/worker/build.gradle | 285 +++++---
.../worker/legacy-worker/build.gradle | 276 -------
.../dataflow/worker/DataflowRunnerHarness.java | 25 -
.../dataflow/harness/util/package-info.java | 24 -
.../dataflow/worker/BatchDataflowWorker.java | 107 +--
.../worker/BeamFnMapTaskExecutorFactory.java | 799 ---------------------
.../worker/DataflowBatchWorkerHarness.java | 7 +
.../worker/DataflowMapTaskExecutorFactory.java | 9 -
.../dataflow/worker/DataflowRunnerHarness.java | 248 -------
...FetchAndFilterStreamingSideInputsOperation.java | 280 --------
.../dataflow/worker/FnApiWindowMappingFn.java | 314 --------
.../worker/GroupAlsoByWindowParDoFnFactory.java | 21 +-
.../worker/IntrinsicMapTaskExecutorFactory.java | 19 -
.../worker/NoOpSourceOperationExecutor.java | 68 --
.../dataflow/worker/SdkHarnessRegistries.java | 286 --------
.../dataflow/worker/SdkHarnessRegistry.java | 68 --
.../worker/SourceOperationExecutorFactory.java | 26 +-
.../dataflow/worker/StreamingDataflowWorker.java | 213 +-----
.../worker/counters/CounterUpdateAggregator.java | 38 -
.../worker/counters/CounterUpdateAggregators.java | 75 --
.../DistributionCounterUpdateAggregator.java | 68 --
.../counters/MeanCounterUpdateAggregator.java | 58 --
.../counters/SumCounterUpdateAggregator.java | 50 --
.../dataflow/worker/fn/BeamFnControlService.java | 125 ----
.../worker/fn/control/BeamFnMapTaskExecutor.java | 646 -----------------
.../control/DataflowSideInputHandlerFactory.java | 163 -----
...ntMonitoringInfoToCounterUpdateTransformer.java | 131 ----
...meMonitoringInfoToCounterUpdateTransformer.java | 148 ----
...piMonitoringInfoToCounterUpdateTransformer.java | 89 ---
...ntMonitoringInfoToCounterUpdateTransformer.java | 136 ----
.../MonitoringInfoToCounterUpdateTransformer.java | 35 -
.../fn/control/ProcessRemoteBundleOperation.java | 165 -----
.../control/RegisterAndProcessBundleOperation.java | 693 ------------------
...onMonitoringInfoToCounterUpdateTransformer.java | 144 ----
...erMonitoringInfoToCounterUpdateTransformer.java | 137 ----
.../worker/fn/data/BeamFnDataGrpcService.java | 256 -------
.../fn/data/RemoteGrpcPortReadOperation.java | 118 ---
.../fn/data/RemoteGrpcPortWriteOperation.java | 255 -------
.../dataflow/worker/fn/grpc/BeamFnService.java | 32 -
.../worker/fn/logging/BeamFnLoggingService.java | 155 ----
.../fn/stream/ServerStreamObserverFactory.java | 103 ---
.../graph/CloneAmbiguousFlattensFunction.java | 149 ----
.../graph/CreateExecutableStageNodeFunction.java | 602 ----------------
.../graph/CreateRegisterFnOperationFunction.java | 318 --------
.../graph/DeduceFlattenLocationsFunction.java | 328 ---------
.../worker/graph/DeduceNodeLocationsFunction.java | 124 ----
...nsertFetchAndFilterStreamingSideInputNodes.java | 176 -----
.../worker/graph/LengthPrefixUnknownCoders.java | 64 --
.../beam/runners/dataflow/worker/graph/Nodes.java | 114 ---
.../worker/graph/RegisterNodeFunction.java | 622 ----------------
.../graph/RemoveFlattenInstructionsFunction.java | 83 ---
.../graph/ReplacePgbkWithPrecombineFunction.java | 87 ---
.../dataflow/worker/BatchDataflowWorkerTest.java | 24 +-
.../worker/DataflowBatchWorkerHarnessTest.java | 21 +-
.../dataflow/worker/FnApiWindowMappingFnTest.java | 182 -----
.../IntrinsicMapTaskExecutorFactoryTest.java | 8 -
.../worker/NoOpSourceOperationExecutorTest.java | 61 --
.../dataflow/worker/SdkHarnessRegistryTest.java | 122 ----
.../worker/SourceOperationExecutorFactoryTest.java | 16 -
.../worker/StreamingDataflowWorkerTest.java | 2 -
.../counters/CounterUpdateAggregatorsTest.java | 96 ---
.../DistributionCounterUpdateAggregatorTest.java | 72 --
.../counters/MeanCounterUpdateAggregatorTest.java | 66 --
.../counters/SumCounterUpdateAggregatorTest.java | 62 --
.../worker/fn/BeamFnControlServiceTest.java | 174 -----
.../fn/control/BeamFnMapTaskExecutorTest.java | 295 --------
.../DataflowSideInputHandlerFactoryTest.java | 173 -----
...nitoringInfoToCounterUpdateTransformerTest.java | 124 ----
...nitoringInfoToCounterUpdateTransformerTest.java | 167 -----
...nitoringInfoToCounterUpdateTransformerTest.java | 95 ---
...nitoringInfoToCounterUpdateTransformerTest.java | 133 ----
.../RegisterAndProcessBundleOperationTest.java | 770 --------------------
.../SingularProcessBundleProgressTrackerTest.java | 148 ----
...nitoringInfoToCounterUpdateTransformerTest.java | 143 ----
...nitoringInfoToCounterUpdateTransformerTest.java | 133 ----
.../worker/fn/data/BeamFnDataGrpcServiceTest.java | 293 --------
.../fn/data/RemoteGrpcPortReadOperationTest.java | 157 ----
.../fn/data/RemoteGrpcPortWriteOperationTest.java | 228 ------
.../fn/logging/BeamFnLoggingServiceTest.java | 235 ------
.../fn/stream/ServerStreamObserverFactoryTest.java | 79 --
.../graph/CloneAmbiguousFlattensFunctionTest.java | 389 ----------
.../CreateRegisterFnOperationFunctionTest.java | 559 --------------
.../graph/DeduceFlattenLocationsFunctionTest.java | 394 ----------
.../graph/DeduceNodeLocationsFunctionTest.java | 324 ---------
...tFetchAndFilterStreamingSideInputNodesTest.java | 259 -------
.../graph/LengthPrefixUnknownCodersTest.java | 91 +--
.../runners/dataflow/worker/graph/NodesTest.java | 86 ---
.../RemoveFlattenInstructionsFunctionTest.java | 382 ----------
.../ReplacePgbkWithPrecombineFunctionTest.java | 153 ----
.../beam/runners/samza/SamzaPipelineOptions.java | 7 +
.../runners/samza/runtime/SamzaDoFnRunners.java | 70 +-
.../runtime/SamzaMetricsBundleProgressHandler.java | 154 ++++
.../SamzaMetricsBundleProgressHandlerTest.java | 187 +++++
.../samza/runtime/SdkHarnessDoFnRunnerTest.java} | 33 +-
sdks/go.mod | 2 +-
sdks/go.sum | 4 +-
sdks/java/extensions/zetasketch/build.gradle | 6 +-
.../beam/sdk/io/gcp/bigquery/AppendClientInfo.java | 115 ++-
.../beam/sdk/io/gcp/bigquery/BigQueryIO.java | 18 +-
.../beam/sdk/io/gcp/bigquery/BigQueryServices.java | 8 +
.../sdk/io/gcp/bigquery/BigQueryServicesImpl.java | 11 +
.../beam/sdk/io/gcp/bigquery/BigQueryUtils.java | 29 -
.../sdk/io/gcp/bigquery/SplittingIterable.java | 61 +-
.../bigquery/StorageApiDynamicDestinations.java | 2 +
.../StorageApiDynamicDestinationsBeamRow.java | 8 +-
.../StorageApiDynamicDestinationsTableRow.java | 24 +-
.../beam/sdk/io/gcp/bigquery/StorageApiLoads.java | 20 +-
.../io/gcp/bigquery/StorageApiWritePayload.java | 25 +-
.../StorageApiWriteRecordsInconsistent.java | 12 +-
.../bigquery/StorageApiWriteUnshardedRecords.java | 146 +++-
.../bigquery/StorageApiWritesShardedRecords.java | 173 +++--
.../io/gcp/bigquery/TableRowToStorageApiProto.java | 110 ++-
.../sdk/io/gcp/testing/FakeDatasetService.java | 63 +-
.../sdk/io/gcp/bigquery/BigQueryIOWriteTest.java | 216 ++++++
.../bigquery/TableRowToStorageApiProtoTest.java | 95 ++-
sdks/java/testing/load-tests/build.gradle | 6 +-
sdks/java/testing/nexmark/build.gradle | 6 +-
sdks/java/testing/tpcds/build.gradle | 6 +-
sdks/java/testing/watermarks/build.gradle | 6 +-
.../kfp/components/preprocessing/requirements.txt | 2 +-
.../kfp/components/train/requirements.txt | 2 +-
.../apache_beam/io/gcp/bigquery_file_loads.py | 2 +-
.../apache_beam/io/gcp/bigquery_file_loads_test.py | 41 ++
settings.gradle.kts | 1 -
132 files changed, 1732 insertions(+), 16274 deletions(-)
delete mode 100644
runners/google-cloud-dataflow-java/worker/legacy-worker/build.gradle
delete mode 100644
runners/google-cloud-dataflow-java/worker/src/main/java/com/google/cloud/dataflow/worker/DataflowRunnerHarness.java
delete mode 100644
runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/harness/util/package-info.java
delete mode 100644
runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/BeamFnMapTaskExecutorFactory.java
delete mode 100644
runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/DataflowRunnerHarness.java
delete mode 100644
runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/FetchAndFilterStreamingSideInputsOperation.java
delete mode 100644
runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/FnApiWindowMappingFn.java
delete mode 100644
runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/NoOpSourceOperationExecutor.java
delete mode 100644
runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/SdkHarnessRegistries.java
delete mode 100644
runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/SdkHarnessRegistry.java
delete mode 100644
runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/counters/CounterUpdateAggregator.java
delete mode 100644
runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/counters/CounterUpdateAggregators.java
delete mode 100644
runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/counters/DistributionCounterUpdateAggregator.java
delete mode 100644
runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/counters/MeanCounterUpdateAggregator.java
delete mode 100644
runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/counters/SumCounterUpdateAggregator.java
delete mode 100644
runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/fn/BeamFnControlService.java
delete mode 100644
runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/fn/control/BeamFnMapTaskExecutor.java
delete mode 100644
runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/fn/control/DataflowSideInputHandlerFactory.java
delete mode 100644
runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/fn/control/ElementCountMonitoringInfoToCounterUpdateTransformer.java
delete mode 100644
runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/fn/control/ExecutionTimeMonitoringInfoToCounterUpdateTransformer.java
delete mode 100644
runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/fn/control/FnApiMonitoringInfoToCounterUpdateTransformer.java
delete mode 100644
runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/fn/control/MeanByteCountMonitoringInfoToCounterUpdateTransformer.java
delete mode 100644
runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/fn/control/MonitoringInfoToCounterUpdateTransformer.java
delete mode 100644
runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/fn/control/ProcessRemoteBundleOperation.java
delete mode 100644
runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/fn/control/RegisterAndProcessBundleOperation.java
delete mode 100644
runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/fn/control/UserDistributionMonitoringInfoToCounterUpdateTransformer.java
delete mode 100644
runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/fn/control/UserMonitoringInfoToCounterUpdateTransformer.java
delete mode 100644
runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/fn/data/BeamFnDataGrpcService.java
delete mode 100644
runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/fn/data/RemoteGrpcPortReadOperation.java
delete mode 100644
runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/fn/data/RemoteGrpcPortWriteOperation.java
delete mode 100644
runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/fn/grpc/BeamFnService.java
delete mode 100644
runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/fn/logging/BeamFnLoggingService.java
delete mode 100644
runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/fn/stream/ServerStreamObserverFactory.java
delete mode 100644
runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/graph/CloneAmbiguousFlattensFunction.java
delete mode 100644
runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/graph/CreateExecutableStageNodeFunction.java
delete mode 100644
runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/graph/CreateRegisterFnOperationFunction.java
delete mode 100644
runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/graph/DeduceFlattenLocationsFunction.java
delete mode 100644
runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/graph/DeduceNodeLocationsFunction.java
delete mode 100644
runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/graph/InsertFetchAndFilterStreamingSideInputNodes.java
delete mode 100644
runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/graph/RegisterNodeFunction.java
delete mode 100644
runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/graph/RemoveFlattenInstructionsFunction.java
delete mode 100644
runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/graph/ReplacePgbkWithPrecombineFunction.java
delete mode 100644
runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/FnApiWindowMappingFnTest.java
delete mode 100644
runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/NoOpSourceOperationExecutorTest.java
delete mode 100644
runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/SdkHarnessRegistryTest.java
delete mode 100644
runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/counters/CounterUpdateAggregatorsTest.java
delete mode 100644
runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/counters/DistributionCounterUpdateAggregatorTest.java
delete mode 100644
runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/counters/MeanCounterUpdateAggregatorTest.java
delete mode 100644
runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/counters/SumCounterUpdateAggregatorTest.java
delete mode 100644
runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/fn/BeamFnControlServiceTest.java
delete mode 100644
runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/fn/control/BeamFnMapTaskExecutorTest.java
delete mode 100644
runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/fn/control/DataflowSideInputHandlerFactoryTest.java
delete mode 100644
runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/fn/control/ElementCountMonitoringInfoToCounterUpdateTransformerTest.java
delete mode 100644
runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/fn/control/ExecutionTimeMonitoringInfoToCounterUpdateTransformerTest.java
delete mode 100644
runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/fn/control/FnApiMonitoringInfoToCounterUpdateTransformerTest.java
delete mode 100644
runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/fn/control/MeanByteCountMonitoringInfoToCounterUpdateTransformerTest.java
delete mode 100644
runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/fn/control/RegisterAndProcessBundleOperationTest.java
delete mode 100644
runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/fn/control/SingularProcessBundleProgressTrackerTest.java
delete mode 100644
runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/fn/control/UserDistributionMonitoringInfoToCounterUpdateTransformerTest.java
delete mode 100644
runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/fn/control/UserMonitoringInfoToCounterUpdateTransformerTest.java
delete mode 100644
runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/fn/data/BeamFnDataGrpcServiceTest.java
delete mode 100644
runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/fn/data/RemoteGrpcPortReadOperationTest.java
delete mode 100644
runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/fn/data/RemoteGrpcPortWriteOperationTest.java
delete mode 100644
runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/fn/logging/BeamFnLoggingServiceTest.java
delete mode 100644
runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/fn/stream/ServerStreamObserverFactoryTest.java
delete mode 100644
runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/graph/CloneAmbiguousFlattensFunctionTest.java
delete mode 100644
runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/graph/CreateRegisterFnOperationFunctionTest.java
delete mode 100644
runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/graph/DeduceFlattenLocationsFunctionTest.java
delete mode 100644
runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/graph/DeduceNodeLocationsFunctionTest.java
delete mode 100644
runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/graph/InsertFetchAndFilterStreamingSideInputNodesTest.java
delete mode 100644
runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/graph/RemoveFlattenInstructionsFunctionTest.java
delete mode 100644
runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/graph/ReplacePgbkWithPrecombineFunctionTest.java
create mode 100644
runners/samza/src/main/java/org/apache/beam/runners/samza/runtime/SamzaMetricsBundleProgressHandler.java
create mode 100644
runners/samza/src/test/java/org/apache/beam/runners/samza/runtime/SamzaMetricsBundleProgressHandlerTest.java
copy
runners/samza/src/{main/java/org/apache/beam/runners/samza/runtime/OutputManagerFactory.java
=>
test/java/org/apache/beam/runners/samza/runtime/SdkHarnessDoFnRunnerTest.java}
(54%)