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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.15 05:17:04