From d31dcb340eba66d6adf6c027675fb9ebb5b18ce2 Mon Sep 17 00:00:00 2001 From: jerryshao Date: Fri, 17 Mar 2017 17:14:56 +0800 Subject: [PATCH 1/4] Using real user to initialize hive SessionState Change-Id: If423f3fdc709ed3284cafc01efd1fe389f635560 --- .../apache/spark/deploy/SparkHadoopUtil.scala | 20 +++++++++++++++ .../security/HiveCredentialProvider.scala | 25 +++---------------- .../sql/hive/client/HiveClientImpl.scala | 13 +++++++++- 3 files changed, 35 insertions(+), 23 deletions(-) diff --git a/core/src/main/scala/org/apache/spark/deploy/SparkHadoopUtil.scala b/core/src/main/scala/org/apache/spark/deploy/SparkHadoopUtil.scala index f475ce87540aa..6829c1f32370c 100644 --- a/core/src/main/scala/org/apache/spark/deploy/SparkHadoopUtil.scala +++ b/core/src/main/scala/org/apache/spark/deploy/SparkHadoopUtil.scala @@ -18,6 +18,7 @@ package org.apache.spark.deploy import java.io.IOException +import java.lang.reflect.UndeclaredThrowableException import java.security.PrivilegedExceptionAction import java.text.DateFormat import java.util.{Arrays, Comparator, Date, Locale} @@ -353,6 +354,25 @@ class SparkHadoopUtil extends Logging { } buffer.toString } + + /** + * Run some code as the real logged in user (which may differ from the current user, for + * example, when using proxying). + */ + private[spark] def doAsRealUser[T](fn: => T): T = { + val currentUser = UserGroupInformation.getCurrentUser() + val realUser = Option(currentUser.getRealUser()).getOrElse(currentUser) + + // For some reason the Scala-generated anonymous class ends up causing an + // UndeclaredThrowableException, even if you annotate the method with @throws. + try { + realUser.doAs(new PrivilegedExceptionAction[T]() { + override def run(): T = fn + }) + } catch { + case e: UndeclaredThrowableException => throw Option(e.getCause()).getOrElse(e) + } + } } object SparkHadoopUtil { diff --git a/resource-managers/yarn/src/main/scala/org/apache/spark/deploy/yarn/security/HiveCredentialProvider.scala b/resource-managers/yarn/src/main/scala/org/apache/spark/deploy/yarn/security/HiveCredentialProvider.scala index 16d8fc32bb42d..9b27632c4eb1e 100644 --- a/resource-managers/yarn/src/main/scala/org/apache/spark/deploy/yarn/security/HiveCredentialProvider.scala +++ b/resource-managers/yarn/src/main/scala/org/apache/spark/deploy/yarn/security/HiveCredentialProvider.scala @@ -17,9 +17,6 @@ package org.apache.spark.deploy.yarn.security -import java.lang.reflect.UndeclaredThrowableException -import java.security.PrivilegedExceptionAction - import scala.reflect.runtime.universe import scala.util.control.NonFatal @@ -30,6 +27,7 @@ import org.apache.hadoop.security.{Credentials, UserGroupInformation} import org.apache.hadoop.security.token.Token import org.apache.spark.SparkConf +import org.apache.spark.deploy.SparkHadoopUtil import org.apache.spark.internal.Logging import org.apache.spark.util.Utils @@ -87,7 +85,7 @@ private[security] class HiveCredentialProvider extends ServiceCredentialProvider classOf[String], classOf[String]) val getHive = hiveClass.getMethod("get", hiveConfClass) - doAsRealUser { + SparkHadoopUtil.get.doAsRealUser { val hive = getHive.invoke(null, conf) val tokenStr = getDelegationToken.invoke(hive, currentUser.getUserName(), principal) .asInstanceOf[String] @@ -108,22 +106,5 @@ private[security] class HiveCredentialProvider extends ServiceCredentialProvider None } - /** - * Run some code as the real logged in user (which may differ from the current user, for - * example, when using proxying). - */ - private def doAsRealUser[T](fn: => T): T = { - val currentUser = UserGroupInformation.getCurrentUser() - val realUser = Option(currentUser.getRealUser()).getOrElse(currentUser) - - // For some reason the Scala-generated anonymous class ends up causing an - // UndeclaredThrowableException, even if you annotate the method with @throws. - try { - realUser.doAs(new PrivilegedExceptionAction[T]() { - override def run(): T = fn - }) - } catch { - case e: UndeclaredThrowableException => throw Option(e.getCause()).getOrElse(e) - } - } + } diff --git a/sql/hive/src/main/scala/org/apache/spark/sql/hive/client/HiveClientImpl.scala b/sql/hive/src/main/scala/org/apache/spark/sql/hive/client/HiveClientImpl.scala index 989fdc5564d39..cb84d6b3c3e20 100644 --- a/sql/hive/src/main/scala/org/apache/spark/sql/hive/client/HiveClientImpl.scala +++ b/sql/hive/src/main/scala/org/apache/spark/sql/hive/client/HiveClientImpl.scala @@ -36,6 +36,7 @@ import org.apache.hadoop.hive.ql.session.SessionState import org.apache.hadoop.security.UserGroupInformation import org.apache.spark.{SparkConf, SparkException} +import org.apache.spark.deploy.SparkHadoopUtil import org.apache.spark.internal.Logging import org.apache.spark.metrics.source.HiveCatalogMetrics import org.apache.spark.sql.AnalysisException @@ -188,7 +189,17 @@ private[hive] class HiveClientImpl( if (clientLoader.cachedHive != null) { Hive.set(clientLoader.cachedHive.asInstanceOf[Hive]) } - SessionState.start(state) + + // When Security is enabled, using real user to initialize Hive SessionState to avoid tgt + // not found issue with proxy user. + if (UserGroupInformation.isSecurityEnabled) { + SparkHadoopUtil.get.doAsRealUser { + SessionState.start(state) + } + } else { + SessionState.start(state) + } + state.out = new PrintStream(outputBuffer, true, "UTF-8") state.err = new PrintStream(outputBuffer, true, "UTF-8") state From 85a4220a0838f2e8e33d44db28a34cfb4b3453a1 Mon Sep 17 00:00:00 2001 From: jerryshao Date: Tue, 21 Mar 2017 13:53:55 +0800 Subject: [PATCH 2/4] Use delegation tokens to authenticate metastore connection Change-Id: I84897d0b14fc69a68a70a6341e64c4c0a8188cba --- .../apache/spark/deploy/SparkHadoopUtil.scala | 20 --------------- .../org/apache/spark/deploy/yarn/Client.scala | 3 +++ .../security/HiveCredentialProvider.scala | 25 ++++++++++++++++--- .../sql/hive/client/HiveClientImpl.scala | 12 +-------- 4 files changed, 26 insertions(+), 34 deletions(-) diff --git a/core/src/main/scala/org/apache/spark/deploy/SparkHadoopUtil.scala b/core/src/main/scala/org/apache/spark/deploy/SparkHadoopUtil.scala index 6829c1f32370c..f475ce87540aa 100644 --- a/core/src/main/scala/org/apache/spark/deploy/SparkHadoopUtil.scala +++ b/core/src/main/scala/org/apache/spark/deploy/SparkHadoopUtil.scala @@ -18,7 +18,6 @@ package org.apache.spark.deploy import java.io.IOException -import java.lang.reflect.UndeclaredThrowableException import java.security.PrivilegedExceptionAction import java.text.DateFormat import java.util.{Arrays, Comparator, Date, Locale} @@ -354,25 +353,6 @@ class SparkHadoopUtil extends Logging { } buffer.toString } - - /** - * Run some code as the real logged in user (which may differ from the current user, for - * example, when using proxying). - */ - private[spark] def doAsRealUser[T](fn: => T): T = { - val currentUser = UserGroupInformation.getCurrentUser() - val realUser = Option(currentUser.getRealUser()).getOrElse(currentUser) - - // For some reason the Scala-generated anonymous class ends up causing an - // UndeclaredThrowableException, even if you annotate the method with @throws. - try { - realUser.doAs(new PrivilegedExceptionAction[T]() { - override def run(): T = fn - }) - } catch { - case e: UndeclaredThrowableException => throw Option(e.getCause()).getOrElse(e) - } - } } object SparkHadoopUtil { diff --git a/resource-managers/yarn/src/main/scala/org/apache/spark/deploy/yarn/Client.scala b/resource-managers/yarn/src/main/scala/org/apache/spark/deploy/yarn/Client.scala index ccb0f8fdbbc21..b754058c07658 100644 --- a/resource-managers/yarn/src/main/scala/org/apache/spark/deploy/yarn/Client.scala +++ b/resource-managers/yarn/src/main/scala/org/apache/spark/deploy/yarn/Client.scala @@ -369,6 +369,9 @@ private[spark] class Client( // Merge credentials obtained from registered providers val nearestTimeOfNextRenewal = credentialManager.obtainCredentials(hadoopConf, credentials) + // Add credentials to current user's UGI, so that following operations don't need to use the + // Kerberos tgt to get delegations again in the client side. + UserGroupInformation.getCurrentUser.addCredentials(credentials) if (credentials != null) { logDebug(YarnSparkHadoopUtil.get.dumpTokens(credentials).mkString("\n")) diff --git a/resource-managers/yarn/src/main/scala/org/apache/spark/deploy/yarn/security/HiveCredentialProvider.scala b/resource-managers/yarn/src/main/scala/org/apache/spark/deploy/yarn/security/HiveCredentialProvider.scala index 9b27632c4eb1e..16d8fc32bb42d 100644 --- a/resource-managers/yarn/src/main/scala/org/apache/spark/deploy/yarn/security/HiveCredentialProvider.scala +++ b/resource-managers/yarn/src/main/scala/org/apache/spark/deploy/yarn/security/HiveCredentialProvider.scala @@ -17,6 +17,9 @@ package org.apache.spark.deploy.yarn.security +import java.lang.reflect.UndeclaredThrowableException +import java.security.PrivilegedExceptionAction + import scala.reflect.runtime.universe import scala.util.control.NonFatal @@ -27,7 +30,6 @@ import org.apache.hadoop.security.{Credentials, UserGroupInformation} import org.apache.hadoop.security.token.Token import org.apache.spark.SparkConf -import org.apache.spark.deploy.SparkHadoopUtil import org.apache.spark.internal.Logging import org.apache.spark.util.Utils @@ -85,7 +87,7 @@ private[security] class HiveCredentialProvider extends ServiceCredentialProvider classOf[String], classOf[String]) val getHive = hiveClass.getMethod("get", hiveConfClass) - SparkHadoopUtil.get.doAsRealUser { + doAsRealUser { val hive = getHive.invoke(null, conf) val tokenStr = getDelegationToken.invoke(hive, currentUser.getUserName(), principal) .asInstanceOf[String] @@ -106,5 +108,22 @@ private[security] class HiveCredentialProvider extends ServiceCredentialProvider None } - + /** + * Run some code as the real logged in user (which may differ from the current user, for + * example, when using proxying). + */ + private def doAsRealUser[T](fn: => T): T = { + val currentUser = UserGroupInformation.getCurrentUser() + val realUser = Option(currentUser.getRealUser()).getOrElse(currentUser) + + // For some reason the Scala-generated anonymous class ends up causing an + // UndeclaredThrowableException, even if you annotate the method with @throws. + try { + realUser.doAs(new PrivilegedExceptionAction[T]() { + override def run(): T = fn + }) + } catch { + case e: UndeclaredThrowableException => throw Option(e.getCause()).getOrElse(e) + } + } } diff --git a/sql/hive/src/main/scala/org/apache/spark/sql/hive/client/HiveClientImpl.scala b/sql/hive/src/main/scala/org/apache/spark/sql/hive/client/HiveClientImpl.scala index cb84d6b3c3e20..ebc0d1199892d 100644 --- a/sql/hive/src/main/scala/org/apache/spark/sql/hive/client/HiveClientImpl.scala +++ b/sql/hive/src/main/scala/org/apache/spark/sql/hive/client/HiveClientImpl.scala @@ -189,17 +189,7 @@ private[hive] class HiveClientImpl( if (clientLoader.cachedHive != null) { Hive.set(clientLoader.cachedHive.asInstanceOf[Hive]) } - - // When Security is enabled, using real user to initialize Hive SessionState to avoid tgt - // not found issue with proxy user. - if (UserGroupInformation.isSecurityEnabled) { - SparkHadoopUtil.get.doAsRealUser { - SessionState.start(state) - } - } else { - SessionState.start(state) - } - + SessionState.start(state) state.out = new PrintStream(outputBuffer, true, "UTF-8") state.err = new PrintStream(outputBuffer, true, "UTF-8") state From 11a10946a575d6ed0f707ea0735b7a0a0024090d Mon Sep 17 00:00:00 2001 From: jerryshao Date: Tue, 21 Mar 2017 13:55:09 +0800 Subject: [PATCH 3/4] Remove unnecessary import Change-Id: Iaa917493ac596e8497394fa89b900d47a94f7da2 --- .../scala/org/apache/spark/sql/hive/client/HiveClientImpl.scala | 1 - 1 file changed, 1 deletion(-) diff --git a/sql/hive/src/main/scala/org/apache/spark/sql/hive/client/HiveClientImpl.scala b/sql/hive/src/main/scala/org/apache/spark/sql/hive/client/HiveClientImpl.scala index ebc0d1199892d..989fdc5564d39 100644 --- a/sql/hive/src/main/scala/org/apache/spark/sql/hive/client/HiveClientImpl.scala +++ b/sql/hive/src/main/scala/org/apache/spark/sql/hive/client/HiveClientImpl.scala @@ -36,7 +36,6 @@ import org.apache.hadoop.hive.ql.session.SessionState import org.apache.hadoop.security.UserGroupInformation import org.apache.spark.{SparkConf, SparkException} -import org.apache.spark.deploy.SparkHadoopUtil import org.apache.spark.internal.Logging import org.apache.spark.metrics.source.HiveCatalogMetrics import org.apache.spark.sql.AnalysisException From e9b55800e9ad02ef07f081a5dfb4943ac5d80523 Mon Sep 17 00:00:00 2001 From: jerryshao Date: Tue, 21 Mar 2017 15:14:10 +0800 Subject: [PATCH 4/4] Fix NPE bug in the code Change-Id: I6be6be7b1e9a4580e5e1eeab8aac451ea830ef8b --- .../main/scala/org/apache/spark/deploy/yarn/Client.scala | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/resource-managers/yarn/src/main/scala/org/apache/spark/deploy/yarn/Client.scala b/resource-managers/yarn/src/main/scala/org/apache/spark/deploy/yarn/Client.scala index b754058c07658..3218d221143e5 100644 --- a/resource-managers/yarn/src/main/scala/org/apache/spark/deploy/yarn/Client.scala +++ b/resource-managers/yarn/src/main/scala/org/apache/spark/deploy/yarn/Client.scala @@ -369,11 +369,11 @@ private[spark] class Client( // Merge credentials obtained from registered providers val nearestTimeOfNextRenewal = credentialManager.obtainCredentials(hadoopConf, credentials) - // Add credentials to current user's UGI, so that following operations don't need to use the - // Kerberos tgt to get delegations again in the client side. - UserGroupInformation.getCurrentUser.addCredentials(credentials) if (credentials != null) { + // Add credentials to current user's UGI, so that following operations don't need to use the + // Kerberos tgt to get delegations again in the client side. + UserGroupInformation.getCurrentUser.addCredentials(credentials) logDebug(YarnSparkHadoopUtil.get.dumpTokens(credentials).mkString("\n")) }