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

PySpark:基于时间与状态存在性校验物流轨迹合规性

物流配送合规校验实现方案

问题背景

现有一份带时间戳的物流全流程数据集,每个物品(tracking_id)需包含三个状态:DEPART_FROM_FACTORY、ARRIVED_AT_DISTRIBUTION、VENDOR_ACCEPTED。需校验以下两个合规条件,满足则Compliance列为True,否则为False:

  • 三个状态均存在
  • DEPART_FROM_FACTORY与ARRIVED_AT_DISTRIBUTION的created_at为同一天

已完成数据预处理及按tracking_id分组生成结构体数组列tbl的步骤,现需实现合规判断逻辑。


方案一:使用Spark内置函数(推荐)

利用Spark原生聚合与条件判断函数实现,性能优于UDF,适合大数据量场景。

实现代码

from pyspark.sql import functions as F

# 透视每个tracking_id的状态,提取各状态对应的时间
status_pivot_df = sample_df.groupBy("tracking_id").pivot("status").agg(F.first("created_at"))

# 计算合规性
compliance_df = status_pivot_df.withColumn(
    "Compliance",
    F.when(
        # 校验三个状态是否都存在
        (F.col("DEPART_FROM_FACTORY").isNotNull()) &
        (F.col("ARRIVED_AT_DISTRIBUTION").isNotNull()) &
        (F.col("VENDOR_ACCEPTED").isNotNull()) &
        # 校验两个状态的日期是否相同
        (F.to_date(F.col("DEPART_FROM_FACTORY")) == F.to_date(F.col("ARRIVED_AT_DISTRIBUTION"))),
        F.lit(True)
    ).otherwise(F.lit(False))
)

# 关联原分组表,保留tbl列
final_result = agg.join(compliance_df, on="tracking_id", how="inner")

# 查看结果
final_result.select("tracking_id", "tbl", "Compliance").show(truncate=False)

方案二:自定义UDF实现

若需遍历tbl数组实现逻辑,可通过自定义UDF完成。

实现代码

from pyspark.sql import functions as F
from pyspark.sql.types import BooleanType

def validate_compliance(tbl_array):
    # 构建状态到时间戳的映射
    status_time_map = {}
    for item in tbl_array:
        status_time_map[item["status"]] = item["created_at"]
    
    # 检查必备状态是否齐全
    required_statuses = {"DEPART_FROM_FACTORY", "ARRIVED_AT_DISTRIBUTION", "VENDOR_ACCEPTED"}
    if not required_statuses.issubset(status_time_map.keys()):
        return False
    
    # 检查两个状态的日期是否一致
    depart_date = status_time_map["DEPART_FROM_FACTORY"].date()
    arrive_date = status_time_map["ARRIVED_AT_DISTRIBUTION"].date()
    return depart_date == arrive_date

# 注册UDF
compliance_udf = F.udf(validate_compliance, BooleanType())

# 应用UDF生成合规列
final_result = agg.withColumn("Compliance", compliance_udf(F.col("tbl")))

# 查看结果
final_result.show(truncate=False)

预期输出

+-----------+---------------------------------------------------------------------------------------------------------------------------+----------+
|tracking_id|tbl                                                                                                                        |Compliance|
+-----------+---------------------------------------------------------------------------------------------------------------------------+----------+
|A001       |[{A001, 2019-10-23 11:06:46, DEPART_FROM_FACTORY}, {A001, 2019-10-23 11:08:35, ARRIVED_AT_DISTRIBUTION}, {A001, 2019-10-25 13:36:14, VENDOR_ACCEPTED}]|true      |
|A002       |[{A002, 2019-10-01 13:06:46, DEPART_FROM_FACTORY}, {A002, 2019-10-02 09:08:35, ARRIVED_AT_DISTRIBUTION}, {A002, 2019-10-03 12:36:14, VENDOR_ACCEPTED}]|false     |
|A003       |[{A003, 2019-10-07 13:06:46, DEPART_FROM_FACTORY}, {A003, 2019-10-08 09:08:35, VENDOR_ACCEPTED}]                            |false     |
+-----------+---------------------------------------------------------------------------------------------------------------------------+----------+

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.05 13:17:09