【问题标题】:How to expire state of dropDuplicates in structured streaming to avoid OOM?如何在结构化流中使 dropDuplicates 的状态过期以避免 OOM?
【发布时间】:2018-01-10 11:19:45
【问题描述】:

我想使用 spark 结构化流计算每天的唯一访问权限,所以我使用以下代码

.dropDuplicates("uuid")

并且在第二天,今天维护的状态应该被删除,这样我就可以获得第二天唯一访问的正确计数并避免 OOM。 spark文档指示使用带有水印的dropDuplicates,例如:

.withWatermark("timestamp", "1 day")
.dropDuplicates("uuid", "timestamp")

但必须在 dropDuplicates 中指定水印列。在这种情况下,uuid 和时间戳将用作组合键,以对具有相同 uuid 和时间戳的元素进行重复数据删除,这不是我所期望的。

那么有完美的解决方案吗?

【问题讨论】:

    标签: apache-spark duplicates apache-spark-sql out-of-memory spark-structured-streaming


    【解决方案1】:

    我发现window函数不起作用所以我选择使用window.start或window.end。

    .select(
       window($"timestamp", "1 day").start,
       $"timestamp",
       $"uuid"
    )
    .withWatermark("window", "1 day")
    .dropDuplicates("uuid", "window")
    

    【讨论】:

      【解决方案2】:

      以下是对 Spark 文档中提出的程序的修改。技巧是操纵事件时间,即将事件时间放入 桶。假设事件时间以毫秒为单位。

      // removes all duplicates that are in 15 minutes tumbling window.
      // doesn't remove duplicates that are in different 15 minutes windows !!!!
      public static Dataset<Row> removeDuplicates(Dataset<Row> df) {
          // converts time in 15 minute buckets
          // timestamp - (timestamp % (15 * 60))
          Column bucketCol = functions.to_timestamp(
                  col("event_time").divide(1000).minus((col("event_time").divide(1000)).mod(15*60)));
          df = df.withColumn("bucket", bucketCol);
      
          String windowDuration = "15 minutes";
          df = df.withWatermark("bucket", windowDuration)
                  .dropDuplicates("uuid", "bucket");
      
          return df.drop("bucket");
      }
      

      【讨论】:

        【解决方案3】:

        经过几天的努力,我终于找到了自己的路。

        在研究watermarkdropDuplicates的源码时,发现watermark除了eventTime列外,还支持window列,所以我们可以使用如下代码:

        .select(
            window($"timestamp", "1 day"),
            $"timestamp",
            $"uuid"
          )
        .withWatermark("window", "1 day")
        .dropDuplicates("uuid", "window")
        

        由于同一天的所有事件都具有相同的窗口,因此这将产生与仅使用 uuid 进行重复数据删除相同的结果。希望可以帮助某人。

        【讨论】:

          猜你喜欢
          • 2020-09-09
          • 2020-12-13
          • 1970-01-01
          • 1970-01-01
          • 2021-12-10
          • 2021-08-06
          • 1970-01-01
          • 2013-02-25
          • 1970-01-01
          相关资源
          最近更新 更多