【问题标题】:split row at specific time在特定时间拆分行
【发布时间】:2021-06-11 11:17:41
【问题描述】:

所以我有一张来自工厂时间序列传感器数据的表格。其中一个传感器负责处理传送带上的原始产品在加工到炼油厂之前的移动(电压/重量秤)。每当 24 小时内的 delta(皮带电压低于或高于正常值/皮带上的重量(每秒得出)低于或高于目标值时(目标 ÷ 86,400 秒) 〜四舍五入到最接近的吨,没有小数)我们将其捕获为新的事件触发器并在我们的仓库数据库中行并移动到数据湖中 我们需要通过工作班次(白班/重班)来找到跨越班次时间的时间段的效率

考虑到 2400 吨的目标,在早上 5:00 到下午 5:00 之间的正常日班和夜班反之亦然,我们需要以下数据框:

开始数据帧

row # event_start event_end operation_status tons_actual tons_target comment
1 2021-02-01 7:00 AM 2021-02-01 9:00 AM normal_run 197 200
2 2021-02-01 9:00 AM 2021-02-01 7:00 PM curtailed 700 1004 shift split here
3 2021-02-01 7:00 PM 2021-02-01 11:00 PM down_for_maintenance 0 301
4 2021-02-01 11:00 PM 2021-02-02 3:00 AM curtailed 320 402
5 2021-02-02 3:00 AM 2021-02-02 8:00 AM over_producing 600 502 shift split here
6 2021-02-02 8:00 AM 2021-02-02 11:00 AM normal_run 280 301
7 2021-02-02 11:00 AM 2021-02-04 4:00 PM broken_belt_unscheduled_loss 0 5323 multiple shift splits here

像这样在换班时间拆分行:

目标数据帧

row # event_start event_end operation_status tons_actual tons_target --------
1 2021-02-01 7:00 AM 2021-02-01 9:00 AM normal_run 197 200
2.1 2021-02-01 9:00 AM 2021-02-01 5:00 PM curtailed 560 804 grave shift split
2.2 2021-02-01 5:00 PM 2021-02-01 7:00 PM curtailed 140 201 grave shift split
3 2021-02-01 7:00 PM 2021-02-01 11:00 PM down_for_maintenance 0 302
4 2021-02-01 11:00 PM 2021-02-02 3:00 AM curtailed 320 402
5.1 2021-02-02 3:00 AM 2021-02-02 5:00 AM over_producing 240 200 day shift split
5.2 2021-02-02 5:00 AM 2021-02-02 8:00 AM over_producing 360 302 day shift split
6 2021-02-02 8:00 AM 2021-02-02 11:00 AM normal_run 280 301
7.1 2021-02-02 11:00 AM 2021-02-02 5:00 PM broken_belt_unscheduled_loss 0 602 shift split
7.2 2021-02-02 5:00 PM 2021-02-03 5:00 AM broken_belt_unscheduled_loss 0 1205 shift split
7.3 2021-02-03 5:00 AM 2021-02-03 5:00 PM broken_belt_unscheduled_loss 0 1205 shift split
7.4 2021-02-03 5:00 PM 2021-02-04 5:00 AM broken_belt_unscheduled_loss 0 1205 shift split
7.5 2021-02-03 5:00 AM 2021-02-04 4:00 PM broken_belt_unscheduled_loss 0 1105 shift split

所以最终结果可以是每班df.groupby(sum : tons)

首先,我知道它需要某种数组在 F.explode() 函数中创建 UDF

