You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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可执行文件,导致任务执行失败。

解决步骤

  1. 检查Executor节点的Python环境

    • 登录到Executor节点(如日志中的cl-hdp-cdp-dn7.cse-cph.int),执行which python3确认是否存在python3。若不存在,需在所有Executor节点安装python3。
    • 若python3存在但路径非默认(如/usr/local/bin/python3),记录完整路径备用。
  2. 修改任务提交命令
    通过--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路径。

  3. 代码优化建议

    • 改用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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.17 21:40:38