Spark SQL中如何生成Column类型Schema以调用from_csv函数?
Spark Java中from_csv的Schema参数问题解决
问题场景
在Windows 11环境下使用spark-3.4.1-hadoop3,尝试生成Schema传入from_csv函数时抛出异常。原代码如下:
import org.apache.spark.sql.Column; import org.apache.spark.sql.Dataset; import org.apache.spark.sql.Row; import org.apache.spark.sql.SparkSession; import static org.apache.spark.sql.functions.col; import static org.apache.spark.sql.functions.from_csv; import static org.apache.spark.sql.functions.not; import java.util.HashMap; import java.util.Map; SparkSession spark = SparkSession.builder().appName("FromCsvStructExample").getOrCreate(); Dataset<Row> df = spark.read().format("csv") .option("header", "true") .option("inferSchema", "true") .load("/path/to/csv/file"); Map<String, String> options = new HashMap<String, String>(); String schemaString = "name string, age int, job string"; // 错误用法:用col()包裹schema字符串 Column schema = from_csv(col("csv"), col(schemaString), options); Dataset<Row> parsed = df.select(schema.as("data")); parsed.printSchema(); spark.close();
抛出的异常:
Exception in thread "main" org.apache.spark.sql.AnalysisException: [INVALID_SCHEMA.NON_STRING_LITERAL] The input schema "name string, age int, job string" is not a valid schema string. The input expression must be string literal and not null. at org.apache.spark.sql.errors.QueryCompilationErrors$.unexpectedSchemaTypeError(QueryCompilationErrors.scala:1055) at org.apache.spark.sql.catalyst.expressions.ExprUtils$.evalTypeExpr(ExprUtils.scala:42) at org.apache.spark.sql.catalyst.expressions.ExprUtils$.evalSchemaExpr(ExprUtils.scala:47) at org.apache.spark.sql.catalyst.expressions.CsvToStructs.<init>(csvExpressions.scala:72) at org.apache.spark.sql.functions$.from_csv(functions.scala:4955) at org.apache.spark.sql.functions.from_csv(functions.scala) at com.aaa.etl.processor.Test_CSV.main(Test_CSV.java:43)
错误原因
col(schemaString)的作用是引用DataFrame中名为name string, age int, job string的列,而非将字符串本身作为Schema字面量传入。而from_csv要求Schema参数为字符串字面量或能解析为字面量的表达式,因此触发异常。
解决方案
方案1:直接传入字符串字面量
使用from_csv的重载方法from_csv(Column e, String schema, java.util.Map<String,String> options),直接将schema字符串作为参数传入,无需用col()包裹:
// 替换原错误行 Column schema = from_csv(col("csv"), schemaString, options);
方案2:用lit()生成Column类型的Schema字面量
如果必须使用Column类型的Schema参数(比如Schema需要动态生成并作为字面量传入),可以用lit()函数将字符串转换为字面量列:
// 先导入lit函数 import static org.apache.spark.sql.functions.lit; // 生成Schema字面量列 Column schemaCol = lit(schemaString); Column schema = from_csv(col("csv"), schemaCol, options);
方案3:使用StructType类型的Schema
无需Scala基础,Java中可以直接构造StructType类型的Schema,调用对应的from_csv重载方法:
// 导入类型相关类 import org.apache.spark.sql.types.DataTypes; import org.apache.spark.sql.types.StructField; import org.apache.spark.sql.types.StructType; // 构造StructType Schema StructType structSchema = DataTypes.createStructType(new StructField[]{ DataTypes.createStructField("name", DataTypes.StringType, true), // true表示允许为空 DataTypes.createStructField("age", DataTypes.IntegerType, true), DataTypes.createStructField("job", DataTypes.StringType, true) }); // 传入StructType类型的Schema Column schema = from_csv(col("csv"), structSchema, options);
内容的提问来源于stack exchange,提问作者Joseph Hwang
相关产品推荐
相关产品推荐

