PySpark加载MongoDB表时遭遇不完整URI错误求助
Hey there! Let's figure out why your HDFS-to-MongoDB data load is still throwing errors even after using the triple-slash HDFS URI. I’ll walk through the most likely issues and fixes based on your code snippet and common pitfalls:
1. Fix the Obvious Parsing Bug in Your Code
Looking at your parseInput function, there’s a clear mistake here:
return Row(userId=int(fields[0]), age=int(fields[1]), gender=fields[2], occupation=fields[3], zip=fields)
You’re assigning the entire fields list to the zip field instead of the 5th element (index 4). This will create malformed data that’s almost certainly causing downstream errors when writing to MongoDB. Correct it to:
return Row(userId=int(fields[0]), age=int(fields[1]), gender=fields[2], occupation=fields[3], zip=fields[4])
2. Double-Check Your HDFS URI Format
Even with triple slashes, make sure your URI follows the right pattern:
- For default HDFS setup (namenode running on default port):
hdfs:///path/to/your/file.txt(three slashes, no hostname/port needed) - If you’re targeting a specific namenode:
hdfs://namenode-host:port/path/to/your/file.txt(two slashes afterhdfs:, then hostname/port, then path)
Verify the file actually exists at that path withhdfs dfs -ls hdfs:///path/to/your/file.txtand that the Spark user has read permissions.
3. Validate MongoDB Spark Connector Configuration
You need to properly configure your SparkSession to talk to MongoDB, and ensure you have the right dependencies:
Example SparkSession Setup
from pyspark.sql import SparkSession spark = SparkSession.builder \ .appName("HDFS-to-MongoDB") \ # Replace with your MongoDB connection details .config("spark.mongodb.output.uri", "mongodb://localhost:27017/your_database.your_collection") \ .getOrCreate()
Add Connector Dependency
When starting your Spark job, include the MongoDB Spark connector package (match versions to your Spark and MongoDB):
spark-submit --packages org.mongodb.spark:mongo-spark-connector_2.12:3.0.1 your_script.py
4. Verify Data Writing Logic
Make sure you’re correctly converting your RDD to a DataFrame and writing to MongoDB:
# Read HDFS text file lines = spark.sparkContext.textFile("hdfs:///path/to/your/file.txt") # Parse and convert to DataFrame users_df = lines.map(parseInput).toDF() # Write to MongoDB (use mode("overwrite") if you want to replace existing data) users_df.write.format("mongo").mode("append").save()
5. Check Permissions & Connectivity
- HDFS Permissions: Ensure the user running Spark has read access to the HDFS file (run
hdfs dfs -chmodif needed). - MongoDB Permissions: Confirm the MongoDB user has
insertpermissions on the target collection. - Network Connectivity: If MongoDB is on a different machine, make sure the Spark cluster can reach the MongoDB port (default 27017).
Start with fixing the parsing bug first—it’s the most straightforward issue, and resolving it might eliminate the error you’re seeing. If you still run into problems, share the exact error message you’re getting, and we can dive deeper!
内容的提问来源于stack exchange,提问作者Rodrigo Ferreira

