Spark create or replace temp view多次不更新问题及CSV循环同步MongoDB需求
Hey there! Let’s break down your two technical challenges and walk through practical solutions for each—they’re both common in ongoing data sync workflows, so I’ve got you covered.
createOrReplaceTempView Not Updating First up, the frustrating issue where re-running createOrReplaceTempView doesn’t overwrite the existing view. Here’s why it happens and how to fix it:
Why It’s Broken
- Session Boundaries: Temp views are tied to the specific
SparkSessionthey’re created in. If you’re accidentally spawning a new session between runs, the old view sticks around in the original context while the new one lives in a separate session. - Cached Data: If you’ve cached the view (either explicitly with
cache()or implicitly via repeated queries), Spark might serve stale cached data instead of pulling from the updated view.
Solutions
Reuse a Single SparkSession
Always stick to one session across operations. Spark’s getOrCreate() method handles this perfectly—it returns an existing session if one exists, or creates a new one if not:
// Scala example val spark = SparkSession.builder() .appName("CSVDataSync") .getOrCreate() // Initial load spark.read.csv("initial_data.csv") .createOrReplaceTempView("csv_view") // Later, update the view with fresh data spark.catalog.uncacheTable("csv_view") // Clear cached data first spark.read.csv("updated_data.csv") .createOrReplaceTempView("csv_view")
Use Global Temp Views (For Cross-Session Access)
If you need the view to persist across multiple sessions (like in a Jupyter notebook with separate cells), switch to createOrReplaceGlobalTempView. Just remember to reference it with the global_temp prefix:
spark.read.csv("data.csv").createOrReplaceGlobalTempView("global_csv_view") // Query it like this spark.sql("SELECT * FROM global_temp.global_csv_view")
Now for your core business workflow: repeating CSV comparison, syncing differences to MongoDB, and updating your static DB file. Here’s a practical Python implementation (easily adaptable to Scala/Spark for larger datasets):
Core Workflow Overview
- Run on a repeating schedule
- Compare static
DB.csvwith dynamicDownloaded.csv - Sync only new/updated records to MongoDB
- Replace
DB.csvwith the latest downloaded file
Full Implementation Code
import pandas as pd from pymongo import MongoClient import shutil import time import os # Configuration - adjust these to match your setup DB_CSV_PATH = "DB.csv" DOWNLOADED_CSV_PATH = "Downloaded.csv" MONGO_URI = "mongodb://localhost:27017/" MONGO_DB_NAME = "data_sync_db" MONGO_COLLECTION_NAME = "updated_records" CHECK_INTERVAL = 3600 # Sync every hour (in seconds) PRIMARY_KEY = "id" # Use your actual unique identifier column # Initialize MongoDB connection client = MongoClient(MONGO_URI) db = client[MONGO_DB_NAME] collection = db[MONGO_COLLECTION_NAME] def sync_csv_to_mongo(): # Guard clause: Ensure both files exist before proceeding if not os.path.exists(DB_CSV_PATH) or not os.path.exists(DOWNLOADED_CSV_PATH): print("One or both CSV files are missing. Skipping this sync cycle.") return # Load CSV data df_static = pd.read_csv(DB_CSV_PATH) df_dynamic = pd.read_csv(DOWNLOADED_CSV_PATH) # Find new/updated records merged = df_static.merge( df_dynamic, on=PRIMARY_KEY, how="outer", indicator=True, suffixes=("_old", "_new") ) # Filter for records only in the downloaded file (new) or changed diff_records = merged[merged["_merge"] != "both"] # Keep only the latest version from the downloaded file diff_records = diff_records[[col for col in df_dynamic.columns]] # Sync to MongoDB if there are changes if not diff_records.empty: collection.insert_many(diff_records.to_dict("records")) print(f"Successfully synced {len(diff_records)} new/updated records to MongoDB") else: print("No changes detected between CSV files.") # Update static DB.csv with latest data shutil.move(DOWNLOADED_CSV_PATH, DB_CSV_PATH) print("Updated static DB.csv with the latest downloaded data.") # Start the sync loop if __name__ == "__main__": print("Starting CSV sync loop... Press Ctrl+C to stop.") while True: try: sync_csv_to_mongo() except Exception as e: print(f"Sync failed with error: {str(e)}") print(f"Waiting {CHECK_INTERVAL/3600} hours before next sync...\n") time.sleep(CHECK_INTERVAL)
Production Tips
- Error Handling: Add checks for file permissions, MongoDB connection timeouts, and CSV schema mismatches to make the loop more robust.
- Spark for Large Files: For big datasets, replace pandas with Spark—use
exceptAllto find differences, then write to MongoDB using the Spark-Mongo connector. - Better Scheduling: Replace the infinite loop with cron (Linux) or Task Scheduler (Windows) for more reliable, managed scheduling.
- Idempotent Writes: If you want to avoid duplicate records, use
update_onewithupsert=Trueinstead ofinsert_manyto overwrite existing entries in MongoDB.
内容的提问来源于stack exchange,提问作者Yashwanth Kambala

