You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

PySpark读写CosmosDB MongoDB API流数据遇微批不支持错误求助

解决Databricks PySpark中CosmosDB MongoDB API流处理的微批不支持问题

错误原因

你遇到的java.lang.UnsupportedOperationException: Data source mongodb does not support microbatch processing错误,核心原因是MongoDB Spark Connector 10.0.0版本不支持Spark Structured Streaming的微批处理模式,仅支持连续处理模式(Continuous Processing)。

修正后的代码

将流处理的触发模式改为连续模式,同时确保输出模式为append(连续模式唯一支持的输出模式),修正后的代码如下:

from pyspark.sql import SparkSession
from pyspark.sql.types import StringType, BooleanType, DateType, StructType, LongType, IntegerType 

spark = SparkSession.\
        builder.\
        appName("streamingExampleRead").\
        config('spark.jars.packages', 'org.mongodb.spark:mongo-spark-connector:10.0.0').\
        getOrCreate()

sourceConnectionString = "<primary connection string of cosmosDB API for MongoDB instance>"
sourceDb = "<your database name>"
sourceCollection =  "<your collection name>"       

# 读取CosmosDB MongoDB Change Stream
dataStreamRead = (
    spark.readStream.format("mongodb")
    .option('spark.mongodb.connection.uri', sourceConnectionString)
    .option('spark.mongodb.database', sourceDb)
    .option('spark.mongodb.collection', sourceCollection)
    .option('spark.mongodb.change.stream.publish.full.document.only','true')
    .load()
)

# 连续模式写流
query2 = (
    dataStreamRead.writeStream 
    .outputMode("append") 
    .format("console") 
    .trigger(continuous='1 second')  # 替换为连续触发模式
    .option("checkpointLocation", "/dbfs/tmp/cosmos-mongo-stream-checkpoint")  # 指定持久化checkpoint路径
    .start().awaitTermination()
)

关键修改说明

  • 触发模式切换:把trigger(processingTime='1 seconds')改为trigger(continuous='1 second'),启用连接器支持的连续处理模式
  • 输出模式限制:连续模式仅支持append输出模式,不能使用update或complete模式
  • Checkpoint优化:移除forceDeleteTempCheckpointLocation,改用Databricks分布式文件系统(DBFS)上的持久化路径,保证流处理的容错性(临时删除checkpoint在生产环境中存在数据丢失风险)

注意事项

  • 连续模式局限性:该模式不支持聚合操作、窗口函数、复杂转换逻辑,仅适用于简单的流数据转发场景
  • 版本兼容:确认MongoDB Spark Connector版本与Databricks的Spark版本匹配(10.0.0对应Spark 3.3.x及以上)
  • CosmosDB配置:确保CosmosDB MongoDB API实例已开启Change Stream功能(默认开启,但需保证目标集合支持)
  • Databricks预览:display(dataStreamRead)在连续模式下无法直接预览,建议先通过console输出验证,再切换到目标sink(如Delta Lake)

内容的提问来源于stack exchange,提问作者Krishna Kumar Sahu

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.17 12:35:18