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

Spark:在flag=1标记区间填充对应id哈希值的实现方案

Spark实现区间哈希值填充方案

要实现将flag=1起始行的id哈希值填充到该区间(当前flag=1到下一个flag=1前)的所有行,可以通过窗口函数分组+哈希计算完成,具体步骤如下:

核心思路

  1. 先将数据按日期排序,确保区间划分符合时间顺序
  2. 用累计求和生成分组标识,每个flag=1会开启一个新分组
  3. 针对每个分组,提取起始行(flag=1)的id并计算哈希,填充到组内所有行

代码实现(Scala版本)

import org.apache.spark.sql.functions._
import org.apache.spark.sql.types.DateType
import org.apache.spark.sql.Window

// 假设原始数据集为df,包含id、date、flg字段
val formattedDF = df
  // 将字符串日期转换为Date类型,确保排序正确
  .withColumn("date", to_date(col("date"), "dd.MM.yyyy"))
  // 按日期排序,保证区间按时间顺序划分
  .orderBy("date")

// 生成分组ID:每遇到一个flg=1,分组ID递增
val groupWindow = Window.orderBy("date").rowsBetween(Window.unboundedPreceding, Window.currentRow)
val groupedDF = formattedDF.withColumn(
  "group_id",
  sum(when(col("flg") === 1, 1).otherwise(0)).over(groupWindow)
)

// 提取每个分组的起始id,计算哈希并填充
val hashWindow = Window.partitionBy("group_id")
val resultDF = groupedDF
  .withColumn(
    "start_id",
    first(when(col("flg") === 1, col("id")).otherwise(null)).over(hashWindow)
  )
  // 计算id的MD5哈希并截取前5位,匹配示例格式
  .withColumn("hash", substring(md5(col("start_id").cast("string")), 1, 5))
  // 删除中间字段
  .drop("group_id", "start_id")

// 查看结果
resultDF.show()

代码实现(Python版本)

from pyspark.sql import functions as F
from pyspark.sql.window import Window

# 处理日期格式并排序
formatted_df = df\
    .withColumn("date", F.to_date(F.col("date"), "dd.MM.yyyy"))\
    .orderBy("date")

# 生成分组标识
group_window = Window.orderBy("date").rowsBetween(Window.unboundedPreceding, Window.currentRow)
grouped_df = formatted_df.withColumn(
    "group_id",
    F.sum(F.when(F.col("flg") == 1, 1).otherwise(0)).over(group_window)
)

# 计算哈希并填充到区间
hash_window = Window.partitionBy("group_id")
result_df = grouped_df\
    .withColumn(
        "start_id",
        F.first(F.when(F.col("flg") == 1, F.col("id")).otherwise(None)).over(hash_window)
    )\
    .withColumn("hash", F.substring(F.md5(F.col("start_id").cast("string")), 1, 5))\
    .drop("group_id", "start_id")

# 展示结果
result_df.show()

关键细节说明

  • 日期转换与排序:原始数据的日期是字符串格式,必须先转为Date类型再排序,否则会出现字符串排序错误(比如"11.01.2024"会错误排在"06.01.2024"之后)
  • 分组标识生成:通过累计求和sum(when(...)),每遇到一个flg=1就给分组ID加1,确保同一个时间区间的行属于同一分组
  • 哈希计算:示例中使用MD5哈希并截取前5位,你也可以根据需求替换为hash()函数(Spark内置哈希)或其他哈希算法

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.21 20:25:54