Spark flatMapGroupsWithState随机丢事件问题求助
作业流程
- 从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

