This is an automated email from the ASF dual-hosted git repository.
lcwik pushed a change to branch master
in repository https://gitbox.apache.org/repos/asf/beam.git.
from 8a1406d Merge pull request #12881 [BEAM-7746] Get mypy passing on
runners.worker
new 80aefdd [BEAM-10670] Make Read use SDF by default. Override in
runners.
new 098817e fixup! Fix unit test failures that were missed.
new 9108d3b [BEAM-10670] Don't start/finish bundles when there are no
timers that are ready.
new fd4190d fixup! Fix spotbugs/checkstyle warning
new 82cfa29 Merge pull request #13006 from lukecwik/beam10670.4
The 29128 revisions listed above as "new" are entirely new to this
repository and will be described in separate emails. The revisions
listed as "add" were already present in the repository and have only
been added to this reference.
Summary of changes:
.../runners/core/construction/ReadTranslation.java | 37 +++--
.../runners/core/construction/SplittableParDo.java | 180 +++++++++++++++++----
.../core/construction/PipelineTranslationTest.java | 40 ++++-
.../core/construction/ReadTranslationTest.java | 6 +-
.../core/construction/SplittableParDoTest.java | 90 ++++++-----
.../construction/graph/QueryablePipelineTest.java | 49 +++---
.../renderer/PipelineDotRendererTest.java | 41 +++--
.../direct/BoundedReadEvaluatorFactory.java | 4 +-
.../apache/beam/runners/direct/DirectRunner.java | 2 +-
.../runners/direct/UnboundedReadDeduplicator.java | 4 +-
.../direct/UnboundedReadEvaluatorFactory.java | 10 +-
.../direct/BoundedReadEvaluatorFactoryTest.java | 11 +-
.../direct/UnboundedReadEvaluatorFactoryTest.java | 31 ++--
.../flink/FlinkBatchPipelineTranslator.java | 9 --
.../org/apache/beam/runners/flink/FlinkRunner.java | 2 +-
.../flink/FlinkStreamingPipelineTranslator.java | 7 -
.../FlinkStreamingTransformTranslatorsTest.java | 14 +-
.../dataflow/DataflowPipelineTranslator.java | 3 +-
.../beam/runners/dataflow/DataflowRunner.java | 8 +-
.../beam/runners/dataflow/ReadTranslator.java | 7 +-
.../dataflow/worker/WorkerCustomSourcesTest.java | 7 +
runners/samza/build.gradle | 1 +
.../org/apache/beam/runners/samza/SamzaRunner.java | 2 +-
.../apache/beam/runners/samza/runtime/DoFnOp.java | 11 +-
.../beam/runners/samza/runtime/GroupByKeyOp.java | 12 +-
.../SplittableParDoProcessKeyedElementsOp.java | 11 +-
.../runners/samza/translation/ReadTranslator.java | 8 +-
.../runners/spark/SparkNativePipelineVisitor.java | 18 ++-
.../org/apache/beam/runners/spark/SparkRunner.java | 3 +
.../beam/runners/spark/SparkRunnerDebugger.java | 2 +
.../apache/beam/runners/spark/io/SourceRDD.java | 3 +-
.../SparkStructuredStreamingRunner.java | 2 +
.../translation/batch/PipelineTranslatorBatch.java | 5 +-
.../streaming/PipelineTranslatorStreaming.java | 5 +-
.../spark/translation/TransformTranslator.java | 9 +-
.../streaming/StreamingTransformTranslator.java | 9 +-
.../beam/runners/twister2/Twister2Runner.java | 4 +-
.../batch/ReadSourceTranslatorBatch.java | 8 +-
.../src/main/java/org/apache/beam/sdk/io/Read.java | 154 +++++++-----------
.../org/apache/beam/sdk/io/TFRecordIOTest.java | 15 +-
.../org/apache/beam/sdk/io/TextIOReadTest.java | 10 +-
.../apache/beam/sdk/runners/TransformTreeTest.java | 16 +-
.../sdk/io/gcp/pubsub/PubsubIOExternalTest.java | 26 +--
.../beam/sdk/io/kafka/KafkaIOExternalTest.java | 31 ++--
44 files changed, 545 insertions(+), 382 deletions(-)