4.5 KiB
Executable File
4.5 KiB
Executable File
In [ ]:
from pyspark.sql import SparkSessionspark = ( SparkSession.builder .appName("lab3-partitioning-schema-evolution") .master("spark://spark-master:7077") .getOrCreate())In [ ]:
from pathlib import Pathsql_path = Path("/opt/src/spark/partitioned_table_demo.sql")sql_text = sql_path.read_text(encoding="utf-8")for statement in sql_text.split(";"): stmt = statement.strip() if stmt: spark.sql(stmt)In [ ]:
spark.sql("""SELECT event_date, COUNT(*) AS cnt, SUM(amount) AS total_amountFROM lakehouse.default.partition_demoGROUP BY event_dateORDER BY event_date""").show()In [ ]:
spark.sql("""SELECT *FROM lakehouse.default.partition_demoWHERE event_date = DATE '2024-01-01'ORDER BY user_id""").show()In [ ]:
sql_path = Path("/opt/src/spark/schema_evolution_demo.sql")sql_text = sql_path.read_text(encoding="utf-8")for statement in sql_text.split(";"): stmt = statement.strip() if stmt: spark.sql(stmt)In [ ]:
spark.sql("DESCRIBE TABLE lakehouse.default.schema_evolution_demo").show(truncate=False)In [ ]:
spark.sql("""SELECT *FROM lakehouse.default.schema_evolution_demoORDER BY id""").show(truncate=False)