You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Scala如何从while/for循环返回值?Spark读取Oracle代码报错求解

问题原因分析
  • 报错ret is already defined as value ret的根源:Scala的for推导式(for/yield语法)中写的ret = ret.concat(...)会被识别为创建新的val变量ret,和你外部定义的var ret重名,并非修改外部变量。
  • 打印空值/提前返回的问题:原代码中return放在了for循环内部,第一次遍历到结果集第一行第一个字段就会直接返回,不会遍历完整的结果集,若结果集为空则直接返回初始的空字符串,也不会走到后续的打印逻辑。
  • Cannot resolve symbol concat的提示是IDE识别到for推导式中错误的变量重定义,认为你调用concat的ret是未初始化的新变量。
修复后的原生JDBC读取代码

如果你坚持使用自定义JDBC读取的方式,修正后的代码如下:

import java.sql.{Connection, Statement, ResultSet}
import oracle.jdbc.pool.OracleDataSource

def read_data(group_id: Int, num_node: Int): String =  {
  val table_name = "table"
  val col_name = "col"
  val query =
    s"""select f1,f2,f3,f4,f5,f6,f7,f8
       |from $table_name 
       |where MOD(TO_NUMBER(substr($col_name, -LEAST(2, LENGTH($col_name)))), $num_node) = $group_id""".stripMargin

  val oracleUser = "ORCL"
  val oraclePassword = "*******"
  val oracleURL = "jdbc:oracle:thin:@//x.x.x.x:1521/ORCLDB"
  
  var con: Connection = null
  var statement: Statement = null
  var rs: ResultSet = null
  val ret = new StringBuilder() // 用StringBuilder比字符串拼接效率高很多

  try {
    val ods = new OracleDataSource()
    ods.setUser(oracleUser)
    ods.setURL(oracleURL)
    ods.setPassword(oraclePassword)
    con = ods.getConnection()
    statement = con.createStatement()
    statement.setFetchSize(1000)
    rs = statement.executeQuery(query)

    while (rs.next()) {
      // 普通for循环直接修改StringBuilder,不用for推导式
      for (i <- 1 to 8) { // 你选了8个字段,应该是1到8不是until 8
        ret.append(rs.getString(i)).append(" ")
      }
      ret.append("\n") // 不同行之间加换行区分
    }
    val result = ret.toString()
    println("ret:", result)
    result
  } finally {
    // 必须关闭资源避免泄漏
    if (rs != null) rs.close()
    if (statement != null) statement.close()
    if (con != null) con.close()
  }
}

// 调用部分注意不要在map里直接println,executor的打印不会输出到driver
val rdd = sc.parallelize(group_list, num_node)
  .map(read_data(_, num_node))
rdd.collect().foreach(println) // 收集到driver再打印
println("rdd计数:", rdd.count())
更推荐的Spark原生JDBC读取方案

你要获取Spark DataFrame完全不需要手写JDBC逻辑,Spark原生JDBC支持自动并行分区读取,性能和稳定性都优于自定义实现,示例代码如下:

val oracleUser = "ORCL"
val oraclePassword = "*******"
val oracleURL = "jdbc:oracle:thin:@//x.x.x.x:1521/ORCLDB"
val tableName = "table"
val numPartitions = num_node // 你要的分区数

val df = spark.read.format("jdbc")
  .option("url", oracleURL)
  .option("dbtable", tableName)
  .option("user", oracleUser)
  .option("password", oraclePassword)
  .option("numPartitions", numPartitions)
  // 自定义分区逻辑,和你原来的MOD分片逻辑一致
  .option("partitionColumn", "col")
  .option("lowerBound", 0)
  .option("upperBound", numPartitions)
  .option("partitionColumn", "MOD(TO_NUMBER(substr(col, -LEAST(2, LENGTH(col)))), " + numPartitions + ")")
  .load()

// 直接选择需要的字段
val resultDf = df.select("f1","f2","f3","f4","f5","f6","f7","f8")
resultDf.show()

内容的提问来源于stack exchange,提问作者M_Gh

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.10.02 09:27:03