使用Spark SQL编写查询示例

Spark SQL是Apache Spark中用于处理结构化数据的模块,它允许您使用SQL或DataFrame API来查询数据。以下是一些Spark SQL查询的示例:

基本查询示例

1. 创建临时视图并查询

from pyspark.sql import SparkSession

# 创建SparkSession
spark = SparkSession.builder \
    .appName("SparkSQLExample") \
    .getOrCreate()

# 示例数据
data = [("Alice", 25), ("Bob", 30), ("Charlie", 35)]
df = spark.createDataFrame(data, ["name", "age"])

# 创建临时视图
df.createOrReplaceTempView("people")

# 使用Spark SQL查询
result = spark.sql("SELECT name, age FROM people WHERE age > 30")
result.show()

2. 聚合查询

# 创建另一个示例DataFrame
sales_data = [("2023-01-01", "A", 100), 
              ("2023-01-01", "B", 200),
              ("2023-01-02", "A", 150),
              ("2023-01-02", "B", 250)]
sales_df = spark.createDataFrame(sales_data, ["date", "product", "amount"])
sales_df.createOrReplaceTempView("sales")

# 按产品和日期分组聚合
spark.sql("""
    SELECT product, date, SUM(amount) as total_sales
    FROM sales
    GROUP BY product, date
    ORDER BY product, date
""").show()

高级查询示例

3. 连接多个表

# 创建第二个表
products = [("A", "Product A", 10.0), 
            ("B", "Product B", 20.0)]
products_df = spark.createDataFrame(products, ["product_id", "product_name", "price"])
products_df.createOrReplaceTempView("products")

# 连接查询
spark.sql("""
    SELECT s.date, p.product_name, s.total_sales, p.price, s.total_sales * p.price as revenue
    FROM (
        SELECT product, date, SUM(amount) as total_sales
        FROM sales
        GROUP BY product, date
    ) s
    JOIN products p ON s.product = p.product_id
    ORDER BY s.date, revenue DESC
""").show()

4. 使用窗口函数

# 窗口函数示例:计算每个产品的销售排名
spark.sql("""
    SELECT 
        date, 
        product, 
        amount,
        RANK() OVER (PARTITION BY date ORDER BY amount DESC) as daily_rank
    FROM sales
""").show()

从文件读取数据并查询

5. 从JSON文件读取并查询

# 假设有一个people.json文件,内容如下:
# [{"name":"Alice","age":25},{"name":"Bob","age":30},{"name":"Charlie","age":35}]

# 从JSON文件读取
people_df = spark.read.json("people.json")
people_df.createOrReplaceTempView("people")

# 查询平均年龄
spark.sql("SELECT AVG(age) as avg_age FROM people").show()

注意事项

  1. 在使用Spark SQL前,确保已创建SparkSession
  2. 对于大型数据集,考虑使用spark.read方法直接从存储系统(如HDFS、S3)读取数据
  3. 可以使用cache()方法缓存频繁查询的表以提高性能
  4. 对于复杂查询,DataFrame API有时比纯SQL更高效

希望这些示例能帮助您开始使用Spark SQL!如果您有特定的查询需求或遇到问题,可以提供更多细节,我可以给出更具体的指导。

Logo

欢迎加入DeepSeek 技术社区。在这里,你可以找到志同道合的朋友,共同探索AI技术的奥秘。

更多推荐