Spark保存Salesforce TASK数据DataFrame报ArrayIndexOutOfBoundsException:1求助
问题排查指导:Salesforce TASK对象写入Parquet抛出ArrayIndexOutOfBoundsException
该错误的核心根因为Spark DataFrame声明的schema字段数大于实际读取到的Row的字段数,Spark访问不存在的字段索引时抛出数组越界
你可以按照以下步骤逐步定位问题:
- 第一步:校验自定义Schema与SOQL返回字段的一致性
分别打印retrievingSOQL中SELECT子句的字段列表,以及你构造的sfdcObjectSchema.fieldNames,对比两者的字段数量、字段顺序、字段名大小写:Salesforce Bulk API返回的字段名大小写可能和你从元数据中取的Name属性大小写不一致,且SpringML的Salesforce连接器对字段顺序、大小写敏感,只要两者不匹配就会触发该错误。 - 第二步:校验元数据过滤逻辑的准确性
针对TASK对象单独打印nonCompoundMetas.size以及nonCompoundMetas.map(_.Name)的完整结果,确认你过滤复合字段的逻辑没有误删普通字段,也没有漏掉Salesforce Bulk API默认返回的隐藏元字段(比如attributes字段)。 - 第三步:通过自动推断Schema快速定位差异
临时注释掉读数据时的.schema(sfdcObjectSchema)配置,让SpringML连接器自动推断返回数据的Schema,打印自动推断的Schema的字段数量、字段列表,和你自定义的Schema做对比,就能快速找到字段不匹配的点。 - 第四步:定位异常数据行
若Schema对比无差异,添加以下测试代码打印异常行的实际字段数:import org.apache.spark.sql.catalyst.expressions.GenericRow // 打印异常行的实际字段数与声明字段数的差异 sfdcObjectDF.rdd.mapPartitions(iter => { iter.map(row => { try { (row.length, "正常行") } catch { case e: ArrayIndexOutOfBoundsException => (row.asInstanceOf[GenericRow].values.length, "异常行") } }) }).filter(_._2 == "异常行").take(10).foreach(println) - 第五步:针对TASK对象特殊字段定位
Salesforce的TASK对象存在大量关联复合字段(比如WhoId、WhatId会关联多类业务对象),你可以先缩小SOQL查询范围:仅查询Id字段测试写入是否正常,再逐步增加查询字段,定位到触发问题的具体字段。
临时修复方案
如果自动推断的Schema符合业务要求,可以直接去掉自定义Schema的配置,使用连接器自动推断的Schema写入即可;如果必须使用自定义Schema,保证自定义Schema的字段顺序、数量、大小写和SOQL查询返回的字段完全一致即可解决该错误。
内容的提问来源于stack exchange,提问作者Kalpana
相关产品推荐
相关产品推荐

