Spark GroupBy报StateSchemaNotCompatible异常:新旧状态键Schema不兼容
Spark流处理EventHub聚合时的StateSchemaNotCompatible异常解决
在Spark中读取并写入EventHub事件时,尝试通过以下代码基于指定键进行聚合:
val df1 = df0 .groupBy( colKey, colTimestamp ) .agg( collect_list( struct( colCreationTimestamp, colRecordId ) ).as("Records") )
运行时报错:
Caused by: org.apache.spark.sql.execution.streaming.state.StateSchemaNotCompatible: Provided schema doesn't match to the schema for existing state! Please note that Spark allow difference of field name: check count of fields and data type of each field. - Provided key schema: StructType(StructField(Key,StringType,true), StructField(Timestamp,TimestampType,true) - Provided value schema: StructType(StructField(buf,BinaryType,true)) - Existing key schema: StructType(StructField(_1,StringType,true), StructField(_2,TimestampType,true)) - Existing value schema: StructType(StructField(buf,BinaryType,true)) If you want to force running query without schema validation, please set spark.sql.streaming.stateStore.stateSchemaCheck to false. Please note running query with incompatible schema could cause indeterministic behavior. at org.apache.spark.sql.execution.streaming.state.StateSchemaCompatibilityChecker.check(StateSchemaCompatibilityChecker.scala:60) at org.apache.spark.sql.execution.streaming.state.StateStore$.$anonfun$getStateStoreProvider$2(StateStore.scala:487) at scala.runtime.java8.JFunction0$mcV$sp.apply(JFunction0$mcV$sp.java:23) at scala.util.Try$.apply(Try.scala:213) at org.apache.spark.sql.execution.streaming.state.StateStore$.$anonfun$getStateStoreProvider$1(StateStore.scala:487) at scala.collection.mutable.HashMap.getOrElseUpdate(HashMap.scala:86)
排查过程
- 异常未包含具体代码行号,通过报错信息里的键Schema定位到上述聚合代码
- 修改groupBy的键列时,错误信息会对应变化
- 尝试在groupBy前用
df0.select()显式获取所需列,问题依然存在
疑问
- Spark如何获取旧状态键Schema?
- 该如何解决此异常?
更新:问题已解决
EventHub Spark库会将流处理的状态数据存储在checkpoint目录中,本次异常是由旧的状态数据导致的Schema不兼容问题。更换为全新的Checkpoint目录即可解决。
内容的提问来源于stack exchange,提问作者rick
相关产品推荐
相关产品推荐

