Following the answer from Douglas I came up with this function to use inside Databricks Notebooks and get some information about the cached RDDs:
I hope it helps. I was looking information about 2h on this topic.
def get_databricks_rdd_info():
import requests, json
# Get Spark Context
sc = spark.sparkContext
# Get App Id (Notebook is attached to it)
app_id = sc._jsc.sc().applicationId()
# Where is my driver
driver_ip = spark.conf.get('spark.driver.host')
port = spark.conf.get("spark.ui.port")
# Compose the query to the Spark UI API
url = f"http://{driver_ip}:{port}/api/v1/applications/{app_id}/storage/rdd"
# Make request
r = requests.get(url, timeout=3.0)
if r.status_code == 200:
# Compose results
df = spark.createDataFrame([json.dumps(r) for r in r.json()], T.StringType())
json_schema = spark.read.json(df.rdd.map(lambda row: row.value)).schema
df = df.withColumn('value', F.from_json(F.col('value'), json_schema))
df = df.selectExpr('value.*')
# Generate summary
df_summary = (df
.withColumn('Name', F.element_at(F.split(F.col('name'), ' '), -1))
.withColumn('Cached', F.round(F.lit(100) * F.col('numCachedPartitions')/F.col('numPartitions'), 2))
.withColumn('Memory GB', F.round(F.col('memoryUsed')*1e-9, 2))
.withColumn('Disk GB', F.round(F.col('diskUsed')*1e-9, 2))
.withColumnRenamed('numPartitions', '# Partitions')
.select([
'Name',
'id',
'Cached',
'Memory GB',
'Disk GB',
'# Partitions',
]))
else:
print('Some error happened, code:', r.status_code)
df = None
df_summary = None
return df, df_summary
You can use it inside Databricks Notebooks as:
df_rdd_info, df_summary = get_databricks_rdd_info()
display(df_summary)