You are viewing a plain text version of this content. The canonical link for it is here.
Posted to commits@beam.apache.org by gi...@apache.org on 2022/09/23 05:32:53 UTC
[beam] branch nightly-refs/heads/master updated (3fe0867cd67 -> 5af5a1fb201)
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 3fe0867cd67 Allow longer Class-Path entries (#23269)
add cf790b50cc5 Do not use .get() on ValueProvider during pipeline creation
add 6b7d8b19299 Merge pull request #23294: SpannerIO - Do not use .get() on ValueProvider during pipeline creation
add 762edd7f3a6 Improved pipeline translation in SparkStructuredStreamingRunner (#22446)
add a6cda1370b3 use avro DataFileReader to read avro container files
add 483a0c95734 Merge pull request #23214: Use avro DataFileReader to read avro container files
add 5af5a1fb201 Change google_cloud_bigdataoss_version to 2.2.8. (#23300)
No new revisions were added by this update.
Summary of changes:
CHANGES.md | 1 +
.../org/apache/beam/gradle/BeamModulePlugin.groovy | 2 +-
.../translation/helpers/EncoderFactory.java | 12 +-
.../utils/{package-info.java => ScalaInterop.java} | 24 +-
.../spark/structuredstreaming/Constants.java | 25 -
.../SparkStructuredStreamingRunner.java | 51 +-
.../io/BoundedDatasetFactory.java | 324 +++++++++++
.../structuredstreaming}/io/package-info.java | 2 +-
.../metrics/WithMetricsSupport.java | 2 +-
.../translation/AbstractTranslationContext.java | 235 --------
.../translation/PipelineTranslator.java | 57 +-
.../translation/TransformTranslator.java | 198 ++++++-
.../translation/TranslationContext.java | 124 ++++-
.../translation/batch/AggregatorCombiner.java | 270 ----------
.../translation/batch/Aggregators.java | 591 +++++++++++++++++++++
.../batch/CombineGloballyTranslatorBatch.java | 121 +++++
.../batch/CombinePerKeyTranslatorBatch.java | 181 ++++---
.../CreatePCollectionViewTranslatorBatch.java | 24 +-
.../translation/batch/DatasetSourceBatch.java | 240 ---------
.../translation/batch/DoFnFunction.java | 164 ------
.../batch/DoFnMapPartitionsFactory.java | 224 ++++++++
.../translation/batch/FlattenTranslatorBatch.java | 60 +--
.../translation/batch/GroupByKeyHelpers.java | 106 ++++
.../batch/GroupByKeyTranslatorBatch.java | 298 +++++++++--
.../translation/batch/ImpulseTranslatorBatch.java | 26 +-
.../translation/batch/ParDoTranslatorBatch.java | 315 ++++++-----
.../translation/batch/PipelineTranslatorBatch.java | 28 +-
.../translation/batch/ProcessContext.java | 138 -----
.../batch/ReadSourceTranslatorBatch.java | 76 +--
.../batch/ReshuffleTranslatorBatch.java | 30 --
.../batch/WindowAssignTranslatorBatch.java | 90 +++-
.../translation/helpers/CoderHelpers.java | 10 +-
.../translation/helpers/EncoderFactory.java | 71 ++-
.../translation/helpers/EncoderHelpers.java | 546 ++++++++++++++++++-
.../translation/helpers/KVHelpers.java | 31 --
.../translation/helpers/MultiOutputCoder.java | 84 ---
.../translation/helpers/RowHelpers.java | 75 ---
.../translation/helpers/SchemaHelpers.java | 39 --
.../translation/helpers/WindowingHelpers.java | 82 ---
.../streaming/DatasetSourceStreaming.java | 25 -
.../streaming/PipelineTranslatorStreaming.java | 93 ----
.../streaming/ReadSourceTranslatorStreaming.java | 87 ---
.../translation/streaming/package-info.java | 20 -
.../translation/utils/ScalaInterop.java | 114 ++++
.../aggregators/metrics/sink/InMemoryMetrics.java | 2 +-
.../translation/batch/AggregatorsTest.java | 370 +++++++++++++
.../translation/batch/CombineGloballyTest.java} | 129 ++---
.../{CombineTest.java => CombinePerKeyTest.java} | 92 ++--
.../translation/batch/ComplexSourceTest.java | 15 +-
.../translation/batch/FlattenTest.java | 12 +-
.../translation/batch/GroupByKeyTest.java | 152 ++++--
.../translation/batch/ParDoTest.java | 54 +-
.../translation/batch/SimpleSourceTest.java | 12 +-
.../translation/batch/WindowAssignTest.java | 12 +-
.../translation/helpers/EncoderHelpersTest.java | 210 +++++++-
.../runners/spark/SparkCommonPipelineOptions.java | 6 +
.../beam/runners/spark/SparkPipelineOptions.java | 6 -
.../java/org/apache/beam/sdk/io/AvroSource.java | 343 +++---------
.../org/apache/beam/sdk/io/AvroSourceTest.java | 166 ------
.../beam/sdk/io/gcp/spanner/BatchSpannerRead.java | 6 +-
.../beam/sdk/io/gcp/spanner/SpannerConfig.java | 13 +-
61 files changed, 4076 insertions(+), 2840 deletions(-)
copy runners/spark/2/src/main/java/org/apache/beam/runners/spark/structuredstreaming/translation/utils/{package-info.java => ScalaInterop.java} (62%)
delete mode 100644 runners/spark/3/src/main/java/org/apache/beam/runners/spark/structuredstreaming/Constants.java
create mode 100644 runners/spark/3/src/main/java/org/apache/beam/runners/spark/structuredstreaming/io/BoundedDatasetFactory.java
copy runners/spark/{src/main/java/org/apache/beam/runners/spark => 3/src/main/java/org/apache/beam/runners/spark/structuredstreaming}/io/package-info.java (93%)
delete mode 100644 runners/spark/3/src/main/java/org/apache/beam/runners/spark/structuredstreaming/translation/AbstractTranslationContext.java
delete mode 100644 runners/spark/3/src/main/java/org/apache/beam/runners/spark/structuredstreaming/translation/batch/AggregatorCombiner.java
create mode 100644 runners/spark/3/src/main/java/org/apache/beam/runners/spark/structuredstreaming/translation/batch/Aggregators.java
create mode 100644 runners/spark/3/src/main/java/org/apache/beam/runners/spark/structuredstreaming/translation/batch/CombineGloballyTranslatorBatch.java
delete mode 100644 runners/spark/3/src/main/java/org/apache/beam/runners/spark/structuredstreaming/translation/batch/DatasetSourceBatch.java
delete mode 100644 runners/spark/3/src/main/java/org/apache/beam/runners/spark/structuredstreaming/translation/batch/DoFnFunction.java
create mode 100644 runners/spark/3/src/main/java/org/apache/beam/runners/spark/structuredstreaming/translation/batch/DoFnMapPartitionsFactory.java
create mode 100644 runners/spark/3/src/main/java/org/apache/beam/runners/spark/structuredstreaming/translation/batch/GroupByKeyHelpers.java
delete mode 100644 runners/spark/3/src/main/java/org/apache/beam/runners/spark/structuredstreaming/translation/batch/ProcessContext.java
delete mode 100644 runners/spark/3/src/main/java/org/apache/beam/runners/spark/structuredstreaming/translation/batch/ReshuffleTranslatorBatch.java
delete mode 100644 runners/spark/3/src/main/java/org/apache/beam/runners/spark/structuredstreaming/translation/helpers/KVHelpers.java
delete mode 100644 runners/spark/3/src/main/java/org/apache/beam/runners/spark/structuredstreaming/translation/helpers/MultiOutputCoder.java
delete mode 100644 runners/spark/3/src/main/java/org/apache/beam/runners/spark/structuredstreaming/translation/helpers/RowHelpers.java
delete mode 100644 runners/spark/3/src/main/java/org/apache/beam/runners/spark/structuredstreaming/translation/helpers/SchemaHelpers.java
delete mode 100644 runners/spark/3/src/main/java/org/apache/beam/runners/spark/structuredstreaming/translation/helpers/WindowingHelpers.java
delete mode 100644 runners/spark/3/src/main/java/org/apache/beam/runners/spark/structuredstreaming/translation/streaming/DatasetSourceStreaming.java
delete mode 100644 runners/spark/3/src/main/java/org/apache/beam/runners/spark/structuredstreaming/translation/streaming/PipelineTranslatorStreaming.java
delete mode 100644 runners/spark/3/src/main/java/org/apache/beam/runners/spark/structuredstreaming/translation/streaming/ReadSourceTranslatorStreaming.java
delete mode 100644 runners/spark/3/src/main/java/org/apache/beam/runners/spark/structuredstreaming/translation/streaming/package-info.java
create mode 100644 runners/spark/3/src/main/java/org/apache/beam/runners/spark/structuredstreaming/translation/utils/ScalaInterop.java
create mode 100644 runners/spark/3/src/test/java/org/apache/beam/runners/spark/structuredstreaming/translation/batch/AggregatorsTest.java
copy runners/spark/{2/src/test/java/org/apache/beam/runners/spark/structuredstreaming/translation/batch/CombineTest.java => 3/src/test/java/org/apache/beam/runners/spark/structuredstreaming/translation/batch/CombineGloballyTest.java} (56%)
rename runners/spark/3/src/test/java/org/apache/beam/runners/spark/structuredstreaming/translation/batch/{CombineTest.java => CombinePerKeyTest.java} (71%)