Spark实现订单商品位置编号:正量编号与退货记录编号赋值
Hey,这个需求我之前在项目里碰到过,刚好可以给你分享Spark下三种语言的实现思路——核心逻辑就是先给正数量的记录分配连续编号,再把负数量的记录和同订单同商品的正记录关联,将对应编号加1000即可,具体实现如下:
Spark SQL 实现方案
这种方式最直观,适合熟悉SQL的同学:
-- 第一步:给正数量的记录生成初始连续编号 WITH positive_records AS ( SELECT order_id, product_id, quantity, -- 这里按product_id排序,你可以根据实际需求修改排序字段 ROW_NUMBER() OVER (PARTITION BY order_id ORDER BY product_id) AS position_number FROM orders WHERE quantity > 0 ), -- 第二步:处理负数量记录,关联对应正记录的编号并加1000 negative_records AS ( SELECT pr.position_number + 1000 AS position_number, n.order_id, n.product_id, n.quantity FROM orders n JOIN positive_records pr ON n.order_id = pr.order_id AND n.product_id = pr.product_id WHERE n.quantity < 0 ) -- 合并两类记录并按order_id和position_number排序 SELECT * FROM positive_records UNION ALL SELECT * FROM negative_records ORDER BY order_id, position_number;
PySpark (Python) 实现方案
用DataFrame API实现,代码可读性强:
from pyspark.sql import SparkSession from pyspark.sql.window import Window from pyspark.sql.functions import row_number, col # 初始化SparkSession spark = SparkSession.builder.appName("OrderPositionAssignment").getOrCreate() # 模拟原始订单数据(实际场景中可以从数据源加载) sample_data = [ ("A", "X", 5), ("A", "Y", 1), ("A", "Z", 3), ("A", "X", -1), ("A", "Z", -1) ] order_df = spark.createDataFrame(sample_data, ["order_id", "product_id", "quantity"]) # 定义窗口规则:按订单分区,按商品ID排序(可按需调整) window_spec = Window.partitionBy("order_id").orderBy("product_id") # 处理正数量记录,生成position_number positive_records = order_df.filter(col("quantity") > 0) \ .withColumn("position_number", row_number().over(window_spec)) # 处理负数量记录,关联正记录的编号并加1000 negative_records = order_df.filter(col("quantity") < 0) \ .join(positive_records, on=["order_id", "product_id"], how="inner") \ .withColumn("position_number", col("position_number") + 1000) \ .select("position_number", "order_id", "product_id", "quantity") # 合并结果并排序(unionByName保证列名匹配,避免顺序问题) final_result = positive_records.unionByName(negative_records) \ .orderBy("order_id", "position_number") # 查看结果 final_result.show()
Scala 实现方案
如果是Scala栈的项目,可以用这个版本:
import org.apache.spark.sql.SparkSession import org.apache.spark.sql.functions.{row_number, col} import org.apache.spark.sql.expressions.Window object OrderPositionGenerator { def main(args: Array[String]): Unit = { val spark = SparkSession.builder.appName("OrderPositionGenerator").getOrCreate() import spark.implicits._ // 模拟原始订单数据 val sampleData = Seq( ("A", "X", 5), ("A", "Y", 1), ("A", "Z", 3), ("A", "X", -1), ("A", "Z", -1) ).toDF("order_id", "product_id", "quantity") // 定义窗口规则 val windowSpec = Window.partitionBy("order_id").orderBy("product_id") // 处理正数量记录 val positiveRecords = sampleData.filter(col("quantity") > 0) .withColumn("position_number", row_number().over(windowSpec)) // 处理负数量记录 val negativeRecords = sampleData.filter(col("quantity") < 0) .join(positiveRecords, Seq("order_id", "product_id"), "inner") .withColumn("position_number", col("position_number") + 1000) .select("position_number", "order_id", "product_id", "quantity") // 合并结果并排序 val finalResult = positiveRecords.unionByName(negativeRecords) .orderBy("order_id", "position_number") // 输出结果 finalResult.show() } }
注意点
- 窗口函数的
ORDER BY字段可以根据实际业务需求调整,比如按quantity大小或者记录插入时间排序 - 如果存在同一商品多次正量+多次负量的情况,这个逻辑依然适用,因为每次负量都会关联到对应商品的正量编号
- 用
unionByName而不是普通union,可以避免因为列顺序不一致导致的错误
内容的提问来源于stack exchange,提问作者Federica
相关产品推荐
相关产品推荐

