Databricks中创建通用函数读取CSV/Parquet文件至DataFrame时出现数据源找不到错误的求助
你遇到的错误java.lang.ClassNotFoundException: Failed to find data source: fileType,核心原因是Spark在尝试加载名为**"fileType"**的数据源,而不是你传入的"csv"或"parquet"。这说明你的format(fileType)调用中,fileType变量并没有被正确解析为你传入的参数值,反而被当成了字面量字符串"fileType"。
下面是具体的排查和解决步骤:
1. 先确认fileType的实际值
在函数内部添加调试打印语句,查看fileType变量的真实内容,这是最直接的排查方式:
def createBronzeDeltaTable(schemaExcelPath, schemaExcelTab ,sourceCsvPath,DestiNationDeltaFilePath,dataBaseLocationPath,delemeterType,databaseTable,dbColumnToOptimize,fileType): # 添加调试打印,确认fileType的值 print(f"Current fileType value: {repr(fileType)}") # 先处理schema,避免eval影响变量 schema_def = create_schema(schemaExcelPath, schemaExcelTab) print(f"Generated schema: {schema_def}") materialMasterDF_schema = spark.read.option("badRecordsPath", mountPoint+"/tmp/badRecordsPath")\ .format(fileType)\ .load(sourceCsvPath, header=True, schema=eval(schema_def), sep=delemeterType) # 后续代码...
运行后如果打印结果是Current fileType value: 'fileType',说明你调用函数时没有正确传递fileType参数;如果是'csv'或'parquet',再继续排查其他问题。
2. 检查函数调用的参数传递
你的函数参数列表较长,很容易出现参数顺序错误或关键字参数拼写错误:
- 确保调用时
fileType是最后一个参数(因为函数定义里它是第9个参数),或者明确使用关键字参数传递,避免顺序混乱:
# 推荐用关键字参数传递,避免顺序错误 createBronzeDeltaTable( schemaExcelPath="/path/to/schema.xlsx", schemaExcelTab="Sheet1", sourceCsvPath="/path/to/source/files", DestiNationDeltaFilePath="/path/to/delta", dataBaseLocationPath="/path/to/db", delemeterType=",", databaseTable="bronze_table", dbColumnToOptimize="id", fileType="csv" # 明确指定fileType参数 )
- 注意参数名的大小写:Python是大小写敏感的,确保调用时的参数名是
fileType,而不是filetype或FileType。
3. 排查是否有变量被意外覆盖
你的代码中使用了eval(create_schema(...)),如果create_schema函数返回的字符串中包含对fileType变量的赋值,会意外覆盖函数参数的值。比如如果create_schema返回的内容里有fileType = "something",eval执行后会把你的fileType参数改掉。
通过打印schema_def的值可以确认这一点,如果发现问题,建议修改create_schema函数,避免返回包含变量赋值的内容,或者改用更安全的方式解析schema(比如用pyspark.sql.types直接构建Schema,而不是用eval)。
4. 额外注意:未定义的mountPoint变量
你的代码中mountPoint变量没有定义,后续运行到badRecordsPath时会报错,建议补充定义这个变量,比如:
mountPoint = "/mnt/your_mount_name" # 根据你的实际挂载路径修改
修改后的完整示例函数
def createBronzeDeltaTable(schemaExcelPath, schemaExcelTab ,sourceCsvPath,DestiNationDeltaFilePath,dataBaseLocationPath,delemeterType,databaseTable,dbColumnToOptimize,fileType): # 定义mountPoint变量 mountPoint = "/mnt/your_mount_name" # 调试打印 print(f"Current fileType value: {repr(fileType)}") # 预生成schema并打印 schema_def = create_schema(schemaExcelPath, schemaExcelTab) print(f"Generated schema definition: {schema_def}") materialMasterDF_schema = spark.read.option("badRecordsPath", f"{mountPoint}/tmp/badRecordsPath")\ .format(fileType)\ .load(sourceCsvPath, header=True, schema=eval(schema_def), sep=delemeterType) materialMasterDF_schema.printSchema() display(materialMasterDF_schema)
内容的提问来源于stack exchange,提问作者sayan nandi

