Syntasa Notebook Utilities
Spark Utilities
spark — dataset → DataFrame + write helpers (requires Spark)
| Method | Purpose |
|---|---|
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.
| Parameter | Type | Default | Purpose |
|---|---|---|---|
df | DataFrame | — | Source DataFrame to write |
datasetName | String | — | Dataset name in database.tablename format |
numPartitions | Int | None (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 |
partitionedDateColumn | String | None (Py) / "" (Scala) | Override the dataset's configured partition column. If set, the dataset is updated to partition by this column before writing |
isOverwrite | Boolean | True | True → INSERT OVERWRITE TABLE (replaces partition data); False → INSERT INTO TABLE (appends) |
overrideProcessMode | Boolean | True | When True, recreates / re-initializes the table even if it exists. Set False to preserve an existing table definition |
fileFormat | FileFormat / String | PARQUET | Used 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=Trueoverwrites at the partition level for partitioned tables, not the entire table. For non-partitioned tables it overwrites the whole table.