【问题标题】:Why does my spark application fail in cluster mode but successful in client mode?为什么我的 spark 应用程序在集群模式下失败但在客户端模式下成功?
【发布时间】:2020-07-17 06:54:16
【问题描述】:

我正在尝试在 pyspark 程序下运行,该程序将从 HDFS 集群中复制文件。

import pyspark
from pyspark import SparkConf
from pyspark.sql import SparkSession

def read_file(spark):
    try:
        csv_data = spark.read.csv('hdfs://hostname:port/user/devuser/example.csv')
        csv_data.write.format('csv').save('/tmp/data')
        print('count of csv_data: {}'.format(csv_data.count()))
    except Exception as e:
        print(e)
        return False
    return True

if __name__ == "__main__":
    spark = SparkSession.builder.master('yarn').config('spark.app.name','dummy_App').config('spark.executor.memory','2g').config('spark.executor.cores','2').config('spark.yarn.keytab','/home/devuser/devuser.keytab').config('spark.yarn.principal','devuser@NAME.COM').config('spark.executor.instances','2').config('hadoop.security.authentication','kerberos').config('spark.yarn.access.hadoopFileSystems','hdfs://hostname:port').getOrCreate()
    if read_file(spark):
        print('Read the file successfully..')
    else:
        print('Reading failed..')

如果我使用 spark-submit 以部署模式作为客户端运行上述代码,则作业运行良好,我可以在目录 /tmp/data 中看到输出

spark-submit --master yarn --num-executors 1 --deploy-mode client --executor-memory 1G --executor-cores 1 --driver-memory 1G --files /home/hdfstest/conn_props/core-site_fs.xml,/home/hdfstest/conn_props/hdfs-site_fs.xml check_con.py

但如果我使用--deploy-mode cluster 运行相同的代码,

spark-submit --master yarn --deploy-mode cluster --num-executors 1 --executor-memory 1G --executor-cores 1 --driver-memory 1G --files /home/hdfstest/conn_props/core-site_fs.xml,/home/hdfstest/conn_props/hdfs-site_fs.xml check_con.py

作业因 kerberos 异常而失败,如下所示:

Caused by: org.apache.hadoop.security.AccessControlException: Client cannot authenticate via:[TOKEN, KERBEROS]
        at org.apache.hadoop.security.SaslRpcClient.selectSaslClient(SaslRpcClient.java:173)
        at org.apache.hadoop.security.SaslRpcClient.saslConnect(SaslRpcClient.java:390)
        at org.apache.hadoop.ipc.Client$Connection.setupSaslConnection(Client.java:615)
        at org.apache.hadoop.ipc.Client$Connection.access$2300(Client.java:411)
        at org.apache.hadoop.ipc.Client$Connection$2.run(Client.java:801)
        at org.apache.hadoop.ipc.Client$Connection$2.run(Client.java:797)
        at java.security.AccessController.doPrivileged(Native Method)
        at javax.security.auth.Subject.doAs(Subject.java:422)
        at org.apache.hadoop.security.UserGroupInformation.doAs(UserGroupInformation.java:1730)
        at org.apache.hadoop.ipc.Client$Connection.setupIOstreams(Client.java:797)
        ... 55 more

代码包含我从中读取文件的集群的 keytab 信息。但是我不明白为什么它在集群模式下失败但在客户端模式下运行。

我应该对代码进行任何配置更改以在集群模式下运行它吗?谁能告诉我如何解决这个问题?

编辑 1:我尝试从 spark-submit 传递 keytab 和原理详细信息,而不是在程序中对其进行硬编码,如下所示:

def read_file(spark):
    try:
        csv_data = spark.read.csv('hdfs://hostname:port/user/devuser/example.csv')
        csv_data.write.format('csv').save('/tmp/data')
        print('count of csv_data: {}'.format(csv_data.count()))
    except Exception as e:
        print(e)
        return False
    return True

if __name__ == "__main__":
    spark = SparkSession.builder.master('yarn').config('spark.app.name','dummy_App').config('spark.executor.memory','2g').config('spark.executor.cores','2').config('spark.executor.instances','2').config('hadoop.security.authentication','kerberos').config('spark.yarn.access.hadoopFileSystems','hdfs://hostname:port').getOrCreate()
    if read_file(spark):
        print('Read the file successfully..')
    else:
        print('Reading failed..')

火花提交:

spark-submit --master yarn --deploy-mode cluster --name checkCon --num-executors 1 --executor-memory 1G --executor-cores 1 --driver-memory 1G --files /home/devuser/conn_props/core-site_fs.xml,/home/devuser/conn_props/hdfs-site_fs.xml --principal devuser@NAME.COM --keytab /home/devuser/devuser.keytab check_con.py

例外:

Job aborted due to stage failure: Task 0 in stage 0.0 failed 4 times, most recent failure: Lost task 0.3 in stage 0.0 (TID 3, executor name, executor 2): java.io.IOException: DestHost:destPort <port given in the csv_data statement> , LocalHost:localPort <localport name>. Failed on local exception: java.io.IOException: org.apache.hadoop.security.AccessControlException: Client cannot authenticate via:[TOKEN, KERBEROS]

