Skip to main content
Syntasa Notebook Utilities

Spark Utilities

spark — dataset → DataFrame + write helpers (requires Spark)​

MethodPurpose
createDataFrame(datasetName, from_date=None, to_date=None)Read dataset into a DataFrame
isTableExists(dataset)True if Hive table exists
createTable(df, name, partitionedDateColumn="")Create Hive table from DataFrame
writeToEventStore(df, datasetName, ...)Write a DataFrame to an event store. Auto-creates the dataset if it doesn't exist.
writeDatasetToEventStore(df, datasetName)Convenience wrapper around writeToEventStore for a registered dataset (uses dataset defaults)
writeFileToEventStore(localPath, eventstorePath)Push a local file into the event store

Python​

df = synutils.spark.createDataFrame("user_events") 
df.show(5)
synutils.spark.writeDatasetToEventStore(df, "user_events")

Scala​

val df = synutils.spark.createDataFrame("user_events") 
df.show(5)
synutils.spark.writeDatasetToEventStore(df, "user_events")

Details​

writeToEventStore — full signature​

Writes a DataFrame to a Hive-managed event store table. If the dataset doesn't already exist in the platform's metadata, it is created on the fly using the supplied fileFormat. Partitioning, compression, and process-mode handling are derived from the dataset definition; this method only orchestrates the Spark-side write.

ParameterTypeDefaultPurpose
dfDataFrame—Source DataFrame to write
datasetNameString—Dataset name in database.tablename format
numPartitionsIntNone (Py) / 0 (Scala)If > 0 and the dataset is partitioned, adds DISTRIBUTE BY <partition_cols>, floor(rand()*numPartitions) to control output file count per partition
partitionedDateColumnStringNone (Py) / "" (Scala)Override the dataset's configured partition column. If set, the dataset is updated to partition by this column before writing
isOverwriteBooleanTrueTrue → INSERT OVERWRITE TABLE (replaces partition data); False → INSERT INTO TABLE (appends)
overrideProcessModeBooleanTrueWhen True, recreates / re-initializes the table even if it exists. Set False to preserve an existing table definition
fileFormatFileFormat / StringPARQUETUsed only when the dataset must be created (404 from the dataset API). Ignored if the dataset already exists. Supported: PARQUET, ORC, AVRO, DELTA, TEXTFILE

Python​

from synutils.file_format
import FileFormat
# Minimal write — dataset must already exist
synutils.spark.writeToEventStore(df, "analytics.user_events")
# Append (don't overwrite existing partitions)
synutils.spark.writeToEventStore(df, "analytics.user_events", isOverwrite=False)
# Control output file count per partition (e.g. 8 files per partition)
synutils.spark.writeToEventStore(df, "analytics.user_events", numPartitions=8)
# Override the partition column for this write only
synutils.spark.writeToEventStore(df, "analytics.user_events", partitionedDateColumn="event_date")
# Auto-create as Avro if dataset doesn't exist yet
synutils.spark.writeToEventStore(df, "analytics.new_avro_dataset", fileFormat=FileFormat.AVRO)
# Preserve existing table definition (don't recreate)
synutils.spark.writeToEventStore(df, "analytics.user_events", overrideProcessMode=False)

Scala​

Details​

import com.syntasa.synutils.FileFormat 

// Minimal write
synutils.spark.writeToEventStore(df, "analytics.user_events")

// Append
synutils.spark.writeToEventStore(df, "analytics.user_events", isOverwrite = false)

// Control output file count per partition
synutils.spark.writeToEventStore(df, "analytics.user_events", numPartitions = 8)

// Override partition column for this write
synutils.spark.writeToEventStore(df, "analytics.user_events", partitionedDateColumn = "event_date")

// Auto-create as Avro
synutils.spark.writeToEventStore(df, "analytics.new_avro_dataset", fileFormat = FileFormat.AVRO.getValue())

// Preserve existing table definition
synutils.spark.writeToEventStore(df, "analytics.user_events", overrideProcessMode = false)

Heads-up: isOverwrite=True overwrites at the partition level for partitioned tables, not the entire table. For non-partitioned tables it overwrites the whole table.