Spark读取Kafka时访问VARIANT列字段遇[INVALID_EXTRACT_BASE_FIELD_TYPE]报错
问题:从Kafka读取JSON数据时提取嵌套字段报错
当不尝试获取嵌套字段时,能得到正常的数据结构。我正从Kafka读取数据并写入表中,问题出在readStream阶段,收到报错:
[INVALID_EXTRACT_BASE_FIELD_TYPE] Can't extract a value from "data". Need a complex type [STRUCT, ARRAY, MAP] but got "VARIANT". SQLSTATE: 42000
以下是我的readStream代码:
df = spark.readStream \ .format("kafka") \ .option("kafka.bootstrap.servers", BOOTSTRAP_SERVERS) \ .option("subscribe", TOPIC) \ .option("startingOffsets", "latest") \ .... \ .load() \ .withColumn("data", parse_json(col("value").cast("string"))) \ .select("data, data:unique_id") .withColumn("timestamp", current_timestamp()) display(df)
解决方法
问题核心是parse_json返回了VARIANT类型部分Spark环境如Databricks默认行为,而.或:的嵌套字段访问语法仅支持STRUCT/ARRAY/MAP这类明确的复杂类型。可以通过两种方式修复:
提前定义Schema,用
from_json替代parse_json
明确指定JSON的结构,将字符串直接解析为STRUCT类型:from pyspark.sql.types import StructType, StructField, StringType # 根据实际JSON结构定义Schema,这里仅示例unique_id字段 data_schema = StructType([ StructField("unique_id", StringType(), nullable=True) # 其他嵌套字段可在此处补充 ]) df = spark.readStream \ .format("kafka") \ .option("kafka.bootstrap.servers", BOOTSTRAP_SERVERS) \ .option("subscribe", TOPIC) \ .option("startingOffsets", "latest") \ .... \ .load() \ .withColumn("data", from_json(col("value").cast("string"), data_schema)) \ .select("data", "data.unique_id") \ .withColumn("timestamp", current_timestamp()) display(df)用
get_json_object直接提取字段
无需提前定义Schema,直接从JSON字符串中定位提取目标字段:df = spark.readStream \ .format("kafka") \ .option("kafka.bootstrap.servers", BOOTSTRAP_SERVERS) \ .option("subscribe", TOPIC) \ .option("startingOffsets", "latest") \ .... \ .load() \ .withColumn("json_str", col("value").cast("string")) \ .select( parse_json(col("json_str")).alias("data"), get_json_object(col("json_str"), "$.unique_id").alias("unique_id") ) \ .withColumn("timestamp", current_timestamp()) display(df)
另外注意原代码中select("data, data:unique_id")的语法错误,Spark中访问嵌套字段应使用.而非:,即data.unique_id。
内容的提问来源于stack exchange,提问作者Climbs_lika_Spyder
相关产品推荐
相关产品推荐