【问题讨论】:

  • 你可以删除所有这些行.config('spark.executor.memory','2g') .config('spark.executor.cores','2') .config('spark.executor.instances','2') .config('hadoop.security.authentication','kerberos') .config('spark.yarn.access.hadoopFileSystems','hdfs://hostname:port') 并尝试
  • 也将这个spark.read.csv('hdfs://hostname:port/user/devuser/example.csv')更改为spark.read.csv('/user/devuser/example.csv')
  • 使用这个 spark-submit 命令 - spark-submit --master yarn --deploy-mode cluster --name checkCon --num-executors 1 --executor-memory 1G --executor-cores 1 --driver-memory 1g -–conf spark.yarn.keytab=/home/devuser/devuser.keytab -–conf spark.yarn.principal=devuser@NAME.COM check_con.py & 让我知道它是否不起作用
  • 删除了所有配置,从 spark-submit 传递了 keytab 和 principal。但这项工作仍然失败,同样的例外。
  • keytab & principal 是否有效?

标签: apache-spark hadoop


【解决方案1】:

因为这个.config('spark.yarn.keytab','/home/devuser/devuser.keytab') conf。 在client 模式下,您的作业将在本地运行且给定路径可用,因此作业成功完成。

cluster 模式下,/home/devuser/devuser.keytab 在数据节点和驱动程序中不可用或无法访问,因此它失败了。

SparkSession\
.builder\
.master('yarn')\
.config('spark.app.name','dummy_App')\
.config('spark.executor.memory','2g')\
.config('spark.executor.cores','2')\
.config('spark.yarn.keytab','/home/devuser/devuser.keytab')\ # This line is causing problem
.config('spark.yarn.principal','devuser@NAME.COM')\
.config('spark.executor.instances','2')\
.config('hadoop.security.authentication','kerberos')\
.config('spark.yarn.access.hadoopFileSystems','hdfs://hostname:port')\
.getOrCreate()

不要硬编码 spark.yarn.keytab & spark.yarn.principal 配置。 将这些配置作为spark-submit 命令的一部分传递。

spark-submit --class ${APP_MAIN_CLASS} \
    --master yarn \
    --deploy-mode cluster \
    --name ${APP_INSTANCE} \
    --files ${APP_BASE_DIR}/conf/${ENV_NAME}/env.conf,${APP_BASE_DIR}/conf/example-application.conf \
    -–conf spark.yarn.keytab=path_to_keytab \
    -–conf spark.yarn.principal=principal@REALM.COM \
    --jars ${JARS} \
    [...]

【讨论】:

  • 还是不行。我从代码中删除了 keytab 详细信息并像这样spark-submit --master yarn --deploy-mode cluster --name checkCon --num-executors 1 --executor-memory 1G --executor-cores 1 --driver-memory 1G --files /path/core-site_fs.xml,/path/hdfs-site_fs.xml --principal devuser@NAME.COM --keytab /home/devuser/devuser.keytab check_con.py 提交,但仍然发现相同的异常。我已经更新了问题中的相同细节。
  • 抱歉,我已更改配置 - -–conf spark.yarn.keytab=path_to_keytab -–conf spark.yarn.principal=principal@REALM.COM 尝试添加这些并再测试一次。
【解决方案2】:

阿法伊克, 以下方法将解决您的问题,

KERBEROS_KEYTAB_PATH=/home/devuser/devuser.keytabKERBEROS_PRINCIPAL=devuser@NAME.COM

方法一:使用 kinit 命令

第 1 步:启动并继续 spark-submit

kinit -kt ${KERBEROS_KEYTAB_PATH} ${KERBEROS_PRINCIPAL}

第 2 步:运行 klist 并验证 Kerberization 对登录的 devuser 是否正常工作。

Ticket cache: FILE:/tmp/krb5cc_XXXXXXXXX_XXXXXX
Default principal: devuser@NAME.COM

Valid starting       Expires              Service principal
07/30/2020 15:52:28  07/31/2020 01:52:28  krbtgt/NAME.COM@NAME.COM
        renew until 08/06/2020 15:52:28

第 3 步:用 spark session 替换 spark 代码

sparkSession = SparkSession.builder().config(sparkConf).appName("TEST1").enableHiveSupport().getOrCreate()

第 4 步:运行 spark-submit,如下所示。

 $SPARK_HOME/bin/spark-submit --class com.test.load.Data \
     --master yarn \
     --deploy-mode cluster \
     --driver-memory 2g \
     --executor-memory 2g \
     --executor-cores 2 --num-executors 2 \
     --conf "spark.driver.cores=2" \
     --name "TEST1" \
     --principal ${KERBEROS_PRINCIPAL} \
     --keytab ${KERBEROS_KEYTAB_PATH} \
     --conf spark.files=$SPARK_HOME/conf/hive-site.xml \
     /home/devuser/sparkproject/Test-jar-1.0.jar 2> /home/devuser/logs/test1.log

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2019-11-15
    • 2015-10-22
    • 1970-01-01
    • 2023-04-08
    • 1970-01-01
    • 2019-09-15
    • 2019-08-13
    • 1970-01-01
    相关资源
    最近更新 更多