如何用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
相关产品推荐
相关产品推荐

