Spark JDBC读取使用变量传值出现type mismatch类型不匹配错误如何解决?
问题原因
两个核心错误共同导致运行失败:
- customerIds处理错误:
collect()返回的结果是Array[Row]类型,直接对Row数组做mkString会生成带Row标识的异常字符串(如[123],[456]),而非你需要的纯ID值。 - SQL拼接类型不匹配:CustomerID是整数类型,你拼接时额外加了单引号包裹,生成的SQL为
IN ('123','456'),等于把数值转为字符串和数据库整数字段匹配,这就是你硬编码带单引号报错、不带单引号可正常执行的原因。
你看到的found : String required: (?, ?)语法报错,是Scala Map构造时的类型匹配问题,改用逐个调用option方法的写法可以规避。
修正后代码
var df = spark.read.format("csv").option("header", "true").option("sep", ",").load("abfss://test@a.dfs.core.windows.net/data/raw") import org.apache.spark.sql df = df.withColumn("CustomerID", $"CustomerID".cast(sql.types.IntegerType)).withColumn("LAST_UPDATED", $"LAST_UPDATED".cast(sql.types.DateType)) // 从Row中提取Int类型的CustomerID值 val customerIds = df.select("CustomerID").collect().map(_.getAs[Int]("CustomerID")) // 整数ID不需要加单引号包裹 val inCondition = customerIds.mkString(",") val jdbcData = spark.read.format("jdbc") .option("url", jdbcUrl) .option("dbtable", s"(SELECT * FROM [dbo].[SinkCustomers] WHERE CustomerID IN ($inCondition)) AS C") .load()
优化建议
如果CustomerID的数量超过1000,不建议使用IN条件查询:一是SQL长度容易超过数据库限制,二是查询性能差。可以将读取CSV得到的df做广播后,和JDBC读取的SinkCustomers全量表做关联,Spark会自动执行谓词下推优化,性能更稳定。
内容的提问来源于stack exchange,提问作者codebot
相关产品推荐
相关产品推荐

