Options
- Mark as New
- Bookmark
- Subscribe
- Mute
- Subscribe to RSS Feed
- Permalink
- Report Inappropriate Content
11-10-2025 06:06 AM
You can customize the below code, that makes use of Spark SQL Server access connector, as per your needs:
def PersistRemoteSQLTableFromDF(
df: DataFrame,
databaseName: str,
tableName: str,
mode: str = "overwrite",
schemaName: str = "",
tableLock: bool = True,
) -> None:
"""
Persist dataframe into remote SQL table using the SQL Server Spark connector.
Improvements:
- input validation and clearer errors
- tolerant handling of mode casing
- avoids double underscores from prefix
- normalizes schema handling
- ensures tableLock is passed as string expected by the connector
ARGS:
- df (DataFrame): dataframe to persist
- databaseName (str): database name (non-empty)
- tableName (str): table name (non-empty)
- mode (str): 'overwrite', 'append', 'ignore', 'error', 'errorifexists', 'default'
- prefix (str): optional prefix (no extra underscore added)
- schemaName (str): optional schema name (no trailing dot required)
- tableLock (bool): if True enables tableLock for better write performance
RAISES:
- TypeError, ValueError on invalid args
"""
# Basic type checks
if not isinstance(df, DataFrame):
raise TypeError("Argument Error - 'df' must be a pyspark.sql.DataFrame")
if not isinstance(databaseName, str) or not isinstance(tableName, str):
raise TypeError(
"Argument Error - 'databaseName' and 'tableName' must be strings"
)
# Normalize and validate string args
mode = (mode or "").lower()
allowed_modes = {
"overwrite",
"append",
"ignore",
"error",
"errorifexists",
"default",
}
if mode not in allowed_modes:
raise ValueError(
f"Argument Error - Mode '{mode}' is not valid. Accepted save modes: {sorted(allowed_modes)}"
)
databaseName = databaseName.strip()
tableName = tableName.strip()
if not databaseName or not tableName:
raise ValueError(
"Argument Error - 'databaseName' and 'tableName' must be non-empty"
)
# Normalise schema name (remove trailing dots) and prefix (avoid double underscores)
schemaName = (schemaName or "").strip().rstrip(".")
dbtable = f"{schemaName + '.' if schemaName else ''}{tableName}"
# tableLock option on connector expects "true"/"false" (string)
tableLock_opt = "true" if bool(tableLock) else "false"
try:
(
df.write.format("sqlserver")
.mode(mode)
.option("host", os.getenv("DEFAULT_SQL_SERVER_NAME"))
.option("port", "1433")
.option("database", databaseName)
.option("dbtable", dbtable)
.option("tableLock", tableLock_opt)
.option("user", os.getenv("DEFAULT_SQL_SERVER_USERNAME"))
.option("password", os.getenv("DEFAULT_SQL_SERVER_PASSWORD"))
.save()
)
except Exception as exc:
# Surface clearer context while preserving original exception
raise RuntimeError(
f"Failed to persist DataFrame to remote table '{dbtable}' in database '{databaseName}': {exc}"
) from exc