You are viewing a plain text version of this content. The canonical link for it is here.
Posted to commits@beam.apache.org by mx...@apache.org on 2019/03/18 11:28:33 UTC

[beam] branch master updated (e3bf740 -> 34fed44)

This is an automated email from the ASF dual-hosted git repository.

mxm pushed a change to branch master
in repository https://gitbox.apache.org/repos/asf/beam.git.


    from e3bf740  Merge pull request #8065: [BEAM-4684] Disable support of @RequiresStableInput on Dataflow runner for now
     new f1ae50b  [BEAM-6751] Add KafkaIO EOS support to Flink via @RequiresStableInput
     new d07316e  Add @NonnullByDefault annotation to stableinput package
     new 34fed44  Merge pull request #7991: [BEAM-6751] Add KafkaIO EOS support to Flink via @RequiresStableInput

The 20567 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/flink/flink_runner.gradle                  |   1 +
 .../runners/flink/FlinkTransformOverrides.java     |   7 -
 .../wrappers/streaming/DoFnOperator.java           |  72 ++++-
 .../BufferedElement.java}                          |  11 +-
 .../streaming/stableinput/BufferedElements.java    | 167 ++++++++++++
 .../streaming/stableinput/BufferingDoFnRunner.java | 209 +++++++++++++++
 .../BufferingElementsHandler.java}                 |  29 +--
 .../stableinput/KeyedBufferingElementsHandler.java | 109 ++++++++
 .../NonKeyedBufferingElementsHandler.java}         |  42 +--
 .../{io => stableinput}/package-info.java          |   8 +-
 .../flink/FlinkRequiresStableInputTest.java        | 252 ++++++++++++++++++
 .../wrappers/streaming/DoFnOperatorTest.java       | 290 +++++++++++++++++++++
 .../stableinput/BufferedElementsTest.java          |  75 ++++++
 .../org/apache/beam/sdk/RequiresStableInputIT.java |  30 ++-
 .../beam/sdk/io/kafka/KafkaExactlyOnceSink.java    |   1 +
 .../java/org/apache/beam/sdk/io/kafka/KafkaIO.java |   3 +-
 16 files changed, 1239 insertions(+), 67 deletions(-)
 copy runners/flink/src/main/java/org/apache/beam/runners/flink/translation/wrappers/streaming/{package-info.java => stableinput/BufferedElement.java} (71%)
 create mode 100644 runners/flink/src/main/java/org/apache/beam/runners/flink/translation/wrappers/streaming/stableinput/BufferedElements.java
 create mode 100644 runners/flink/src/main/java/org/apache/beam/runners/flink/translation/wrappers/streaming/stableinput/BufferingDoFnRunner.java
 copy runners/flink/src/main/java/org/apache/beam/runners/flink/translation/wrappers/streaming/{PushedBackElementsHandler.java => stableinput/BufferingElementsHandler.java} (56%)
 create mode 100644 runners/flink/src/main/java/org/apache/beam/runners/flink/translation/wrappers/streaming/stableinput/KeyedBufferingElementsHandler.java
 copy runners/flink/src/main/java/org/apache/beam/runners/flink/translation/wrappers/streaming/{NonKeyedPushedBackElementsHandler.java => stableinput/NonKeyedBufferingElementsHandler.java} (57%)
 copy runners/flink/src/main/java/org/apache/beam/runners/flink/translation/wrappers/streaming/{io => stableinput}/package-info.java (77%)
 create mode 100644 runners/flink/src/test/java/org/apache/beam/runners/flink/FlinkRequiresStableInputTest.java
 create mode 100644 runners/flink/src/test/java/org/apache/beam/runners/flink/translation/wrappers/streaming/stableinput/BufferedElementsTest.java