【问题讨论】:

    标签: python-3.x apache-spark pyspark apache-spark-sql


    【解决方案1】:

    您可以使用flatMap 将单行转换为多行。

    第 1 步:解析日期列(如有必要,取决于数据源):

    df = spark.read....
    
    dateformat = "yyyy-MM-dd h:mm a"
    df = df.withColumn("event_start", F.to_timestamp(F.col("event_start"), dateformat)) \
           .withColumn("event_end", F.to_timestamp(F.col("event_end"), dateformat))
    

    第二步:定义实际逻辑。函数split_shifts 接受row 并根据需要生成尽可能多的新行:

    def split_shifts(r):
        import datetime as dt
        def get_next_end(start):
            if start.hour < 5:
                end = dt.datetime(year=start.year, month=start.month, day=start.day, hour=5)
            elif start.hour < 17:
                end = dt.datetime(year=start.year, month=start.month, day=start.day, hour=17)
            else:
                next_day = start + dt.timedelta(days=1)
                end = dt.datetime(year=next_day.year, month=next_day.month, day=next_day.day, hour=5)
            return end
        def calc_tons(start, end, current_start, current_end, tons):
            return (current_end-current_start)/(end-start)*tons
    
        row = r['row']
        event_start = r['event_start']
        event_end = r['event_end']
        operation_status = r['operation_status']
        tons_actual = r['tons_actual']
        tons_target = r['tons_target']
    
        current_event_start = event_start
        expected_event_end = get_next_end(current_event_start)
        while( expected_event_end < event_end):
            yield Row(row=row, event_start=current_event_start, event_end=expected_event_end, operation_status=operation_status, tons_actual=calc_tons(event_start, event_end, current_event_start, expected_event_end, tons_actual), tons_target=calc_tons(event_start, event_end, current_event_start, expected_event_end, tons_target))
            current_event_start = expected_event_end
            expected_event_end = get_next_end(current_event_start)
        yield Row(row=row, event_start=current_event_start, event_end=event_end, operation_status=operation_status, tons_actual=calc_tons(event_start, event_end, current_event_start, event_end, tons_actual), tons_target=calc_tons(event_start, event_end, current_event_start, event_end, tons_target))
    

    第 3 步:使用 flatMap 应用 split_shifts 函数:

    df = df.rdd.flatMap(lambda r: split_shifts(r)).toDF()
    

    第 4 步:必要时将日期列格式化为字符串:

    df = df.withColumn("event_start", F.date_format(F.col("event_start"), dateformat)) \
           .withColumn("event_end", F.date_format(F.col("event_end"), dateformat))
    

    输出:

    +---+-------------------+-------------------+--------------------+-----------+------------------+
    |row|        event_start|          event_end|    operation_status|tons_actual|       tons_target|
    +---+-------------------+-------------------+--------------------+-----------+------------------+
    |  1| 2021-02-01 7:00 AM| 2021-02-01 9:00 AM|          normal_run|      197.0|             200.0|
    |  2| 2021-02-01 9:00 AM| 2021-02-01 5:00 PM|           curtailed|      560.0|             803.2|
    |  2| 2021-02-01 5:00 PM| 2021-02-01 7:00 PM|           curtailed|      140.0|             200.8|
    |  3| 2021-02-01 7:00 PM|2021-02-01 11:00 PM|down_for_maintenance|        0.0|             301.0|
    |  4|2021-02-01 11:00 PM| 2021-02-02 3:00 AM|           curtailed|      320.0|             402.0|
    |  5| 2021-02-02 3:00 AM| 2021-02-02 5:00 AM|      over_producing|      240.0|             200.8|
    |  5| 2021-02-02 5:00 AM| 2021-02-02 8:00 AM|      over_producing|      360.0|             301.2|
    |  6| 2021-02-02 8:00 AM|2021-02-02 11:00 AM|          normal_run|      280.0|             301.0|
    |  7|2021-02-02 11:00 AM| 2021-02-02 5:00 PM|broken_belt_unsch...|        0.0| 602.6037735849056|
    |  7| 2021-02-02 5:00 PM| 2021-02-03 5:00 AM|broken_belt_unsch...|        0.0|1205.2075471698113|
    |  7| 2021-02-03 5:00 AM| 2021-02-03 5:00 PM|broken_belt_unsch...|        0.0|1205.2075471698113|
    |  7| 2021-02-03 5:00 PM| 2021-02-04 5:00 AM|broken_belt_unsch...|        0.0|1205.2075471698113|
    |  7| 2021-02-04 5:00 AM| 2021-02-04 4:00 PM|broken_belt_unsch...|        0.0|1104.7735849056605|
    +---+-------------------+-------------------+--------------------+-----------+------------------+
    

    【讨论】:

    • 谢谢我今天会这样做!
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2021-02-21
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2013-08-28
    • 1970-01-01
    相关资源
    最近更新 更多