使用Avro4s生成AoState Avro Schema时遇StackOverflowError求助
问题描述
在Spark中使用Avro4s存储GroupState(需转换为Array[Byte]格式),编写测试用例生成AoState的Avro Schema时,始终触发StackOverflowError。
相关Case Class定义
private[objects] final case class AoState( conf: AnomalyObjectConf, timedScores: Array[TimedScores] = Array(), events: Option[AnomalyObject] = None ) extends CommonGroupState[AnomalyObjectConf, AnomalyObject, Option[AnomalyObject]] private[spe] final case class AnomalyObjectConf( pmuId: Int, asThreshold: Double, aasThreshold: Double, durationThresholdSeconds: Int, resourceIds: Seq[Long], contextIds: Seq[Long], version: Long, updatedTimestamp: Instant, deleteFlag: Boolean = false ) extends CommonConf private[spe] case class TimedScores( creationDate: Instant, asScores: MatrixScore, aasScores: MatrixScore = emptyMatrix ) private[spe] type MatrixScore = util.List[java.lang.Double]
AnomalyObject的Avro Schema
{ "namespace": "com.mycomosi.dto.aiops", "protocol": "aiops", "doc": "AIOPS specific Data Transfer Objects", "types": [ { "type": "record", "name": "AnomalySignaturesSnap", "doc": "Anomaly signatures at specific time", "fields": [ {"name": "timestamp", "type": ["null", { "type": "long", "logicalType": "timestamp-millis"}], "doc": "Measurement date time"}, {"name": "anomalyScores", "type": ["null", {"type": "array", "items": ["null","double"], "default": null}], "doc": "Matrix of anomaly score value encoded into an Array. The matrix is array encoded. It requires resourceIds & contextIds arrays for decoding AS.The value AS for the Xth resource and Yth context is available in the array cell Y+X*len(contextIds). The AAS for this PMU is the latest value from the array"} ] }, { "type": "record", "name": "AnomalySignatures", "doc": "Anomaly signatures over time with resource and context ids", "fields": [ {"name": "resourceIds", "type": {"type": "array", "items": "long"}, "doc": "Ids of the resources"}, {"name": "contextIds", "type": {"type": "array", "items": "long"}, "doc": "Ids of the measurement contexts"}, {"name": "anomalyScores", "type": { "type": "array", "items": ["null","com.mycomosi.dto.aiops.AnomalySignaturesSnap"]}, "doc": "Anomaly signatures over time"} ] }, { "type": "record", "name": "AnomalyObject", "doc": "Anomaly object status with signature", "fields": [ {"name": "aoId", "type": { "type": "string", "logicalType": "UUID"}, "doc": "Anomaly object unique Id"}, {"name": "pmuId", "type": "int", "doc": "PMU unique Id"}, {"name": "status", "type": {"type": "enum", "name": "AoStatus", "symbols" : ["OPEN", "CLOSED"]}, "doc": "Anomaly object status"}, {"name": "startedTimestamp", "type": { "type": "long", "logicalType": "timestamp-millis"}, "doc": "Anomaly object opening date"}, {"name": "closedTimestamp", "type": ["null", { "type": "long", "logicalType": "timestamp-millis"}], "default": null, "doc": "Anomaly object closure date"}, {"name": "updatedTimestamp", "type": { "type": "long", "logicalType": "timestamp-millis"}, "doc": "Update timestamp"}, {"name": "pmuAnomalyScore", "type":"double", "doc": "Latest PMU aggregated anomaly score"}, {"name": "signatureDtls", "type":"com.mycomosi.dto.aiops.AnomalySignatures", "doc": "Anomaly signature"} ] } ] }
测试代码
import com.mycomosi.bda.spe.spark.anomaly.objects.AnomalyObjectsDtos.AoState import com.sksamuel.avro4s.{AvroSchema, SchemaFor} import org.scalatest.funsuite.AnyFunSuite class AnomalyObjectAvroTest extends AnyFunSuite { test("Test AOState Avro Schema") { implicit lazy val schemaForAoState: SchemaFor[AoState] = SchemaFor[AoState] val schema = AvroSchema[AoState] println(schema.toString(true)) } }
报错信息
java.lang.StackOverflowError at com.mycomosi.bda.spe.spark.anomaly.objects.AnomalyObjectAvroTest.schemaForAoState$1(AnomalyObjectAvroTest.scala:10) at com.mycomosi.bda.spe.spark.anomaly.objects.AnomalyObjectAvroTest.schemaForAoState$lzycompute$1(AnomalyObjectAvroTest.scala:10) at com.mycomosi.bda.spe.spark.anomaly.objects.AnomalyObjectAvroTest.schemaForAoState$1(AnomalyObjectAvroTest.scala:10) at com.mycomosi.bda.spe.spark.anomaly.objects.AnomalyObjectAvroTest.schemaForAoState$lzycompute$1(AnomalyObjectAvroTest.scala:10) ... 重复递归调用 ...
报错原因
测试代码中错误定义了循环依赖的隐式SchemaFor[AoState]:
implicit lazy val schemaForAoState: SchemaFor[AoState] = SchemaFor[AoState]
这段代码在初始化时会不断调用自身来获取SchemaFor[AoState],导致无限递归,最终触发栈溢出。此外,Instant、MatrixScore等非基础类型可能需要显式Schema映射,否则也会影响Avro4s的自动推导流程。
解决方案
1. 移除错误的隐式SchemaFor定义
直接删除测试中自定义的schemaForAoState,让Avro4s自动推导Schema:
import com.mycomosi.bda.spe.spark.anomaly.objects.AnomalyObjectsDtos.AoState import com.sksamuel.avro4s.AvroSchema import org.scalatest.funsuite.AnyFunSuite class AnomalyObjectAvroTest extends AnyFunSuite { test("Test AOState Avro Schema") { val schema = AvroSchema[AoState] println(schema.toString(true)) } }
2. 为特殊类型提供显式SchemaFor
如果Avro4s无法自动推导Instant、MatrixScore等类型的Schema,需手动定义隐式实例:
import java.time.Instant import java.util import com.sksamuel.avro4s.{AvroSchema, SchemaFor} import org.apache.avro.Schema // 为Instant映射timestamp-millis类型 implicit val instantSchemaFor: SchemaFor[Instant] = SchemaFor[Instant]( Schema.create(Schema.Type.LONG), (instant: Instant) => instant.toEpochMilli, (millis: Long) => Instant.ofEpochMilli(millis) ) // 为MatrixScore映射double数组类型 implicit val matrixScoreSchemaFor: SchemaFor[util.List[java.lang.Double]] = SchemaFor.list[java.lang.Double]
3. 关联手动定义的AnomalyObject Schema
如果AnomalyObject已有手动定义的Avro Schema,需让Avro4s复用该Schema,避免重复推导:
// 加载本地的AnomalyObject Schema文件 val anomalyObjectSchema = new Schema.Parser().parse(getClass.getResourceAsStream("/anomaly-object.avsc")) // 为AnomalyObject提供SchemaFor实例 implicit val anomalyObjectSchemaFor: SchemaFor[AnomalyObject] = SchemaFor[AnomalyObject](anomalyObjectSchema)
将上述隐式定义放在测试类内部或伴生对象中,确保Avro4s能正常获取,即可生成正确的AoState Schema。
内容的提问来源于stack exchange,提问作者Chandan Gawri
相关产品推荐
相关产品推荐

