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

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.

1. Fixing Spark 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 SparkSession they’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")
2. Building a Loop for CSV Sync to MongoDB

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.csv with dynamic Downloaded.csv
  • Sync only new/updated records to MongoDB
  • Replace DB.csv with 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 exceptAll to 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_one with upsert=True instead of insert_many to overwrite existing entries in MongoDB.

内容的提问来源于stack exchange,提问作者Yashwanth Kambala

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.19 08:34:00