You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何用Java编写Spark应用找出各品类销量最低的产品

Java版Spark实现各品类销量最低产品方案

1. 依赖准备

如果用Maven管理项目,需添加Spark SQL依赖(版本可根据实际环境调整):

<dependency>
    <groupId>org.apache.spark</groupId>
    <artifactId>spark-sql_2.12</artifactId>
    <version>3.3.0</version>
</dependency>

2. 核心实现步骤

步骤1:初始化SparkSession

SparkSession是Spark SQL的核心入口,先完成初始化:

SparkSession spark = SparkSession.builder()
        .appName("LowestSalesProductByCategory")
        .master("local[*]") // 本地调试用,生产环境需移除该配置
        .getOrCreate();

步骤2:读取两个数据源

假设数据为带表头的CSV格式,读取生成DataFrame:

// 读取产品信息文件
Dataset<Row> productDF = spark.read()
        .option("header", "true")
        .option("inferSchema", "true")
        .csv("path/to/product_file.csv");

// 读取订单信息文件
Dataset<Row> orderDF = spark.read()
        .option("header", "true")
        .option("inferSchema", "true")
        .csv("path/to/order_file.csv");

步骤3:关联数据并计算产品销量

通过Product ID关联两张表,计算每个产品的总销量(以下以订单数量为例,若需统计总订单金额,将count替换为sum即可):

import static org.apache.spark.sql.functions.*;

// 关联表并聚合计算每个产品的销量
Dataset<Row> productSalesDF = productDF.join(orderDF, "Product ID")
        .groupBy("Product ID", "Product Name", "Product Category")
        .agg(count("Order ID").alias("TotalSales"));

步骤4:用窗口函数筛选各品类最低销量产品

通过窗口函数按品类分区、销量升序排序,标记每个品类内的最低销量产品:

import org.apache.spark.sql.expressions.Window;
import org.apache.spark.sql.expressions.WindowSpec;

// 定义窗口规则:按品类分区,按销量升序排列
WindowSpec windowSpec = Window.partitionBy("Product Category")
        .orderBy(col("TotalSales").asc());

// 添发行号并筛选出每个品类的最低销量产品
Dataset<Row> lowestSalesByCategoryDF = productSalesDF.withColumn("row_num", row_number().over(windowSpec))
        .filter(col("row_num") == 1)
        .drop("row_num");

步骤5:输出结果

可选择打印到控制台或写入文件:

// 控制台打印结果
lowestSalesByCategoryDF.show();

// 写入输出文件(可选)
lowestSalesByCategoryDF.write()
        .option("header", "true")
        .csv("path/to/output_result.csv");

// 关闭SparkSession
spark.stop();

3. 注意事项

  • 数据源适配:若数据格式为JSON、Parquet等,替换对应读取方法(如json()、parquet())即可。
  • 空值与关联逻辑:若存在Product ID空值或不匹配的情况,可通过join(..., JoinType.INNER)确保只保留有订单记录的产品,或根据业务需求选择左关联。
  • 并列销量处理:若需保留同品类中销量相同的最低产品,将row_number()替换为rank()或dense_rank()即可。

内容的提问来源于stack exchange,提问作者heisenberg11231

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.22 04:18:14