PySpark读取HDFS中Sqoop导入数据失败及相关技术咨询
问题分析与解决
一、Sqoop导入命令验证
你的Sqoop导入命令逻辑没问题,但可以做以下确认:
-m 1指定单Map任务,生成单个part-m-00000文件,和后续读取路径匹配,无需调整。- 先通过HDFS命令确认文件是否成功导入且有内容:
如果文件为空或不存在,说明Sqoop导入失败,需检查MySQL连接信息、表权限或网络连通性,可添加# 查看目录下文件列表 hdfs dfs -ls /user/root/Spar_Nord # 查看文件前10行内容 hdfs dfs -head /user/root/Spar_Nord/part-m-00000--verbose参数重新执行导入,查看详细报错日志。
二、PySpark读取卡顿的解决思路
你的读取代码存在几个性能隐患,调整后可解决卡顿问题:
- 关闭自动Schema推断
inferSchema = True会让Spark全量扫描文件推断字段类型,大文件场景下直接导致卡顿。先关闭该参数快速读取,确认数据格式后再手动定义Schema:
确认格式后,手动指定Schema(示例):# 先快速加载数据 df = spark.read.csv("/user/root/Spar_Nord/part-m-00000", header=False, inferSchema=False) # 查看前5行确认格式 df.show(5)from pyspark.sql.types import StructType, StructField, StringType, DoubleType, TimestampType custom_schema = StructType([ StructField("trans_id", StringType(), nullable=True), StructField("account_num", StringType(), nullable=True), StructField("amount", DoubleType(), nullable=True), StructField("trans_time", TimestampType(), nullable=True) # 其他字段按实际数据定义 ]) df = spark.read.csv("/user/root/Spar_Nord/part-m-00000", header=False, schema=custom_schema) - 读取目录而非单个文件
Spark支持直接读取目录下所有part文件,无需指定单个文件名:df = spark.read.csv("/user/root/Spar_Nord", header=False, inferSchema=False) - 调整Driver内存
Jupyter内核卡顿可能是Driver内存不足,启动SparkSession时可增加内存配置:from pyspark.sql import SparkSession spark = SparkSession.builder \ .appName("ReadSqoopData") \ .config("spark.driver.memory", "4g") \ .getOrCreate()
三、确定Sqoop导入文件类型的方法
Sqoop默认导出分隔文本文件,可通过以下方式确认具体格式:
- 查看文件内容
用HDFS命令查看文件片段,判断字段分隔符:hdfs dfs -head /user/root/Spar_Nord/part-m-00000- 若字段用逗号分隔,就是CSV格式;若用制表符(
\t)分隔,就是TSV,读取时需指定sep="\t"。
- 若字段用逗号分隔,就是CSV格式;若用制表符(
- 确认Sqoop导出参数
Sqoop默认用--as-textfile导出文本文件,若未指定--as-parquetfile/--as-sequencefile等参数,就是文本格式。也可在导入时显式指定分隔符:sqoop import --connect jdbc:mysql://upgraddetest.cyaielc9bmnf.us-east-1.rds.amazonaws.com/testdatabase --table SRC_ATM_TRANS --username student --password STUDENT123 --target-dir /user/root/Spar_Nord -m 1 --fields-terminated-by ',' --lines-terminated-by '\n' - 检查文件扩展名
若导出为Parquet等格式,文件会带有.parquet后缀,此时需用对应方式读取:df = spark.read.parquet("/user/root/Spar_Nord")
内容的提问来源于stack exchange,提问作者Manjari
相关产品推荐
相关产品推荐

