PySpark读取HDFS中CSV文件报错求助(找不到python3)
问题描述
HDFS的/user/spark路径下已存在archivo.csv文件,但运行PySpark代码读取该文件时报错,现提供执行命令、代码、错误日志及文件列表,请求排查原因及代码问题。
执行命令
export PYSPARK_PYTHON=python3 export PYSPARK_DRIVER_PYTHON=python3 spark-submit --queue=OID Proceso_Match1.py
Python代码
import os import sys from pyspark.sql import HiveContext, Row from pyspark import SparkContext, SparkConf from pyspark.sql.functions import * from pyspark.sql import functions as F from pyspark.sql.types import * if __name__ =='__main__': conf=SparkConf().setAppName("Spark RDD").set("spark.speculation","true") sc=SparkContext(conf=conf) sc.setLogLevel("OFF") sqlContext = HiveContext(sc) #rddCentral = sc.textFile("hdfs:///user/spark/archivo.csv") rddCentral = sc.textFile("/user/spark/archivo.csv") rddCentralMap = rddCentral.map(lambda line : line.split(",")) print('paso 1') dfCentral = sqlContext.createDataFrame(rddCentralMap, ["ROWID_CDR","DURACION","FECHA_LLAMADA","FECHA_LLAMADA_2","MATCH"]) dfCentral=dfCentral.withColumn("FECHA_LLAMADA_NUM",dfCentral.FECHA_LLAMADA_2.cast(IntegerType())) dfCentral=dfCentral.withColumn("DURACION_NUM",dfCentral.DURACION.cast(IntegerType())) dfCentral=dfCentral.withColumn("MATCH_NUM",dfCentral.MATCH.cast(IntegerType())) sc.stop()
错误日志
22/09/30 12:49:14 INFO cluster.YarnClientSchedulerBackend: SchedulerBackend is ready for scheduling beginning after reached minRegisteredResourcesRatio: 0.8 paso 1 /usr/local/bin/python3/lib/python3.7/site-packages/pandas/compat/__init__.py:124: UserWarning: Could not import the lzma module. Your installed Python is incomplete. Attempting to use lzma compression will result in a RuntimeError. warnings.warn(msg) Traceback (most recent call last): File "/home/aic_proceso_vfs/rjaimea/vfs_504/bin/Proceso_Match1.py", line 21, in <module> dfCentral = sqlContext.createDataFrame(rddCentralMap, ["ROWID_CDR","DURACION","FECHA_LLAMADA","FECHA_LLAMADA_2","MATCH"]) : org.apache.spark.SparkException: 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, cl-hdp-cdp-dn7.cse-cph.int, executor 1): java.io.IOException: Cannot run program "python3": error=2, No such file or directory at java.lang.ProcessBuilder.start(ProcessBuilder.java:1048) Caused by: java.io.IOException: error=2, No such file or directory ... 16 more
HDFS文件列表
hdfs dfs -ls /user/spark Found 3 items drwxr-xr-x - spark hdfs 0 2022-07-25 10:11 /user/spark/.sparkStaging -rw------- 3 hadoopadmin hdfs 21 2022-09-30 12:25 /user/spark/archivo.csv drwxrwxrwt - spark spark 0 2022-09-30 12:33 /user/spark/driverLogs
问题排查与解决
核心报错原因
错误日志明确提示Cannot run program "python3": error=2, No such file or directory,说明Spark Executor节点上未安装python3,或python3的路径未被Executor识别。
提交命令中设置的PYSPARK_PYTHON=python3仅在Driver节点生效,Executor节点无法继承该环境变量,且本身找不到python3可执行文件,导致任务执行失败。
解决步骤
检查Executor节点的Python环境
- 登录到Executor节点(如日志中的
cl-hdp-cdp-dn7.cse-cph.int),执行which python3确认是否存在python3。若不存在,需在所有Executor节点安装python3。 - 若python3存在但路径非默认(如
/usr/local/bin/python3),记录完整路径备用。
- 登录到Executor节点(如日志中的
修改任务提交命令
通过--conf参数将Python路径传递给所有Executor节点,确保全局生效,示例:export PYSPARK_PYTHON=python3 export PYSPARK_DRIVER_PYTHON=python3 spark-submit --queue=OID \ --conf spark.executorEnv.PYSPARK_PYTHON=/usr/bin/python3 \ --conf spark.executorEnv.PYSPARK_DRIVER_PYTHON=/usr/bin/python3 \ Proceso_Match1.py注意替换
/usr/bin/python3为实际的python3路径。代码优化建议
- 改用
SparkSession替代SparkContext+HiveContext,这是Spark 2.x及以后的标准API,更简洁且功能完善:from pyspark.sql import SparkSession from pyspark.sql.types import * if __name__ == '__main__': spark = SparkSession.builder \ .appName("Spark RDD") \ .config("spark.speculation", "true") \ .enableHiveSupport() \ .getOrCreate() spark.sparkContext.setLogLevel("OFF") # 直接读取CSV并指定Schema,避免手动RDD转换 dfCentral = spark.read.csv( "/user/spark/archivo.csv", header=False, schema=StructType([ StructField("ROWID_CDR", StringType()), StructField("DURACION", StringType()), StructField("FECHA_LLAMADA", StringType()), StructField("FECHA_LLAMADA_2", StringType()), StructField("MATCH", StringType()) ]) ) # 类型转换逻辑保持不变 dfCentral = dfCentral.withColumn("FECHA_LLAMADA_NUM", dfCentral.FECHA_LLAMADA_2.cast(IntegerType())) dfCentral = dfCentral.withColumn("DURACION_NUM", dfCentral.DURACION.cast(IntegerType())) dfCentral = dfCentral.withColumn("MATCH_NUM", dfCentral.MATCH.cast(IntegerType())) spark.stop() - 使用
spark.read.csv直接读取文件,无需手动处理RDD拆分,减少出错概率,同时可直接指定Schema,避免类型推断问题。
- 改用
内容的提问来源于stack exchange,提问作者Jeniffer Lorena Guillen Castro
相关产品推荐
相关产品推荐

