HDInsight Spark能否通过MongoDB API连接Cosmos DB?
Absolutely, connecting HDInsight Spark to Cosmos DB using the MongoDB API is totally feasible—the org.bson.BsonInvalidOperationException you’re seeing is almost always tied to misconfiguration or version mismatches. Let’s break down the steps to get this working:
1. Verify Cosmos DB MongoDB API Compatibility
First, make sure your Cosmos DB account is configured with a MongoDB-compatible version that aligns with your Spark MongoDB connector. Cosmos DB supports MongoDB 3.6, 4.0, 4.2, and 5.0—check this in your Cosmos DB account’s Settings > Features pane. Pick a connector version that matches this compatibility:
- For MongoDB 4.0+/Spark 3.x: Use connector version
10.x.x - For MongoDB 3.6/Spark 2.4.x: Use connector version
2.4.x
2. Correctly Add Connector Dependencies
When submitting your Spark job (or configuring your cluster), include the right Maven dependency for the connector. For example:
- Spark 3.x + Scala 2.12:
--packages org.mongodb.spark:mongo-spark-connector_2.12:10.1.1 - Spark 2.4.x + Scala 2.11:
--packages org.mongodb.spark:mongo-spark-connector_2.11:2.4.2
3. Critical Configuration Parameters
Your Spark configuration needs specific Cosmos DB MongoDB API URIs and settings—this is where most people slip up. Here’s the correct format for input/output URIs:
// Scala example val spark = SparkSession.builder() .appName("CosmosDBSparkMongo") .config("spark.mongodb.input.uri", "mongodb://<ACCOUNT_NAME>:<PRIMARY_KEY>@<ACCOUNT_NAME>.mongo.cosmos.azure.com:10255/?ssl=true&replicaSet=globaldb&retrywrites=false&maxIdleTimeMS=120000&appName=@<ACCOUNT_NAME>@") .config("spark.mongodb.output.uri", "mongodb://<ACCOUNT_NAME>:<PRIMARY_KEY>@<ACCOUNT_NAME>.mongo.cosmos.azure.com:10255/?ssl=true&replicaSet=globaldb&retrywrites=false&maxIdleTimeMS=120000&appName=@<ACCOUNT_NAME>@") .getOrCreate()
Key settings to note:
ssl=true: Mandatory for Cosmos DB’s secure connectionreplicaSet=globaldb: Fixed value for Cosmos DB MongoDB APIretrywrites=false: Cosmos DB doesn’t support MongoDB’s retryable writes—omitting this often triggers BSON-related errors- Replace
<ACCOUNT_NAME>and<PRIMARY_KEY>with your actual Cosmos DB credentials
4. Troubleshoot the BsonInvalidOperationException
If you still hit this error, check these common issues:
- Invalid BSON types: Ensure your DataFrame doesn’t contain data types that aren’t supported by MongoDB (e.g., complex nested types that don’t map cleanly to BSON). Try writing a simple test DataFrame first to rule this out.
- Connector version mismatch: Double-check that your connector version matches both your Spark version and Cosmos DB’s MongoDB compatibility version. Mismatches here can cause low-level BSON parsing errors.
- URI typos: A missing parameter (like
ssl=true) or incorrect credentials can lead to unexpected BSON errors during connection setup.
Example: Read/Write Data
Once configured, you can read and write data like this:
// Read data from Cosmos DB collection val df = spark.read.format("mongo").option("collection", "myCollection").load() // Write data to Cosmos DB collection df.write.format("mongo").option("collection", "myCollection").mode("append").save()
内容的提问来源于stack exchange,提问作者MaSe

