frompyspark.sqlimportSparkSession# Create a SparkSessionspark=SparkSession.builder \
.appName("MyApp") \
.getOrCreate()# Read a CSV filedf=spark.read.csv("path/to/file.csv",header=True,inferSchema=True)# Show DataFrameprint(df.show())
# Select specific columnsdf.select("column1","column2").show()# Filter rowsdf.filter(df["column"]>100).show()# Using SQL-like syntaxdf.createOrReplaceTempView("table_name")spark.sql("SELECT * FROM table_name WHERE column > 100").show()
# Add a new columndf=df.withColumn("new_column",df["existing_column"]*2)# Rename a columndf=df.withColumnRenamed("old_column","new_column")# Drop a columndf=df.drop("column_to_drop")
# Group by and aggregatedf.groupBy("column").count().show()df.groupBy("column").agg({"another_column":"sum"}).show()frompyspark.sql.functionsimportavg,maxdf.groupBy("column").agg(avg("another_column"),max("another_column")).show()
frompyspark.sql.functionsimportcol,lit,when# Use `col` to reference columnsdf.select(col("column1")*2).show()# Add a constant valuedf=df.withColumn("constant_column",lit(42))# Conditional column (like CASE WHEN)df=df.withColumn("new_column",when(df["column"]>10,"Yes").otherwise("No"))