【发布时间】:2016-02-05 12:00:12
【问题描述】:
我有一个火花流作业,它从 Kafka 读取数据,并在再次写入 Postrges 之前与 Postgres 中的现有表进行一些比较。这就是它的样子:
val message = KafkaUtils.createStream(...).map(_._2)
message.foreachRDD( rdd => {
if (!rdd.isEmpty){
val kafkaDF = sqlContext.read.json(rdd)
println("First")
kafkaDF.foreachPartition(
i =>{
val jdbcDF = sqlContext.read.format("jdbc").options(
Map("url" -> "jdbc:postgresql://...",
"dbtable" -> "table", "user" -> "user", "password" -> "pwd" )).load()
createConnection()
i.foreach(
row =>{
println("Second")
connection.sendToTable()
}
)
closeConnection()
}
)
这段代码在 val jbdcDF = ... 行给了我 NullPointerException
我做错了什么?此外,我的日志"First" 有效,但"Second" 没有出现在日志中的任何位置。我用kafkaDF.collect().foreach(...) 尝试了整个代码,它运行良好,但性能很差。我希望用foreachPartition 替换它。
谢谢
【问题讨论】:
标签: scala apache-spark spark-streaming