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