databricks
// Industrial Lakehouse
x1-spark-cluster (4 Cores / 16GB / Spark 3.5.1)
PySpark Notebook
Delta Lake Catalog
Spark Streaming
Delta Lake Records
--
▲ Live Sync (PostgreSQL & InfluxDB)
Fleet Power Draw (Avg)
-- kW
480V Balanced SCADA Bus
Motor RPM Health
--
Nominal target: 1750 RPM
PySpark Micro-Batch Rate
12.4 msg/s
EMQX UNS Sparkplug B
scada_delta_lake_telemetry_analysis.py
Run Cell
# Databricks Lakehouse: Live Delta Table Feature Aggregation
from
pyspark.sql.functions
import
col, avg, round, count, desc df_silver = spark.read.table(
"scada_features_silver"
) display( df_silver.filter(col(
"power_kw"
) >
80.0
) .groupBy(
"gateway"
) .agg( round(avg(
"voltage"
),
2
).alias(
"avg_voltage_v"
), round(avg(
"power_kw"
),
2
).alias(
"avg_power_kw"
), round(avg(
"rpm"
),
0
).alias(
"avg_rpm"
), round(avg(
"temp_c"
),
2
).alias(
"avg_temp_c"
), count(
"*"
).alias(
"sample_count"
) ) )
Execution Output (Spark Action Completed in
22 ms
)
Source: PostgreSQL / Delta Lake
LOG ID
TIMESTAMP (UTC)
GATEWAY IDENTIFIER
VOLTAGE (V)
CURRENT (A)
POWER (kW)
RPM
TEMP (°C)
Loading live Spark telemetry...
Unity Catalog // Sovereign Delta Tables
TABLE NAME
TIER
STORAGE FORMAT
PARTITION KEY
RECORD COUNT
UPSTREAM CONNECTOR
scada_telemetry_bronze
BRONZE
Delta / Parquet
date(timestamp)
5,000+
EMQX MQTT UNS Broker
scada_features_silver
SILVER
Delta / Parquet
gateway
5,000+
PostgreSQL scada_telemetry_log
predictive_maintenance_gold
GOLD
Delta / Parquet
cluster_node
Hourly Aggregated
PySpark ML Anomaly Detector
cyber_threat_lake
SILVER
Delta / Parquet
severity
482+ CISA KEVs
PostgreSQL cyber_security_intel
Structured Streaming Pipeline Status
Direct micro-batch streaming from EMQX UNS Broker (10.43.96.238:1883) to Delta Lake.
Status:
RUNNING
(QueryID: q-uns-iot-telemetry-01)
Source: MQTT (Topic: spBv1.0/Enterprise/XWORKS/+)
Sink: Delta Lake (/srv/nfs/k3s-storage/delta-lakehouse/scada_bronze)
Batch Interval: 2000 ms • Trigger: ProcessingTime('2 seconds')
Input Rate: 12.4 rows/sec • Process Rate: 18.2 rows/sec • Backpressure: None