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
相关产品推荐
相关产品推荐

