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

Spark flatMapGroupsWithState随机丢事件问题求助

Spark流处理作业事件丢失与结果不一致问题

作业流程

  • 从Delta Lake读取静态dataFrame;
  • 从Delta Lake读取流式dataFrame;
  • 将流数据与静态数据关联;
  • 执行flatMapGroupsWithState操作;
  • 输出结果。

问题现象

flatMapGroupsWithState环节出现事件丢失,输出结果随机,相同输入重新运行得到不同输出。本地模式下给写入操作添加.coalesce(1)后能得到预期结果,但集群模式下无效。

相关代码

主流程代码

val entityScheduleSlots = data
  .withColumn("products", concat(col("batteries"), col("photovoltaics")))
  .drop("photovoltaics", "batteries", "labels")
  .join(
    entities,
    array_contains(entities("entity_delivery_points"), col("delivery_point_id")))
  .withColumn("now", current_timestamp())
  .withWatermark("now", "5 minutes")
  .as(Encoders.product[enrichedDeliveryPointSchedule])
  .groupByKey(e => e.timestamp.toString + e.entity_id.toString + e.schedule_id)(
    Encoders.STRING)
  .flatMapGroupsWithState(
    outputMode = OutputMode.Append,
    timeoutConf = GroupStateTimeout.EventTimeTimeout)(
    Function.computeExplodedEntityScheduleSlots)(
    Encoders.kryo[Function.State],
    Encoders.product[EntityScheduleSlot])

entityScheduleSlots为输出结果,测试在本地模式进行。

状态处理代码

object Function {
  case class ProductState(
      var count: Int,
      var quantity: Double,
      var price: Double,
      val sellable: Boolean)
  case class State(var delivery_points_count: Int, var products: mutable.Map[Long, ProductState])
  private def computeExplodedEntityScheduleSlots(
      uid: String,
      ss: Iterator[enrichedDeliveryPointSchedule],
      state: GroupState[State]): Iterator[EntityScheduleSlot] = {
    if (state.hasTimedOut) {
      state.remove()
      return Iterator.empty
    }
    val schedules = ss.toList
    val newState: State =
      state.getOption.getOrElse(State(0, mutable.Map()))
    schedules.foreach(s => {
      newState.delivery_points_count = newState.delivery_points_count + 1
      val qualificationsProductsIDs =
        if (s.entity_qualifications != null) s.entity_qualifications.map(q => q.product)
        else List()
      if (s.products != null) {
        s.products.foreach(p => {
          if (qualificationsProductsIDs.contains(p.product)) {
            val productState =
              newState.products.getOrElse(p.product, ProductState(0, 0.0, 0.0, p.sellable))
            val factor =
              if (productState.count == 0) 1
              else p.quantity / (productState.quantity / productState.count)
            productState.quantity += p.quantity
            productState.price =
              (productState.price * productState.count + p.price * factor) / (productState.count + 1)
            productState.count += 1
            newState.products.update(p.product, productState)
          }
        })
      }
    })
    if (newState.delivery_points_count == schedules.head.entity_delivery_points.length) {
      state.remove()
      return Iterator(
        EntityScheduleSlot(
          timestamp = schedules.head.timestamp,
          entity = schedules.head.entity_id,
          schedule_timestamp = schedules.head.schedule_timestamp,
          schedule_id = schedules.head.schedule_id,
          products =
            if (schedules.head.entity_qualifications != null)
              schedules.head.entity_qualifications
                .map(q => {
                  val product =
                    newState.products.getOrElse(q.product, ProductState(0, 0.0, 0.0, false))
                  EntityScheduleSlotProduct(
                    q.product,
                    product.quantity,
                    product.price,
                    product.sellable)
                })
            else List()))
    }
    state.update(newState)
    val currentWatermarkMs =
      if (state.getCurrentWatermarkMs() > 0) state.getCurrentWatermarkMs()
      else System.currentTimeMillis()
    state.setTimeoutTimestamp(currentWatermarkMs, "2 minutes")
    Iterator.empty
  }
}

case class enrichedDeliveryPointSchedule(
    timestamp: java.sql.Timestamp,
    schedule_timestamp: java.sql.Timestamp,
    schedule_id: String,
    delivery_point_id: Long,
    products: List[DeliveryPointScheduleSlotProduct],
    entity_id: Long,
    entity_delivery_points: List[Long],
    entity_qualifications: List[EntityQualification])

问题根源与解决方案

1. 分组键的非确定性

你使用e.timestamp.toString + e.entity_id.toString + e.schedule_id作为分组键,java.sql.Timestamp.toString()的输出依赖JVM时区配置,集群中不同节点时区不一致会导致同一事件生成不同分组键,分散到不同状态处理器,造成聚合错误和结果混乱。

修复:
改用固定格式的时间字符串或时间戳毫秒数生成分组键,避免时区影响:

// 方法1:用date_format统一时间格式
.withColumn("timestamp_str", date_format(col("timestamp"), "yyyy-MM-dd'T'HH:mm:ss.SSS"))
.groupByKey(e => e.timestamp_str + e.entity_id.toString + e.schedule_id)(Encoders.STRING)

// 方法2:直接用时间戳毫秒数
.groupByKey(e => e.timestamp.getTime.toString + e.entity_id.toString + e.schedule_id)(Encoders.STRING)

2. 可变状态的线程安全隐患

代码中使用mutable.Map和可变ProductState,虽然单分组处理是单线程的,但集群环境下状态序列化/反序列化可能出现数据错乱,导致状态更新丢失。

修复:
改用不可变数据结构,每次状态更新生成新对象:

// 定义不可变状态类
case class ProductState(
    count: Int,
    quantity: Double,
    price: Double,
    sellable: Boolean)
case class State(delivery_points_count: Int, products: Map[Long, ProductState], entityDeliveryPointsLength: Int)

// 状态更新逻辑示例
val updatedProductState = productState.copy(
  count = productState.count + 1,
  quantity = productState.quantity + p.quantity,
  price = (productState.price * productState.count + p.price * factor) / (productState.count + 1)
)
val updatedProducts = newState.products.updated(p.product, updatedProductState)
val updatedState = newState.copy(
  delivery_points_count = newState.delivery_points_count + 1,
  products = updatedProducts
)

3. Watermark与超时逻辑错误

你用current_timestamp()生成的now列作为watermark依据,这是系统执行时间而非事件时间,会导致watermark不合理前进,提前触发状态超时丢失事件。

修复:
改用流数据中的实际事件时间字段(比如timestamp)作为watermark:

.withWatermark("timestamp", "5 minutes")

同时,超时时间应基于事件时间计算,确保状态有足够时间接收全部分组数据。

4. 依赖schedules.head的风险

直接取schedules.head获取entity_delivery_points.length,集群下批次内事件顺序不确定,会导致聚合完成的判断条件错误。

修复:
首次初始化状态时保存entity_delivery_points的长度,后续判断依赖状态中存储的值:

// 初始化状态时保存目标长度
val initialState = State(0, Map.empty, schedules.head.entity_delivery_points.length)
// 判断聚合完成时用状态中的值
if (newState.delivery_points_count == newState.entityDeliveryPointsLength) {
  // 生成输出逻辑
}

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.04 03:15:47