【发布时间】:2020-06-19 19:15:43
【问题描述】:
我的 dStream.foreachRDD 方法中有一个处理块,该处理包括使用 spark sql 持久化到 mysql。 发布后,我将最新处理的偏移量保存在另一个模式/表中。我想让整个块交易(scala)。如何做到这一点? 以下是代码的相关摘录:
foreachRDD(rdd => {
...........
...................................
df.write.mode("append") .jdbc(url + rawstore_schema +"?rewriteBatchedStatements=true",tablesToFetch(index),connectionProperties)
....................
metricsStatement.executeUpdate("Insert into metrics.txn_offsets (topic,part,off,date_updated) values (...........................
}
由于写入操作(已处理数据和偏移数据)都是在两个不同的数据库/连接上完成的,如何使它们具有事务性?
谢谢
【问题讨论】:
标签: mysql scala transactions apache-spark-sql