Spark中mapPartitionsWithIndex执行JDBC查询返回结果异常问题
问题:mapPartitionsWithIndex结合JDBC并行查询返回空结果的排查与修复
我之前用mapPartitionsWithIndex写过一个简单示例,运行完全正常:
val rdd1 = sc.parallelize(List(1,2,3,4,5,6,7,8,9,10), 3) def myfunc(index: Int, iter: Iterator[Int]) : Iterator[String] = { iter.map(x => index + "," + (x, x, x+100)) } rdd1.mapPartitionsWithIndex(myfunc).collect()
现在我想把这个方法用到JDBC并行查询的场景里——因为有些数据源的处理逻辑太复杂,用DataFrame或者常规RDD操作搞不定。我写了模拟调用的代码,但执行后总是返回空数据,没法正确拿到数据库的Any类型结果:
import java.sql.DriverManager import java.util.Properties val rdd1 = sc.parallelize(List("G%", "C%", "I%", "B%", "X%", "F%", "J%"), 3) def myfunc(index: Int, iter: Iterator[String]) : Iterator[Any] = { val jdbcHostname = "mysql-rfam-public.ebi.ac.uk" val jdbcPort = 4497 val jdbcDatabase = "Rfam" val jdbcUrl = s"jdbc:mysql://${jdbcHostname}:${jdbcPort}/${jdbcDatabase}" val jdbcUsername = "rfamro" val jdbcPassword = "" val connectionProperties = new Properties() connectionProperties.put("user", s"${jdbcUsername}") connectionProperties.put("password", s"${jdbcPassword}") val connection = DriverManager.getConnection(jdbcUrl, jdbcUsername, jdbcPassword) iter.map { x => val val1 = x; val statement = connection.createStatement() val resultSet = statement.executeQuery(s"""(select DISTINCT type from family where type like '${val1}' ) """) while ( resultSet.next() ) { val hInType = resultSet.getString("type") } } } rdd1.mapPartitionsWithIndex(myfunc).collect()
后来我试着用ListBuffer来收集结果,结果还是返回null。我怀疑是不是编译器只执行了最后一行语句?另外我也不确定是不是每个分区真的只建立了一次数据库连接,想问问这个方案到底可行不可行,以及该怎么修改代码?我尝试的代码片段是这样的:
... var fruits = new ListBuffer[String]() iter.map { x => val val1 = x; println (x) val statement = connection.createStatement() val resultSet = statement.executeQuery(s"""(select DISTINCT type from family where type like '${val1}' ) """) while ( resultSet.next() ) { val hInType = resultSet.getString("type") fruits += hInType } } return fruits.toList.toIterator
内容的提问来源于stack exchange,提问作者Ged
相关产品推荐
相关产品推荐

