You are viewing a plain text version of this content. The canonical link for it is here.
Posted to commits@beam.apache.org by am...@apache.org on 2017/01/21 15:17:03 UTC
[1/3] beam git commit: [BEAM-1291] KafkaIO: don't log warnig in
offset fetcher while closing.
Repository: beam
Updated Branches:
refs/heads/master f799a57af -> 09d131ced
[BEAM-1291] KafkaIO: don't log warnig in offset fetcher while closing.
add partition to the message. break when the reader is closed.
Project: http://git-wip-us.apache.org/repos/asf/beam/repo
Commit: http://git-wip-us.apache.org/repos/asf/beam/commit/9396c5d4
Tree: http://git-wip-us.apache.org/repos/asf/beam/tree/9396c5d4
Diff: http://git-wip-us.apache.org/repos/asf/beam/diff/9396c5d4
Branch: refs/heads/master
Commit: 9396c5d4dad5cdd38945a368f2ba9951a67bc9f4
Parents: f799a57
Author: Raghu Angadi <ra...@google.com>
Authored: Fri Jan 20 13:24:09 2017 -0800
Committer: Sela <an...@paypal.com>
Committed: Sat Jan 21 16:57:03 2017 +0200
----------------------------------------------------------------------
.../src/main/java/org/apache/beam/sdk/io/kafka/KafkaIO.java | 6 +++++-
1 file changed, 5 insertions(+), 1 deletion(-)
----------------------------------------------------------------------
http://git-wip-us.apache.org/repos/asf/beam/blob/9396c5d4/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 735b8e7..2de2174 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
@@ -1058,7 +1058,11 @@ public class KafkaIO {
long offset = offsetConsumer.position(p.topicPartition);
p.setLatestOffset(offset);
} catch (Exception e) {
- LOG.warn("{}: exception while fetching latest offsets. ignored.", this, e);
+ if (closed.get()) {
+ break;
+ }
+ LOG.warn("{}: exception while fetching latest offset for partition {}. will be retried.",
+ this, p.topicPartition, e);
p.setLatestOffset(UNINITIALIZED_OFFSET); // reset
}
[2/3] beam git commit: added comment on ignore exception.
Posted by am...@apache.org.
added comment on ignore exception.
Project: http://git-wip-us.apache.org/repos/asf/beam/repo
Commit: http://git-wip-us.apache.org/repos/asf/beam/commit/2e9cde24
Tree: http://git-wip-us.apache.org/repos/asf/beam/tree/2e9cde24
Diff: http://git-wip-us.apache.org/repos/asf/beam/diff/2e9cde24
Branch: refs/heads/master
Commit: 2e9cde243ba2841f4f9a2864d2a9292d0f672814
Parents: 9396c5d
Author: Sela <an...@paypal.com>
Authored: Sat Jan 21 17:01:14 2017 +0200
Committer: Sela <an...@paypal.com>
Committed: Sat Jan 21 17:01:14 2017 +0200
----------------------------------------------------------------------
.../kafka/src/main/java/org/apache/beam/sdk/io/kafka/KafkaIO.java | 3 ++-
1 file changed, 2 insertions(+), 1 deletion(-)
----------------------------------------------------------------------
http://git-wip-us.apache.org/repos/asf/beam/blob/2e9cde24/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 2de2174..36ab1fd 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
@@ -1058,11 +1058,12 @@ public class KafkaIO {
long offset = offsetConsumer.position(p.topicPartition);
p.setLatestOffset(offset);
} catch (Exception e) {
+ // An exception is expected if we've closed the reader in another thread. Ignore and exit.
if (closed.get()) {
break;
}
LOG.warn("{}: exception while fetching latest offset for partition {}. will be retried.",
- this, p.topicPartition, e);
+ this, p.topicPartition, e);
p.setLatestOffset(UNINITIALIZED_OFFSET); // reset
}
[3/3] beam git commit: This closes #1804
Posted by am...@apache.org.
This closes #1804
Project: http://git-wip-us.apache.org/repos/asf/beam/repo
Commit: http://git-wip-us.apache.org/repos/asf/beam/commit/09d131ce
Tree: http://git-wip-us.apache.org/repos/asf/beam/tree/09d131ce
Diff: http://git-wip-us.apache.org/repos/asf/beam/diff/09d131ce
Branch: refs/heads/master
Commit: 09d131ced040bfe58d067632b0818f67c097022d
Parents: f799a57 2e9cde2
Author: Sela <an...@paypal.com>
Authored: Sat Jan 21 17:02:45 2017 +0200
Committer: Sela <an...@paypal.com>
Committed: Sat Jan 21 17:02:45 2017 +0200
----------------------------------------------------------------------
.../src/main/java/org/apache/beam/sdk/io/kafka/KafkaIO.java | 7 ++++++-
1 file changed, 6 insertions(+), 1 deletion(-)
----------------------------------------------------------------------