Spark Streaming写入Azure Cosmos DB遇Update模式不支持报错,求UPSERT方案
问题:Spark Streaming聚合数据写入Azure Cosmos DB时Update模式不支持,如何实现UPSERT?
我尝试用Spark Streaming将聚合数据写入Azure Cosmos DB,程序从网络控制台获取输入,做词频聚合后写入流,但运行时报错:
java.lang.IllegalArgumentException: requirement failed: com.azure.cosmos.spark.items.bua-cosmos.buabookkeeping.wordscount does not support Update mode.
完整堆栈跟踪:
java.lang.IllegalArgumentException: requirement failed: com.azure.cosmos.spark.items.bua-cosmos.buabookkeeping.wordscount does not support Update mode. at scala.Predef$.require(Predef.scala:281) at org.apache.spark.sql.execution.datasources.v2.V2Writes$.org$apache$spark$sql$execution$datasources$v2$V2Writes$$buildWriteForMicroBatch(V2Writes.scala:121) at org.apache.spark.sql.execution.datasources.v2.V2Writes$$anonfun$apply$1.applyOrElse(V2Writes.scala:90) at org.apache.spark.sql.execution.datasources.v2.V2Writes$$anonfun$apply$1.applyOrElse(V2Writes.scala:43) at org.apache.spark.sql.catalyst.trees.TreeNode.$anonfun$transformDownWithPruning$1(TreeNode.scala:584) at org.apache.spark.sql.catalyst.trees.CurrentOrigin$.withOrigin(TreeNode.scala:176) at org.apache.spark.sql.catalyst.trees.TreeNode.transformDownWithPruning(TreeNode.scala:584) at org.apache.spark.sql.catalyst.plans.logical.LogicalPlan.org$apache$spark$sql$catalyst$plans$logical$AnalysisHelper$$super$transformDownWithPruning(LogicalPlan.scala:30) at org.apache.spark.sql.catalyst.plans.logical.AnalysisHelper.transformDownWithPruning(AnalysisHelper.scala:267) at org.apache.spark.sql.catalyst.plans.logical.AnalysisHelper.transformDownWithPruning$(AnalysisHelper.scala:263) at org.apache.spark.sql.catalyst.plans.logical.LogicalPlan.transformDownWithPruning(LogicalPlan.scala:30) at org.apache.spark.sql.catalyst.plans.logical.LogicalPlan.transformDownWithPruning(LogicalPlan.scala:30) at org.apache.spark.sql.catalyst.trees.TreeNode.transformDown(TreeNode.scala:560) at org.apache.spark.sql.execution.datasources.v2.V2Writes$.apply(V2Writes.scala:43) at org.apache.spark.sql.execution.datasources.v2.V2Writes$.apply(V2Writes.scala:39) at org.apache.spark.sql.catalyst.rules.RuleExecutor.$anonfun$execute$2(RuleExecutor.scala:211) at scala.collection.LinearSeqOptimized.foldLeft(LinearSeqOptimized.scala:126) at scala.collection.LinearSeqOptimized.foldLeft$(LinearSeqOptimized.scala:122) at scala.collection.immutable.List.foldLeft(List.scala:91) at org.apache.spark.sql.catalyst.rules.RuleExecutor.$anonfun$execute$1(RuleExecutor.scala:208) at org.apache.spark.sql.catalyst.rules.RuleExecutor.$anonfun$execute$1$adapted(RuleExecutor.scala:200) at scala.collection.immutable.List.foreach(List.scala:431) at org.apache.spark.sql.catalyst.rules.RuleExecutor.execute(RuleExecutor.scala:200) at org.apache.spark.sql.catalyst.rules.RuleExecutor.$anonfun$executeAndTrack$1(RuleExecutor.scala:179) at org.apache.spark.sql.catalyst.QueryPlanningTracker$.withTracker(QueryPlanningTracker.scala:88) at org.apache.spark.sql.catalyst.rules.RuleExecutor.executeAndTrack(RuleExecutor.scala:179) at org.apache.spark.sql.execution.streaming.IncrementalExecution.$anonfun$optimizedPlan$1(IncrementalExecution.scala:81) at org.apache.spark.sql.catalyst.QueryPlanningTracker.measurePhase(QueryPlanningTracker.scala:111) at org.apache.spark.sql.execution.QueryExecution.$anonfun$executePhase$2(QueryExecution.scala:185) at org.apache.spark.sql.execution.QueryExecution$.withInternalError(QueryExecution.scala:510) at org.apache.spark.sql.execution.QueryExecution.$anonfun$executePhase$1(QueryExecution.scala:185) at org.apache.spark.sql.SparkSession.withActive(SparkSession.scala:779) at org.apache.spark.sql.execution.QueryExecution.executePhase(QueryExecution.scala:184) at org.apache.spark.sql.execution.streaming.IncrementalExecution.optimizedPlan$lzycompute(IncrementalExecution.scala:82) at org.apache.spark.sql.execution.streaming.IncrementalExecution.optimizedPlan(IncrementalExecution.scala:79) at org.apache.spark.sql.execution.QueryExecution.assertOptimized(QueryExecution.scala:136) at org.apache.spark.sql.execution.QueryExecution.executedPlan$lzycompute(QueryExecution.scala:154) at org.apache.spark.sql.execution.QueryExecution.executedPlan(QueryExecution.scala:151) at org.apache.spark.sql.execution.streaming.MicroBatchExecution.$anonfun$runBatch$15(MicroBatchExecution.scala:656) at org.apache.spark.sql.execution.streaming.ProgressReporter.reportTimeTaken(ProgressReporter.scala:375) at org.apache.spark.sql.execution.streaming.ProgressReporter.reportTimeTaken$(ProgressReporter.scala:373) at org.apache.spark.sql.execution.streaming.StreamExecution.reportTimeTaken(StreamExecution.scala:68) at org.apache.spark.sql.execution.streaming.MicroBatchExecution.runBatch(MicroBatchExecution.scala:646) at org.apache.spark.sql.execution.streaming.MicroBatchExecution.$anonfun$runActivatedStream$2(MicroBatchExecution.scala:256) at scala.runtime.java8.JFunction0$mcV$sp.apply(JFunction0$mcV$sp.java:23) at org.apache.spark.sql.execution.streaming.ProgressReporter.reportTimeTaken(ProgressReporter.scala:375) at org.apache.spark.sql.execution.streaming.ProgressReporter.reportTimeTaken$(ProgressReporter.scala:373) at org.apache.spark.sql.execution.streaming.StreamExecution.reportTimeTaken(StreamExecution.scala:68) at org.apache.spark.sql.execution.streaming.MicroBatchExecution.$anonfun$runActivatedStream$1(MicroBatchExecution.scala:219) at org.apache.spark.sql.execution.streaming.ProcessingTimeExecutor.execute(TriggerExecutor.scala:67) at org.apache.spark.sql.execution.streaming.MicroBatchExecution.runActivatedStream(MicroBatchExecution.scala:213) at org.apache.spark.sql.execution.streaming.StreamExecution.$anonfun$runStream$1(StreamExecution.scala:307) at scala.runtime.java8.JFunction0$mcV$sp.apply(JFunction0$mcV$sp.java:23) at org.apache.spark.sql.SparkSession.withActive(SparkSession.scala:779) at org.apache.spark.sql.execution.streaming.StreamExecution.org$apache$spark$sql$execution$streaming$StreamExecution$$runStream(StreamExecution.scala:285) at org.apache.spark.sql.execution.streaming.StreamExecution$$anon$1.run(StreamExecution.scala:208)
原代码:
import org.apache.spark.sql.SparkSession object SparkStreamCosmos extends App{ val cfgMap = Map("spark.cosmos.accountEndpoint" -> "https://xxx-cosmos.documents.azure.com:443/", "spark.cosmos.accountKey" -> "xxxxxxx==", "spark.cosmos.database" -> "buabookkeeping", "spark.cosmos.container" -> "wordscount", "spark.cosmos.write.strategy" -> "ItemOverwrite" ) // start the local host n/w server by using the command : nc -lk 9999 val spark = SparkSession .builder .appName("StructuredNetworkWordCount") .master("local[2]") .getOrCreate() spark.sparkContext.setLogLevel("ERROR") import spark.implicits._ // Create DataFrame representing the stream of input lines from connection to localhost: 9999 val lines = spark.readStream .format("socket") .option("host", "localhost") .option("port", 9999) .load() // Split the lines into words val words = lines.as[String].flatMap(_.split(" ")) // Generate running word count var wordCounts = words.groupBy("value").count() wordCounts = wordCounts.withColumnRenamed("value", "id") // write to cosmos DB wordCounts.writeStream. format("cosmos.oltp"). outputMode("update") .options(cfgMap) .option("checkpointLocation", "/Users/k0d03gd/project/code/krushna/spark-hello/checkpointLocation") .start() .awaitTermination(100000) }
解决方案
Azure Cosmos DB的Spark流写入连接器目前不支持update输出模式,仅支持append或complete模式。但complete模式会全量写入所有聚合结果,而我们需要仅对更新的行执行UPSERT,因此可以通过foreachBatch自定义批量写入逻辑,利用Cosmos DB批量写入的ItemOverwrite策略实现UPSERT。
修改后的核心代码如下:
// 替换原有的writeStream部分 wordCounts.writeStream .foreachBatch { (batchDF: org.apache.spark.sql.DataFrame, batchId: Long) => // 对每个批次的DataFrame执行批量写入,用ItemOverwrite实现UPSERT batchDF.write .format("cosmos.oltp") .options(cfgMap) .mode("append") // 批量写入用append模式,结合ItemOverwrite策略实现UPSERT .save() } .outputMode("update") // 流处理层用update模式,仅输出变化的行 .option("checkpointLocation", "/Users/k0d03gd/project/code/krushna/spark-hello/checkpointLocation") .start() .awaitTermination(100000)
说明:
foreachBatch允许我们在每个微批次中处理结果DataFrame,这里用批量写入API操作Cosmos DB,批量API支持ItemOverwrite策略(存在则更新,不存在则插入)。- 流处理层保持
update输出模式,确保每个批次仅输出发生变化的词频记录,减少写入的数据量。 - 批量写入时使用
append模式,配合spark.cosmos.write.strategy设置的ItemOverwrite,即可实现目标UPSERT效果。
内容的提问来源于stack exchange,提问作者Krushna Dash
相关产品推荐
相关产品推荐

