【发布时间】:2021-07-16 04:22:22
【问题描述】:
scala/spark 新手在这里。我继承了一个旧代码,我已经重构并尝试使用它来从 Scylla 检索数据。代码如下:
val TEST_QUERY = s"SELECT user_id FROM test_table WHERE name = ? AND id_type = 'test_type';"
var selectData = List[Row]()
dataRdd.foreachPartition {
iter => {
// Build up a cluster that we can connect to
// Start a session with the cluster by connecting to it.
val cluster = ScyllaConnector.getCluster(clusterIpString, scyllaPreferredDc, scyllaUsername, scyllaPassword)
var batchCounter = 0
val session = cluster.connect(tableConfig.keySpace)
val preparedStatement: PreparedStatement = session.prepare(TEST_QUERY)
iter.foreach {
case (test_name: String) => {
// Get results
val testResults = session.execute(preparedStatement.bind(test_name))
if (testResults != null){
val testResult = testResults.one()
if(testResult != null){
val user_id = testResult.getString("user_id")
selectData ::= Row(user_id, test_name)
}
}
}
}
session.close()
cluster.close()
}
}
println("Head is =======> ")
println(selectData.head)
上面没有返回任何数据,并且由于空指针异常而失败,因为selectedData 列表是空的,尽管其中肯定有与 select 语句匹配的数据。我觉得我的做法不正确,但不知道需要改变什么才能解决这个问题,因此非常感谢任何帮助。
PS:我使用列表来保存结果的整个想法是,我可以使用该列表来创建数据框。如果您能在这里指出正确的方向,我将不胜感激。
【问题讨论】:
标签: scala apache-spark datastax scylla