You are viewing a plain text version of this content. The canonical link for it is here.
Posted to commits@beam.apache.org by jk...@apache.org on 2017/04/20 23:18:56 UTC

[1/2] beam git commit: Fix erroneous use of .expand() in KafkaIO

Repository: beam
Updated Branches:
  refs/heads/master 8be1dacab -> aa899e4ce


Fix erroneous use of .expand() in KafkaIO


Project: http://git-wip-us.apache.org/repos/asf/beam/repo
Commit: http://git-wip-us.apache.org/repos/asf/beam/commit/02591269
Tree: http://git-wip-us.apache.org/repos/asf/beam/tree/02591269
Diff: http://git-wip-us.apache.org/repos/asf/beam/diff/02591269

Branch: refs/heads/master
Commit: 0259126920058ffcd08a235fee68f7cfd3d6ffe4
Parents: 8be1dac
Author: Eugene Kirpichov <ki...@google.com>
Authored: Thu Apr 20 15:07:58 2017 -0700
Committer: Eugene Kirpichov <ki...@google.com>
Committed: Thu Apr 20 16:18:14 2017 -0700

----------------------------------------------------------------------
 .../src/main/java/org/apache/beam/sdk/io/kafka/KafkaIO.java      | 4 ++--
 1 file changed, 2 insertions(+), 2 deletions(-)
----------------------------------------------------------------------


http://git-wip-us.apache.org/repos/asf/beam/blob/02591269/sdks/java/io/kafka/src/main/java/org/apache/beam/sdk/io/kafka/KafkaIO.java
----------------------------------------------------------------------
diff --git a/sdks/java/io/kafka/src/main/java/org/apache/beam/sdk/io/kafka/KafkaIO.java b/sdks/java/io/kafka/src/main/java/org/apache/beam/sdk/io/kafka/KafkaIO.java
index 47657bb..fbd96eb 100644
--- a/sdks/java/io/kafka/src/main/java/org/apache/beam/sdk/io/kafka/KafkaIO.java
+++ b/sdks/java/io/kafka/src/main/java/org/apache/beam/sdk/io/kafka/KafkaIO.java
@@ -561,8 +561,8 @@ public class KafkaIO {
 
     @Override
     public PCollection<KV<K, V>> expand(PBegin begin) {
-      return read
-          .expand(begin)
+      return begin
+          .apply(read)
           .apply("Remove Kafka Metadata",
               ParDo.of(new DoFn<KafkaRecord<K, V>, KV<K, V>>() {
                 @ProcessElement


[2/2] beam git commit: This closes #2617

Posted by jk...@apache.org.
This closes #2617


Project: http://git-wip-us.apache.org/repos/asf/beam/repo
Commit: http://git-wip-us.apache.org/repos/asf/beam/commit/aa899e4c
Tree: http://git-wip-us.apache.org/repos/asf/beam/tree/aa899e4c
Diff: http://git-wip-us.apache.org/repos/asf/beam/diff/aa899e4c

Branch: refs/heads/master
Commit: aa899e4ce4cb1d360bfdc53fab318f6bb14c42aa
Parents: 8be1dac 0259126
Author: Eugene Kirpichov <ki...@google.com>
Authored: Thu Apr 20 16:18:26 2017 -0700
Committer: Eugene Kirpichov <ki...@google.com>
Committed: Thu Apr 20 16:18:26 2017 -0700

----------------------------------------------------------------------
 .../src/main/java/org/apache/beam/sdk/io/kafka/KafkaIO.java      | 4 ++--
 1 file changed, 2 insertions(+), 2 deletions(-)
----------------------------------------------------------------------