Azure Databricks写Snowflake遇无字段信息的SnowflakeSQLException如何排查
Databricks 写入 Snowflake 报错无具体列的高效排查方案
1. 写入前预校验schema匹配性
先拉取Snowflake目标表的schema,和待写入DataFrame的schema做对比,提前排查非空约束不匹配、数据类型不匹配的问题,该操作仅拉取0行数据,不会占用过多资源:
// 配置Snowflake连接参数 val sfOptions = Map( "sfURL" -> "<你的Snowflake访问域名>", "sfUser" -> "<用户名>", "sfPassword" -> "<密码>", "sfDatabase" -> "<目标库名>", "sfSchema" -> "<目标schema名>", "sfWarehouse" -> "<计算仓库名>", "dbtable" -> "<目标表名>" ) // 拉取Snowflake目标表的schema val sfTargetSchema = spark.read.format("net.snowflake.spark.snowflake") .options(sfOptions) .load() .limit(0) .schema // 待写入的DataFrame,替换为你自己的DF变量名 val writeDF: org.apache.spark.sql.DataFrame = ???
1.1 排查非空约束不匹配列
val nonNullMismatchCols = sfTargetSchema.fields .filter(f => !f.nullable && writeDF.schema(f.name).nullable) .map(_.name) if(nonNullMismatchCols.nonEmpty) { println(s"检测到非空约束不匹配列:${nonNullMismatchCols.mkString("、")}") }
1.2 排查数据类型不匹配列
val typeMismatchCols = sfTargetSchema.fields .filter(f => sfTargetSchema(f.name).dataType != writeDF.schema(f.name).dataType) .map(_.name) if(typeMismatchCols.nonEmpty) { println(s"检测到数据类型不匹配列:${typeMismatchCols.mkString("、")}") }
2. 逐列校验非法值
针对类型不匹配的列单独扫描异常值,比如时间戳解析失败的场景,可批量检查所有字符串列是否存在无法转换为时间格式的非法值:
import org.apache.spark.sql.functions._ // 批量校验所有字符串列是否存在无法解析为时间戳的异常值 val stringCols = writeDF.schema.fields .filter(_.dataType.typeName == "string") .map(_.name) stringCols.foreach{ colName => val invalidCnt = writeDF .filter(to_timestamp(col(colName)).isNull && col(colName).isNotNull) .count() if(invalidCnt > 0) { println(s"列【$colName】存在 $invalidCnt 条无法解析为时间戳的非法值") } }
也可以根据报错类型调整校验逻辑,比如数值类型列就校验是否存在非数字字符串。
3. 开启连接器DEBUG日志
在Notebook开头调整Snowflake Spark连接器的日志级别为DEBUG,报错时可从Spark Executor日志中获取更详细的字段处理上下文:
import org.apache.log4j.{Level, Logger} Logger.getLogger("net.snowflake.spark.snowflake").setLevel(Level.DEBUG)
4. 小批量分组写入定位
如果前面的方法仍未定位到问题列,可将待写入DF按列拆分分组,逐组写入测试缩小排查范围:
注意:测试写入建议使用Snowflake临时表,避免脏数据写入生产表
// 每5列分为一组测试写入 val colGroups = writeDF.columns.grouped(5).toList colGroups.foreach{ group => val testDF = writeDF.select(group.map(col):_*) try { testDF.write.format("net.snowflake.spark.snowflake") .options(sfOptions) .mode("append") .save() println(s"列组【${group.mkString("、")}】写入正常") } catch { case e: net.snowflake.client.jdbc.SnowflakeSQLException => println(s"列组【${group.mkString("、")}】写入报错,问题列在该组范围内") } }
内容的提问来源于stack exchange,提问作者testbg testbg
相关产品推荐
相关产品推荐

