You are viewing a plain text version of this content. The canonical link for it is here.
Posted to commits@spark.apache.org by zs...@apache.org on 2017/03/10 01:42:12 UTC
spark git commit: [SPARK-19886] Fix reportDataLoss if statement in SS
KafkaSource
Repository: spark
Updated Branches:
refs/heads/master f79371ad8 -> 82138e09b
[SPARK-19886] Fix reportDataLoss if statement in SS KafkaSource
## What changes were proposed in this pull request?
Fix the `throw new IllegalStateException` if statement part.
## How is this patch tested
Regression test
Author: Burak Yavuz <br...@gmail.com>
Closes #17228 from brkyvz/kafka-cause-fix.
Project: http://git-wip-us.apache.org/repos/asf/spark/repo
Commit: http://git-wip-us.apache.org/repos/asf/spark/commit/82138e09
Tree: http://git-wip-us.apache.org/repos/asf/spark/tree/82138e09
Diff: http://git-wip-us.apache.org/repos/asf/spark/diff/82138e09
Branch: refs/heads/master
Commit: 82138e09b9ad8d9609d5c64d6c11244b8f230be7
Parents: f79371a
Author: Burak Yavuz <br...@gmail.com>
Authored: Thu Mar 9 17:42:10 2017 -0800
Committer: Shixiong Zhu <sh...@databricks.com>
Committed: Thu Mar 9 17:42:10 2017 -0800
----------------------------------------------------------------------
.../sql/kafka010/CachedKafkaConsumer.scala | 33 +++++++++++--------
.../sql/kafka010/CachedKafkaConsumerSuite.scala | 34 ++++++++++++++++++++
2 files changed, 54 insertions(+), 13 deletions(-)
----------------------------------------------------------------------
http://git-wip-us.apache.org/repos/asf/spark/blob/82138e09/external/kafka-0-10-sql/src/main/scala/org/apache/spark/sql/kafka010/CachedKafkaConsumer.scala
----------------------------------------------------------------------
diff --git a/external/kafka-0-10-sql/src/main/scala/org/apache/spark/sql/kafka010/CachedKafkaConsumer.scala b/external/kafka-0-10-sql/src/main/scala/org/apache/spark/sql/kafka010/CachedKafkaConsumer.scala
index 15b2825..6d76904 100644
--- a/external/kafka-0-10-sql/src/main/scala/org/apache/spark/sql/kafka010/CachedKafkaConsumer.scala
+++ b/external/kafka-0-10-sql/src/main/scala/org/apache/spark/sql/kafka010/CachedKafkaConsumer.scala
@@ -273,19 +273,7 @@ private[kafka010] case class CachedKafkaConsumer private(
message: String,
cause: Throwable = null): Unit = {
val finalMessage = s"$message ${additionalMessage(failOnDataLoss)}"
- if (failOnDataLoss) {
- if (cause != null) {
- throw new IllegalStateException(finalMessage)
- } else {
- throw new IllegalStateException(finalMessage, cause)
- }
- } else {
- if (cause != null) {
- logWarning(finalMessage)
- } else {
- logWarning(finalMessage, cause)
- }
- }
+ reportDataLoss0(failOnDataLoss, finalMessage, cause)
}
private def close(): Unit = consumer.close()
@@ -398,4 +386,23 @@ private[kafka010] object CachedKafkaConsumer extends Logging {
consumer
}
}
+
+ private def reportDataLoss0(
+ failOnDataLoss: Boolean,
+ finalMessage: String,
+ cause: Throwable = null): Unit = {
+ if (failOnDataLoss) {
+ if (cause != null) {
+ throw new IllegalStateException(finalMessage, cause)
+ } else {
+ throw new IllegalStateException(finalMessage)
+ }
+ } else {
+ if (cause != null) {
+ logWarning(finalMessage, cause)
+ } else {
+ logWarning(finalMessage)
+ }
+ }
+ }
}
http://git-wip-us.apache.org/repos/asf/spark/blob/82138e09/external/kafka-0-10-sql/src/test/scala/org/apache/spark/sql/kafka010/CachedKafkaConsumerSuite.scala
----------------------------------------------------------------------
diff --git a/external/kafka-0-10-sql/src/test/scala/org/apache/spark/sql/kafka010/CachedKafkaConsumerSuite.scala b/external/kafka-0-10-sql/src/test/scala/org/apache/spark/sql/kafka010/CachedKafkaConsumerSuite.scala
new file mode 100644
index 0000000..7aa7dd0
--- /dev/null
+++ b/external/kafka-0-10-sql/src/test/scala/org/apache/spark/sql/kafka010/CachedKafkaConsumerSuite.scala
@@ -0,0 +1,34 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.spark.sql.kafka010
+
+import org.scalatest.PrivateMethodTester
+
+import org.apache.spark.sql.test.SharedSQLContext
+
+class CachedKafkaConsumerSuite extends SharedSQLContext with PrivateMethodTester {
+
+ test("SPARK-19886: Report error cause correctly in reportDataLoss") {
+ val cause = new Exception("D'oh!")
+ val reportDataLoss = PrivateMethod[Unit]('reportDataLoss0)
+ val e = intercept[IllegalStateException] {
+ CachedKafkaConsumer.invokePrivate(reportDataLoss(true, "message", cause))
+ }
+ assert(e.getCause === cause)
+ }
+}
---------------------------------------------------------------------
To unsubscribe, e-mail: commits-unsubscribe@spark.apache.org
For additional commands, e-mail: commits-help@spark.apache.org