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 0b6057d0645 Merge pull request #39718 from apache/website-2-76
add e6d7f23de91 SolaceIO: support binary content (text, bytes) data
payload (#39876)
add b5863dceb48 Bump cloud.google.com/go/bigquery from 1.81.0 to 1.82.0 in
/sdks (#39945)
add 30689dcb683 Bump github.com/aws/aws-sdk-go-v2/service/s3 in /sdks
(#39946)
add 33ce824d67e Bump cloud.google.com/go/storage from 1.65.0 to 1.66.0 in
/sdks (#39944)
add a93fb2c4e64 Bump cloud.google.com/go/bigtable from 1.52.0 to 1.53.0 in
/sdks (#39948)
add a8f3e208ed7 Bump github.com/aws/aws-sdk-go-v2/feature/s3/manager in
/sdks (#39947)
add 4855a35e92d [Python] Honor fixed object detection batch size (#39953)
add 0ba44c65b80 Bump the KafkaToPubsubE2ETest emulator image to an
existing tag (#39956)
add 261ed247388 [Spark 4] Reuse the legacy maxRecordsPerBatch option in
the Structured Streaming runner (#39952)
add e671af6c779 [Spark][#36841] Register Spark's streaming internals with
Kryo for the Structured Streaming runner (#39939)
add 6e1cb9683f3 upgrade spotless to 7.2.1 (#39936)
add 3b476c42095 [Flink] Cache batch side input materialization in Flink 2
(#39867)
add 3e206896f0a Bump browserslist from 4.16.3 to 4.28.8 in /website/www
(#39964)
add 4c02fcff485 Bump transformers (#39962)
add d3ccb3a9a75 upgrade grpc to 1.82.1 for test infra and playground
(#39915)
No new revisions were added by this update.
Summary of changes:
.test-infra/metrics/build.gradle | 3 +-
.../metrics/src/test/groovy/ProberTests.groovy | 1 -
.test-infra/mock-apis/go.mod | 50 ++--
.test-infra/mock-apis/go.sum | 272 ++++++------------
.../beam/testinfra/mockapis/echo/v1/Echo.java | 30 ++
CHANGES.md | 2 +
buildSrc/build.gradle.kts | 12 +-
.../org/apache/beam/gradle/BeamDockerPlugin.groovy | 6 +-
.../org/apache/beam/gradle/BeamModulePlugin.groovy | 48 +++-
.../apache/beam/examples/RateLimiterSimple.java | 1 +
.../apache/beam/examples/WindowedWordCount.java | 1 +
.../java/org/apache/beam/examples/WordCount.java | 3 +
.../beam/examples/complete/game/GameStats.java | 1 +
.../beam/examples/complete/game/LeaderBoard.java | 1 +
.../beam/examples/complete/game/UserScore.java | 1 +
.../options/KafkaToPubsubOptions.java | 3 +-
.../apache/beam/examples/snippets/Snippets.java | 4 +
.../subprocess/kernel/SubProcessIOFiles.java | 4 +-
.../examples/complete/game/LeaderBoardTest.java | 1 +
.../kafkatopubsub/KafkaToPubsubE2ETest.java | 2 +-
.../ReadFromTwitterDoFn.java | 1 +
.../apache/beam/it/common/PipelineLauncher.java | 3 +-
.../common/dataflow/DefaultPipelineLauncher.java | 4 +-
.../beam/it/gcp/bigquery/BigQueryStreamingLT.java | 3 +-
playground/backend/cmd/server/controller_test.go | 2 +-
playground/backend/containers/java/Dockerfile | 2 +-
playground/backend/containers/python/Dockerfile | 2 +-
playground/backend/containers/python/build.gradle | 2 +-
playground/backend/containers/router/Dockerfile | 15 +-
playground/backend/containers/router/build.gradle | 2 +-
playground/backend/containers/scio/Dockerfile | 2 +-
playground/backend/go.mod | 65 +++--
playground/backend/go.sum | 138 ++++++---
.../code_processing/code_processing_test.go | 13 +-
.../beam/runners/core/InMemoryTimerInternals.java | 8 +-
...TimeBoundedSplittableProcessElementInvoker.java | 3 +-
.../org/apache/beam/runners/core/ReduceFn.java | 2 +
.../apache/beam/runners/core/ReduceFnRunner.java | 4 +-
.../apache/beam/runners/core/SimpleDoFnRunner.java | 6 +-
.../core/SplittableParDoViaKeyedWorkItems.java | 3 +-
.../org/apache/beam/runners/core/StateTags.java | 8 +-
.../apache/beam/runners/core/TimerInternals.java | 14 +-
.../beam/runners/core/metrics/BoundedTrieData.java | 1 +
.../beam/runners/core/metrics/MetricUpdates.java | 1 +
.../runners/core/metrics/MetricsContainerImpl.java | 32 ++-
.../core/metrics/MonitoringInfoConstants.java | 4 +-
.../core/metrics/MonitoringInfoMetricName.java | 4 +-
.../core/metrics/SimpleMonitoringInfoBuilder.java | 1 +
.../runners/core/triggers/TriggerStateMachine.java | 3 +-
.../core/metrics/MonitoringInfoTestUtil.java | 8 +-
.../core/triggers/AfterPaneStateMachineTest.java | 1 +
.../beam/runners/direct/CommittedBundle.java | 4 +-
.../beam/runners/direct/DirectTimerInternals.java | 8 +-
.../runners/direct/DirectTransformExecutor.java | 5 +-
.../beam/runners/direct/MultiStepCombine.java | 1 +
.../runners/direct/TransformEvaluatorFactory.java | 3 +-
.../beam/runners/direct/TransformResult.java | 3 +-
.../beam/runners/direct/UncommittedBundle.java | 3 +-
.../beam/runners/direct/WatermarkManager.java | 1 +
.../beam/runners/direct/DirectRunnerTest.java | 3 +-
.../wrappers/streaming/DoFnOperator.java | 18 +-
.../FlinkStreamingPortablePipelineTranslator.java | 4 +-
.../wrappers/streaming/CachedSideInputReader.java | 93 ++++++
.../wrappers/streaming/DoFnOperator.java | 49 +++-
.../wrappers/streaming/SideInputCache.java | 117 ++++++++
.../streaming/io/UnboundedSourceWrapper.java | 2 +
.../runners/flink/FlinkPipelineOptionsTest.java | 1 -
.../streaming/FlinkCachedSideInputReaderTest.java | 283 +++++++++++++++++++
.../wrappers/streaming/DoFnOperator.java | 49 +++-
.../runners/flink/FlinkPipelineOptionsTest.java | 1 -
runners/flink/build.gradle | 6 +-
.../beam/runners/flink/FlinkJobServerDriver.java | 3 +-
.../beam/runners/flink/FlinkPipelineOptions.java | 4 +-
.../FlinkStreamingPortablePipelineTranslator.java | 4 +-
.../flink/translation/utils/CheckpointStats.java | 1 +
.../wrappers/streaming/DoFnOperator.java | 18 +-
.../streaming/ExecutableStageDoFnOperator.java | 5 +
.../streaming/io/UnboundedSourceWrapper.java | 2 +
.../streaming/stableinput/BufferingDoFnRunner.java | 10 +
.../streaming/ExecutableStageDoFnOperatorTest.java | 14 +-
.../beam/runners/dataflow/BatchViewOverrides.java | 4 +
.../beam/runners/dataflow/DataflowMetrics.java | 4 +-
.../beam/runners/dataflow/DataflowPipelineJob.java | 9 +-
.../beam/runners/dataflow/DataflowRunner.java | 10 +-
.../beam/runners/dataflow/TestDataflowRunner.java | 3 +-
.../options/DataflowPipelineWorkerPoolOptions.java | 24 +-
.../options/DataflowStreamingPipelineOptions.java | 8 +-
.../runners/dataflow/util/DataflowTemplateJob.java | 3 +-
...DefaultCoderCloudObjectTranslatorRegistrar.java | 2 +
.../beam/runners/dataflow/util/PackageUtil.java | 2 +
.../beam/runners/dataflow/util/PropertyNames.java | 1 +
.../beam/runners/dataflow/DataflowRunnerTest.java | 7 +-
.../runners/dataflow/util/CloudObjectsTest.java | 2 +-
.../dataflow/worker/BatchDataflowWorker.java | 14 +
.../dataflow/worker/DataflowMapTaskExecutor.java | 4 +-
.../worker/DataflowPortabilityPCollectionView.java | 3 +-
.../dataflow/worker/DoFnInstanceManagers.java | 1 +
.../dataflow/worker/GroupAlsoByWindowFnRunner.java | 1 +
.../dataflow/worker/IsmSideInputReader.java | 2 +
.../runners/dataflow/worker/OperationalLimits.java | 2 +
.../beam/runners/dataflow/worker/PubsubReader.java | 3 +-
.../worker/RemoveSafeDeltaCounterCell.java | 1 +
.../runners/dataflow/worker/SimpleParDoFn.java | 3 +-
.../dataflow/worker/StreamingDataflowWorker.java | 1 +
.../StreamingKeyedWorkItemSideInputParDoFn.java | 3 +-
.../worker/StreamingModeExecutionContext.java | 10 +-
.../worker/StreamingStepMetricsContainer.java | 1 +
.../dataflow/worker/UngroupedWindmillReader.java | 3 +-
.../dataflow/worker/WindowingWindmillReader.java | 6 +-
.../runners/dataflow/worker/counters/Counter.java | 45 +--
.../worker/graph/LengthPrefixUnknownCoders.java | 5 +-
.../logging/DataflowWorkerLoggingHandler.java | 1 +
.../logging/DataflowWorkerLoggingInitializer.java | 3 +-
.../streaming/harness/StreamingCounters.java | 2 +
.../dataflow/worker/util/MemoryMonitor.java | 14 +-
.../common/worker/BatchingShuffleEntryReader.java | 1 +
.../worker/util/common/worker/GroupingTables.java | 1 +
.../util/common/worker/ShuffleReadCounter.java | 1 +
.../worker/windmill/client/WindmillStreamPool.java | 4 +-
.../worker/windmill/client/commits/Commit.java | 3 +-
.../client/getdata/ApplianceGetDataClient.java | 3 +-
.../client/grpc/GrpcDirectGetWorkStream.java | 3 +-
.../client/grpc/GrpcWindmillStreamFactory.java | 4 +-
.../grpc/observers/DirectStreamObserver.java | 4 +-
.../client/grpc/stubs/FailoverChannel.java | 4 +
.../worker/windmill/state/PagingIterable.java | 8 +-
.../worker/windmill/state/ToIterableFunction.java | 9 +-
.../worker/windmill/state/WindmillBag.java | 1 +
.../windmill/state/WindmillStateInternals.java | 1 +
.../worker/windmill/state/WindmillStateReader.java | 19 +-
.../worker/windmill/state/WindmillValue.java | 2 +
.../windmill/state/WindmillWatermarkHold.java | 1 +
.../work/processing/StreamingWorkScheduler.java | 9 +-
.../dataflow/worker/AvroByteReaderFactoryTest.java | 10 +-
.../dataflow/worker/ConcatReaderFactoryTest.java | 2 +-
.../dataflow/worker/FakeWindmillServer.java | 4 +-
.../dataflow/worker/InMemoryReaderFactoryTest.java | 2 +-
.../IntrinsicMapTaskExecutorFactoryTest.java | 8 +-
.../dataflow/worker/IsmSideInputReaderTest.java | 2 +-
.../runners/dataflow/worker/ReaderFactoryTest.java | 5 +-
.../dataflow/worker/ShuffleReaderFactoryTest.java | 6 +-
.../runners/dataflow/worker/SinkRegistryTest.java | 2 +-
.../worker/StreamingDataflowWorkerTest.java | 19 +-
.../worker/StreamingModeExecutionContextTest.java | 8 +-
...eamingPCollectionViewWriterDoFnFactoryTest.java | 2 +-
.../dataflow/worker/WorkerCustomSourcesTest.java | 22 +-
.../graph/LengthPrefixUnknownCodersTest.java | 42 +--
.../logging/DataflowWorkerLoggingHandlerTest.java | 4 +-
.../harness/WindmillStreamSenderTest.java | 9 +-
.../worker/util/BoundedQueueExecutorTest.java | 4 +-
.../dataflow/worker/util/CloudSourceUtilsTest.java | 2 +-
.../worker/util/KeyGroupWorkQueueTest.java | 2 +-
.../dataflow/worker/util/TimerOrElementTest.java | 2 +-
.../common/worker/WorkProgressUpdaterTest.java | 3 +-
.../client/grpc/FakeWindmillGrpcService.java | 3 +-
.../client/grpc/GrpcCommitWorkStreamTest.java | 6 +-
.../processing/StreamingCommitFinalizerTest.java | 4 +-
.../failures/WorkFailureProcessorTest.java | 4 +-
.../work/refresh/ActiveWorkRefresherTest.java | 4 +-
.../artifact/ArtifactStagingService.java | 3 +-
.../fnexecution/control/BundleSplitHandler.java | 3 +-
.../control/DefaultJobBundleFactory.java | 1 +
.../control/ProcessBundleDescriptors.java | 12 +-
.../runners/fnexecution/data/GrpcDataService.java | 5 +-
.../status/BeamWorkerStatusGrpcService.java | 3 +-
.../runners/jobsubmission/InMemoryJobService.java | 6 +-
.../jobsubmission/PortablePipelineJarCreator.java | 4 +-
.../java/org/apache/beam/runners/local/Bundle.java | 3 +-
.../beam/runners/portability/PortableRunner.java | 2 +
.../apache/beam/runners/prism/PrismRegistrar.java | 1 +
.../SparkKryoRegistratorStreamingTest.java | 130 +++++++++
runners/spark/build.gradle | 6 +-
.../runners/spark/SparkNativePipelineVisitor.java | 4 +-
.../runners/spark/TestSparkPipelineOptions.java | 3 +-
.../SparkStructuredStreamingPipelineOptions.java | 15 +-
.../translation/SparkSessionFactory.java | 36 ++-
.../translation/batch/Aggregators.java | 3 +-
...parkStructuredStreamingPipelineOptionsTest.java | 57 ++++
.../beam/runners/twister2/Twister2Runner.java | 1 +
sdks/go.mod | 16 +-
sdks/go.sum | 32 +--
.../resources/beam/checkstyle/suppressions.xml | 2 +
.../main/java/org/apache/beam/sdk/Pipeline.java | 1 +
.../java/org/apache/beam/sdk/coders/Coder.java | 1 +
.../java/org/apache/beam/sdk/coders/ZstdCoder.java | 1 +
.../sdk/fn/data/BeamFnDataGrpcMultiplexer.java | 14 +-
.../sdk/fn/data/BeamFnDataInboundObserver.java | 5 +-
.../apache/beam/sdk/fn/server/GrpcFnServer.java | 4 +-
.../apache/beam/sdk/fn/server/ServerFactory.java | 1 +
.../sdk/io/BoundedReadFromUnboundedSource.java | 3 +-
.../org/apache/beam/sdk/io/CompressedSource.java | 4 +-
.../org/apache/beam/sdk/io/CountingSource.java | 5 +
.../apache/beam/sdk/io/DefaultFilenamePolicy.java | 1 +
.../java/org/apache/beam/sdk/io/FileBasedSink.java | 13 +-
.../main/java/org/apache/beam/sdk/io/FileIO.java | 2 +
.../java/org/apache/beam/sdk/io/FileSystems.java | 3 +-
.../org/apache/beam/sdk/io/OffsetBasedSource.java | 4 +-
.../main/java/org/apache/beam/sdk/io/Source.java | 4 +-
.../java/org/apache/beam/sdk/io/TFRecordIO.java | 12 +-
.../main/java/org/apache/beam/sdk/io/TextIO.java | 12 +-
.../java/org/apache/beam/sdk/io/WriteFiles.java | 5 +-
.../java/org/apache/beam/sdk/io/fs/ResourceId.java | 3 +-
.../apache/beam/sdk/lineage/LineageOptions.java | 3 +-
.../java/org/apache/beam/sdk/metrics/Lineage.java | 4 +-
.../beam/sdk/metrics/MetricsEnvironment.java | 3 +-
.../java/org/apache/beam/sdk/options/Default.java | 2 +
.../beam/sdk/options/ExperimentalOptions.java | 3 +-
.../apache/beam/sdk/options/PipelineOptions.java | 3 +-
.../beam/sdk/options/PortablePipelineOptions.java | 9 +-
.../apache/beam/sdk/options/StreamingOptions.java | 3 +-
.../beam/sdk/schemas/FieldAccessDescriptor.java | 6 +-
.../beam/sdk/schemas/FieldTypeDescriptors.java | 2 +
.../apache/beam/sdk/schemas/FieldValueGetter.java | 3 +-
.../sdk/schemas/FieldValueTypeInformation.java | 7 +-
.../sdk/schemas/GetterBasedSchemaProvider.java | 9 +-
.../java/org/apache/beam/sdk/schemas/Schema.java | 39 +--
.../apache/beam/sdk/schemas/SchemaProvider.java | 9 +-
.../org/apache/beam/sdk/schemas/io/Failure.java | 1 +
.../sdk/schemas/logicaltypes/EnumerationType.java | 1 +
.../apache/beam/sdk/schemas/transforms/Cast.java | 3 +-
.../apache/beam/sdk/schemas/transforms/Join.java | 3 +-
.../beam/sdk/schemas/transforms/RenameFields.java | 1 +
.../apache/beam/sdk/schemas/transforms/Select.java | 1 +
.../beam/sdk/schemas/utils/ByteBuddyUtils.java | 3 +-
.../beam/sdk/schemas/utils/SelectHelpers.java | 1 +
.../java/org/apache/beam/sdk/state/Timers.java | 3 +-
.../java/org/apache/beam/sdk/state/ValueState.java | 3 +-
.../beam/sdk/testing/TestPipelineOptions.java | 3 +-
.../beam/sdk/transforms/ApproximateUnique.java | 16 +-
.../apache/beam/sdk/transforms/BatchElements.java | 2 +
.../org/apache/beam/sdk/transforms/Create.java | 1 +
.../apache/beam/sdk/transforms/Deduplicate.java | 1 +
.../java/org/apache/beam/sdk/transforms/DoFn.java | 1 +
.../org/apache/beam/sdk/transforms/DoFnTester.java | 116 ++++++--
.../beam/sdk/transforms/GroupIntoBatches.java | 3 +-
.../java/org/apache/beam/sdk/transforms/ParDo.java | 1 +
.../java/org/apache/beam/sdk/transforms/Top.java | 15 +-
.../java/org/apache/beam/sdk/transforms/View.java | 4 +-
.../java/org/apache/beam/sdk/transforms/Watch.java | 3 +-
.../sdk/transforms/errorhandling/ErrorHandler.java | 3 +-
.../beam/sdk/transforms/join/CoGbkResult.java | 1 +
.../reflect/ByteBuddyDoFnInvokerFactory.java | 10 +-
.../beam/sdk/transforms/reflect/DoFnInvoker.java | 15 +-
.../beam/sdk/transforms/reflect/DoFnSignature.java | 11 +-
.../sdk/transforms/reflect/DoFnSignatures.java | 2 +
.../sdk/transforms/windowing/DefaultTrigger.java | 4 +-
.../org/apache/beam/sdk/util/HistogramData.java | 5 +
.../main/java/org/apache/beam/sdk/util/Holder.java | 3 +-
.../util/construction/DefaultArtifactResolver.java | 12 +-
.../beam/sdk/util/construction/Environments.java | 6 +-
.../beam/sdk/util/construction/External.java | 10 +-
.../sdk/util/construction/ExternalTranslation.java | 4 +-
.../util/construction/PTransformTranslation.java | 10 +-
.../sdk/util/construction/ParDoTranslation.java | 9 +
.../beam/sdk/util/construction/SdkComponents.java | 5 +-
.../sdk/util/construction/SplittableParDo.java | 1 +
.../sdk/util/construction/TransformUpgrader.java | 4 +-
.../UnboundedReadFromBoundedSource.java | 7 +-
.../util/construction/graph/ExecutableStage.java | 3 +-
.../graph/ImmutableExecutableStage.java | 3 +-
.../construction/graph/OutputDeduplicator.java | 4 +-
.../util/construction/graph/ProtoOverrides.java | 3 +-
.../construction/graph/SideInputReference.java | 2 +
.../graph/SplittableParDoExpander.java | 3 +-
.../util/construction/graph/TimerReference.java | 1 +
.../construction/graph/UserStateReference.java | 2 +
.../apache/beam/sdk/values/PCollectionList.java | 1 +
.../apache/beam/sdk/values/PCollectionView.java | 4 +-
.../apache/beam/sdk/values/PCollectionViews.java | 1 +
.../java/org/apache/beam/sdk/values/PValues.java | 3 +-
.../java/org/apache/beam/sdk/values/RowUtils.java | 6 +-
.../beam/sdk/values/ValueInSingleWindow.java | 3 +-
.../org/apache/beam/sdk/values/WindowedValue.java | 9 +-
.../org/apache/beam/sdk/values/WindowedValues.java | 4 +-
.../beam/sdk/schemas/transforms/FilterTest.java | 3 +-
.../beam/sdk/schemas/transforms/SelectTest.java | 3 +-
.../TypedSchemaTransformProviderTest.java | 3 +-
.../org/apache/beam/sdk/transforms/CreateTest.java | 1 +
.../util/construction/ExternalTranslationTest.java | 4 +-
.../util/construction/ValidateRunnerXlangTest.java | 1 +
.../graph/GreedyPipelineFuserTest.java | 35 +--
.../construction/graph/GreedyStageFuserTest.java | 45 +--
.../sdk/expansion/service/ExpansionService.java | 4 +-
.../apache/beam/sdk/extensions/avro/io/AvroIO.java | 10 +-
.../beam/sdk/extensions/avro/io/AvroSource.java | 3 +-
.../extensions/avro/schemas/utils/AvroUtils.java | 3 +-
.../beam/sdk/extensions/avro/io/AvroIOTest.java | 4 +-
.../client/accumulators/AccumulatorProvider.java | 1 +
.../euphoria/core/client/type/TypeAware.java | 4 +-
.../euphoria/core/client/type/TypeAwareness.java | 4 +-
.../extensions/euphoria/core/client/util/Fold.java | 4 +-
.../core/client/util/PCollectionLists.java | 4 +-
.../core/translate/BeamAccumulatorProvider.java | 4 +-
.../euphoria/core/translate/FlatMapTranslator.java | 4 +-
.../core/translate/provider/CompositeProvider.java | 4 +-
.../sdk/extensions/euphoria/core/util/IOUtils.java | 4 +-
.../sdk/extensions/gcp/auth/CredentialFactory.java | 3 +-
.../extensions/gcp/auth/GcpCredentialFactory.java | 3 +-
.../sdk/extensions/gcp/options/GcpOptions.java | 17 +-
.../sdk/extensions/gcp/options/GcsOptions.java | 23 +-
.../beam/sdk/extensions/gcp/util/GcsUtil.java | 44 ++-
.../beam/sdk/extensions/gcp/util/GcsUtilV1.java | 32 ++-
.../beam/sdk/extensions/ml/DLPReidentifyText.java | 1 +
.../beam/sdk/extensions/ml/CloudVisionTest.java | 3 +-
...GcpAuthAutoConfigurationCustomizerProvider.java | 12 +-
.../ordered/ContiguousSequenceRange.java | 8 +-
.../beam/sdk/extensions/ordered/EventExaminer.java | 3 +-
.../ordered/combiner/SequenceRangeAccumulator.java | 8 +-
.../apache/beam/sdk/extensions/ordered/Event.java | 12 +-
.../extensions/protobuf/ProtoBeamConverter.java | 9 +-
.../protobuf/ProtoDynamicMessageSchema.java | 4 +-
.../extensions/python/PythonExternalTransform.java | 4 +-
.../datacatalog/DataCatalogTableProvider.java | 10 +-
.../meta/provider/iceberg/IcebergMetastore.java | 3 +-
.../apache/beam/sdk/extensions/sql/BeamSqlCli.java | 1 +
.../beam/sdk/extensions/sql/SqlTransform.java | 1 +
.../beam/sdk/extensions/sql/impl/BeamSqlEnv.java | 1 +
.../extensions/sql/impl/CatalogManagerSchema.java | 4 +-
.../sdk/extensions/sql/impl/CatalogSchema.java | 1 +
.../sql/impl/UdfImplReflectiveFunctionBase.java | 1 +
.../sdk/extensions/sql/impl/cep/CEPLiteral.java | 3 +-
.../extensions/sql/impl/cep/PatternCondition.java | 3 +-
.../sql/impl/parser/SqlCreateDatabase.java | 4 +-
.../sdk/extensions/sql/impl/rel/BeamRelNode.java | 3 +-
.../sdk/extensions/sql/impl/rel/BeamUnnestRel.java | 1 +
.../impl/transform/BeamBuiltinAggregations.java | 3 +
.../sdk/extensions/sql/meta/catalog/Catalog.java | 3 +-
.../sql/meta/catalog/CatalogManager.java | 3 +-
.../sql/meta/provider/mongodb/MongoDbTable.java | 1 +
.../sql/meta/provider/test/TestBoundedTable.java | 1 +
.../sql/meta/provider/test/TestUnboundedTable.java | 1 +
.../sql/meta/store/InMemoryMetaStore.java | 4 +-
.../sql/impl/utils/CalciteUtilsTest.java | 2 +-
.../provider/TestSchemaIOTableProviderWrapper.java | 1 +
.../sql/meta/provider/kafka/KafkaTestTable.java | 1 +
.../beam/sdk/extensions/yaml/YamlTransform.java | 1 +
.../zetasketch/ApproximateCountDistinctTest.java | 3 +
.../apache/beam/fn/harness/FnApiDoFnRunner.java | 5 +-
.../org/apache/beam/fn/harness/HandlesSplits.java | 3 +-
.../fn/harness/control/ExecutionStateSampler.java | 1 +
.../apache/beam/fn/harness/control/Metrics.java | 2 +
.../fn/harness/control/ProcessBundleHandler.java | 10 +-
.../apache/beam/fn/harness/state/BagUserState.java | 9 +-
.../beam/fn/harness/state/MultimapUserState.java | 3 +-
.../fn/harness/state/StateFetchingIterators.java | 3 +-
.../beam/fn/harness/status/MemoryMonitor.java | 14 +-
.../beam/fn/harness/BeamFnDataReadRunnerTest.java | 4 +-
.../beam/fn/harness/BeamFnDataWriteRunnerTest.java | 4 +-
.../apache/beam/fn/harness/CombineRunnersTest.java | 1 +
.../beam/fn/harness/FnApiDoFnRunnerTest.java | 4 +-
...bleTruncateSizedRestrictionsDoFnRunnerTest.java | 4 +-
.../fn/harness/state/StateBackedIterableTest.java | 3 +-
.../beam/sdk/io/aws2/dynamodb/DynamoDBIO.java | 3 +-
.../sdk/io/aws2/kinesis/EFOShardSubscriber.java | 1 +
.../io/aws2/kinesis/EFOShardSubscribersPool.java | 3 +-
.../io/aws2/kinesis/WatermarkPolicyFactory.java | 3 +-
.../apache/beam/sdk/io/aws2/options/AwsModule.java | 3 +-
.../apache/beam/sdk/io/aws2/options/S3Options.java | 6 +-
.../org/apache/beam/sdk/io/aws2/sqs/SqsIO.java | 4 +-
.../beam/sdk/io/aws2/common/ObjectPoolTest.java | 3 +-
.../sdk/io/aws2/kinesis/testing/KinesisIOIT.java | 6 +-
.../beam/sdk/io/aws2/s3/S3FileSystemTest.java | 3 +-
.../apache/beam/sdk/io/aws2/s3/S3TestUtils.java | 3 +-
.../apache/beam/sdk/io/azure/cosmos/CosmosIO.java | 3 +-
.../beam/sdk/io/azure/cosmos/CosmosOptions.java | 6 +-
.../sdk/io/azure/blobstore/AzfsResourceId.java | 1 +
.../sdk/io/azure/options/BlobstoreOptions.java | 12 +-
sdks/java/io/cassandra/build.gradle | 2 +-
.../org/apache/beam/sdk/io/cassandra/ReadFn.java | 9 +-
.../beam/sdk/io/cassandra/CassandraIOTest.java | 1 +
.../beam/sdk/io/clickhouse/ClickHouseIO.java | 1 +
.../beam/sdk/io/common/IOTestPipelineOptions.java | 6 +-
.../io/contextualtextio/ContextualTextIOTest.java | 1 +
.../apache/beam/sdk/io/csv/CsvRowConversions.java | 3 +-
.../apache/beam/io/debezium/OffsetRetainer.java | 3 +-
.../io/debezium/DebeziumIOMySqlConnectorIT.java | 1 +
.../io/elasticsearch/ElasticsearchIOTestUtils.java | 1 +
.../beam/sdk/io/elasticsearch/ElasticsearchIO.java | 6 +-
.../io/common/FileBasedIOTestPipelineOptions.java | 3 +-
.../CsvWriteSchemaTransformFormatProviderTest.java | 3 +-
.../FileReadSchemaTransformFormatProviderTest.java | 3 +-
...FileWriteSchemaTransformFormatProviderTest.java | 9 +-
.../io/fileschematransform/XmlRowAdapterTest.java | 20 +-
.../beam/sdk/io/googleads/GoogleAdsOptions.java | 6 +-
.../beam/sdk/io/gcp/bigquery/AppendClientInfo.java | 3 +-
.../sdk/io/gcp/bigquery/BigQueryAvroUtils.java | 1 +
.../beam/sdk/io/gcp/bigquery/BigQueryHelpers.java | 3 +-
.../beam/sdk/io/gcp/bigquery/BigQueryIO.java | 37 ++-
.../beam/sdk/io/gcp/bigquery/BigQueryServices.java | 16 +-
.../sdk/io/gcp/bigquery/BigQueryServicesImpl.java | 22 +-
.../io/gcp/bigquery/BigQueryStorageAvroReader.java | 2 +-
.../gcp/bigquery/BigQueryStorageQuerySource.java | 4 +-
.../gcp/bigquery/BigQueryStorageTableSource.java | 6 +-
.../sdk/io/gcp/bigquery/DynamicDestinations.java | 3 +-
.../StorageApiDynamicDestinationsBeamRow.java | 9 +-
.../StorageApiDynamicDestinationsProto.java | 3 +-
.../StorageApiDynamicDestinationsTableRow.java | 6 +-
.../bigquery/StorageApiWriteUnshardedRecords.java | 13 +-
.../bigquery/StorageApiWritesShardedRecords.java | 17 +-
.../io/gcp/bigquery/StreamingInsertsMetrics.java | 1 +
.../sdk/io/gcp/bigquery/TableRowJsonCoder.java | 3 +-
.../io/gcp/bigquery/TableRowToStorageApiProto.java | 21 +-
.../beam/sdk/io/gcp/bigquery/TableSchemaCache.java | 9 +-
.../io/gcp/bigquery/TableSchemaUpdateUtils.java | 3 +-
.../sdk/io/gcp/bigquery/UpgradeTableSchema.java | 3 +-
.../beam/sdk/io/gcp/bigquery/WriteTables.java | 1 +
.../beam/sdk/io/gcp/bigtable/BigtableConfig.java | 21 +-
.../io/gcp/bigtable/BigtableConfigTranslator.java | 3 +-
.../beam/sdk/io/gcp/bigtable/BigtableIO.java | 17 +-
.../changestreams/ChangeStreamMetrics.java | 1 +
.../beam/sdk/io/gcp/datastore/DatastoreV1.java | 2 +
.../sdk/io/gcp/firestore/FirestoreOptions.java | 6 +-
.../beam/sdk/io/gcp/firestore/FirestoreV1.java | 6 +-
.../beam/sdk/io/gcp/firestore/RpcQosImpl.java | 1 +
.../beam/sdk/io/gcp/firestore/RpcQosOptions.java | 1 +
.../apache/beam/sdk/io/gcp/healthcare/DicomIO.java | 1 +
.../apache/beam/sdk/io/gcp/healthcare/FhirIO.java | 7 +
.../io/gcp/healthcare/FhirIOPatientEverything.java | 2 +
.../sdk/io/gcp/healthcare/FhirSearchParameter.java | 2 +
.../apache/beam/sdk/io/gcp/healthcare/HL7v2IO.java | 3 +
.../beam/sdk/io/gcp/healthcare/HL7v2Message.java | 1 +
.../sdk/io/gcp/healthcare/HealthcareApiClient.java | 1 +
.../sdk/io/gcp/pubsub/PreparePubsubWriteDoFn.java | 5 +-
.../apache/beam/sdk/io/gcp/pubsub/PubsubIO.java | 6 +-
.../sdk/io/gcp/pubsub/PubsubSchemaIOProvider.java | 3 +-
.../sdk/io/gcp/pubsub/PubsubUnboundedSink.java | 7 +-
.../beam/sdk/io/gcp/spanner/ReadOperation.java | 4 +-
.../apache/beam/sdk/io/gcp/spanner/SpannerIO.java | 7 +-
.../changestreams/ChangeStreamsConstants.java | 1 +
.../action/DataChangeRecordAction.java | 4 +-
.../dao/PartitionMetadataAdminDao.java | 10 +
.../mapper/ChangeStreamRecordMapper.java | 17 +-
.../beam/sdk/io/gcp/testing/TableContainer.java | 3 +-
.../gcp/bigquery/BigQueryIOStorageQueryTest.java | 88 +++---
.../io/gcp/bigquery/BigQueryIOStorageReadTest.java | 4 +-
.../StorageApiDataTriggeredSchemaUpdateIT.java | 3 +-
.../gcp/bigquery/TableRowToStorageApiProtoIT.java | 6 +-
.../gcp/bigquery/providers/BigQueryManagedIT.java | 3 +-
.../io/gcp/bigtable/BigtableSharedClientTest.java | 3 +-
.../gcp/firestore/BaseFirestoreV1WriteFnTest.java | 7 +-
...storeV1FnBatchWriteWithDeadLetterQueueTest.java | 3 +-
.../FirestoreV1FnBatchWriteWithSummaryTest.java | 3 +-
.../firestore/FirestoreV1FnPartitionQueryTest.java | 3 +-
.../gcp/firestore/FirestoreV1FnRunQueryTest.java | 10 +-
.../beam/sdk/io/gcp/firestore/QueryUtilsTest.java | 3 +-
.../beam/sdk/io/gcp/firestore/RpcQosTest.java | 3 +-
.../sdk/io/gcp/firestore/it/BaseFirestoreIT.java | 5 +-
.../sdk/io/gcp/healthcare/HL7v2IOTestUtil.java | 2 +
.../beam/sdk/io/gcp/pubsub/PubsubIOTest.java | 3 +-
.../beam/sdk/io/gcp/pubsub/PubsubWriteIT.java | 3 +-
.../sdk/io/gcp/spanner/SpannerAccessorTest.java | 36 +--
.../sdk/io/gcp/spanner/SpannerIOWriteTest.java | 9 +-
.../beam/sdk/io/gcp/spanner/SpannerReadIT.java | 3 +-
.../beam/sdk/io/gcp/spanner/SpannerWriteIT.java | 3 +-
.../beam/sdk/io/gcp/spanner/StructUtilsTest.java | 3 +-
.../it/ChangeStreamTestPipelineOptions.java | 3 +-
.../it/SpannerChangeStreamPlacementTableIT.java | 8 +-
...pannerChangeStreamPlacementTablePostgresIT.java | 8 +-
.../it/SpannerChangeStreamPostgresIT.java | 8 +-
sdks/java/io/hadoop-file-system/build.gradle | 2 +-
sdks/java/io/hadoop-format/build.gradle | 2 +-
.../beam/sdk/io/hadoop/format/HadoopFormatIO.java | 6 +-
.../hadoop/format/HadoopFormatIOCassandraIT.java | 1 +
.../io/hadoop/format/HadoopFormatIOElasticIT.java | 1 +
.../hadoop/format/HadoopFormatIOElasticTest.java | 1 +
.../io/hadoop/format/HadoopFormatIOReadTest.java | 1 +
.../apache/beam/sdk/io/hcatalog/HCatalogIO.java | 3 +-
.../org/apache/beam/sdk/io/iceberg/AddFiles.java | 3 +-
.../iceberg/AssignDestinationsAndPartitions.java | 3 +-
.../apache/beam/sdk/io/iceberg/FilterUtils.java | 1 +
.../org/apache/beam/sdk/io/iceberg/IcebergIO.java | 3 +-
.../beam/sdk/io/iceberg/RecordWriterManager.java | 1 +
.../beam/sdk/io/iceberg/SerializableDataFile.java | 4 +-
.../sdk/io/iceberg/cdc/ApplyWatermarkColumn.java | 3 +-
.../beam/sdk/io/iceberg/cdc/CdcReadUtils.java | 3 +-
.../beam/sdk/io/iceberg/cdc/ChangelogScanner.java | 12 +-
.../beam/sdk/io/iceberg/cdc/OverlapRange.java | 3 +-
.../beam/sdk/io/iceberg/cdc/ResolveChanges.java | 3 +-
.../apache/beam/sdk/io/iceberg/ReadUtilsTest.java | 7 +-
.../sdk/io/iceberg/SerializableTableSpecTest.java | 3 +-
.../io/iceberg/catalog/IcebergCatalogBaseIT.java | 3 +-
.../apache/beam/sdk/io/influxdb/InfluxDbIO.java | 1 +
.../java/org/apache/beam/sdk/io/jdbc/JdbcIO.java | 6 +-
.../beam/sdk/io/jdbc/JdbcSchemaIOProvider.java | 3 +-
.../org/apache/beam/sdk/io/jdbc/SchemaUtil.java | 4 +-
.../ReadFromPostgresSchemaTransformProvider.java | 3 +-
.../apache/beam/sdk/io/jms/JmsCheckpointMark.java | 9 +-
.../java/org/apache/beam/sdk/io/jms/JmsIO.java | 10 +-
.../java/org/apache/beam/sdk/io/jms/CommonJms.java | 8 +-
.../beam/sdk/io/kafka/KafkaExactlyOnceSink.java | 4 +-
.../java/org/apache/beam/sdk/io/kafka/KafkaIO.java | 26 +-
.../KafkaIOReadImplementationCompatibility.java | 3 +-
.../org/apache/beam/sdk/io/kafka/KafkaMetrics.java | 3 +-
.../org/apache/beam/sdk/io/kafka/KafkaIOIT.java | 3 +-
.../org/apache/beam/sdk/io/kafka/KafkaIOTest.java | 1 +
.../beam/sdk/io/kafka/ReadFromKafkaDoFnTest.java | 3 +-
.../apache/beam/sdk/io/mongodb/UpdateField.java | 3 +-
.../apache/beam/sdk/io/mongodb/MongoDBIOIT.java | 7 +-
.../java/org/apache/beam/sdk/io/neo4j/Neo4jIO.java | 8 +-
.../org/apache/beam/sdk/io/parquet/ParquetIO.java | 4 +-
.../beam/sdk/io/rabbitmq/ExchangeTestPlan.java | 4 +-
.../beam/io/requestresponse/RequestResponseIO.java | 13 +-
.../RedisExternalResourcesRule.java | 3 +-
.../singlestore/SingleStoreDefaultRowMapper.java | 21 +-
.../apache/beam/sdk/io/snowflake/SnowflakeIO.java | 3 +-
.../org/apache/beam/sdk/io/solace/SolaceIO.java | 11 +-
.../solace/broker/BasicAuthSempClientFactory.java | 1 +
.../sdk/io/solace/broker/MessageProducerUtils.java | 16 +-
.../solace/broker/SempBasicAuthClientExecutor.java | 4 +-
.../io/solace/broker/SolaceMessageProducer.java | 4 +-
.../org/apache/beam/sdk/io/solace/data/Solace.java | 189 +++++++++++--
.../beam/sdk/io/solace/read/WatermarkPolicy.java | 4 +-
.../beam/sdk/io/solace/SolaceIOWriteTest.java | 83 ++++++
.../sdk/io/solace/data/SolaceRecordMapperTest.java | 313 +++++++++++++++++++++
.../beam/sdk/io/solace/data/SolaceRecordTest.java} | 31 +-
.../it/BasicAuthMultipleSempClientFactory.java | 1 +
.../beam/sdk/io/sparkreceiver/HasOffset.java | 3 +-
.../beam/sdk/io/sparkreceiver/SparkConsumer.java | 3 +-
.../sdk/io/sparkreceiver/SparkReceiverIOIT.java | 5 +-
.../beam/sdk/io/splunk/HttpEventPublisher.java | 10 +-
.../sdk/io/splunk/CustomX509TrustManagerTest.java | 7 +-
.../beam/sdk/io/splunk/HttpEventPublisherTest.java | 20 +-
.../sdk/io/synthetic/SyntheticBoundedSource.java | 3 +-
.../beam/sdk/io/thrift/TestThriftStruct.java | 12 +-
.../apache/beam/sdk/io/thrift/TestThriftUnion.java | 12 +-
.../java/org/apache/beam/sdk/io/xml/XmlIO.java | 8 +-
.../org/apache/beam/sdk/managed/ManagedTest.java | 6 +-
.../ml/inference/openai/OpenAIModelParameters.java | 1 +
.../java/org/apache/beam/sdk/jpmstests/JpmsIT.java | 1 +
.../beam/sdk/loadtests/CoGroupByKeyLoadTest.java | 3 +-
.../apache/beam/sdk/loadtests/LoadTestOptions.java | 15 +-
.../apache/beam/sdk/nexmark/NexmarkLauncher.java | 5 +
.../apache/beam/sdk/nexmark/NexmarkOptions.java | 159 ++++-------
.../sdk/nexmark/queries/AbstractSimulator.java | 3 +-
.../apache/beam/sdk/nexmark/queries/Query10.java | 7 +-
.../apache/beam/sdk/nexmark/queries/Query8.java | 3 +-
.../beam/sdk/nexmark/queries/WinningBids.java | 3 +-
.../testutils/publishing/InfluxDBPublisher.java | 8 +-
.../org/apache/beam/sdk/tpcds/TpcdsOptions.java | 18 +-
.../sdk/transformservice/ExpansionService.java | 6 +-
.../anomaly_detection_pipeline/setup.py | 2 +-
.../inference/pytorch_image_object_detection.py | 3 +-
website/www/yarn.lock | 65 +++--
543 files changed, 3506 insertions(+), 1826 deletions(-)
create mode 100644
runners/flink/2.0/src/main/java/org/apache/beam/runners/flink/translation/wrappers/streaming/CachedSideInputReader.java
create mode 100644
runners/flink/2.0/src/main/java/org/apache/beam/runners/flink/translation/wrappers/streaming/SideInputCache.java
create mode 100644
runners/flink/2.0/src/test/java/org/apache/beam/runners/flink/translation/wrappers/streaming/FlinkCachedSideInputReaderTest.java
create mode 100644
runners/spark/4/src/test/java/org/apache/beam/runners/spark/structuredstreaming/translation/SparkKryoRegistratorStreamingTest.java
create mode 100644
runners/spark/src/test/java/org/apache/beam/runners/spark/structuredstreaming/SparkStructuredStreamingPipelineOptionsTest.java
create mode 100644
sdks/java/io/solace/src/test/java/org/apache/beam/sdk/io/solace/data/SolaceRecordMapperTest.java
copy
sdks/java/{core/src/test/java/org/apache/beam/sdk/coders/StructuralByteArrayTest.java
=>
io/solace/src/test/java/org/apache/beam/sdk/io/solace/data/SolaceRecordTest.java}
(56%)