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

基于独立flink-cep实现无Keyed Stream的患者心率事件分组

因为没法用Flink的Keyed Stream,我们可以手动实现按患者ID分组和5秒时间窗口过滤的逻辑,再结合Flink CEP的核心API完成匹配。

前提补充

假设你的Event类包含时间戳字段(事件时间,毫秒级),定义如下:

data class Event(
    val type: String,
    val properties: Map<String, Any>,
    val timestamp: Long // 事件发生的时间戳,必须有这个才能做时间窗口
)

核心实现逻辑

  1. 用线程安全的Map维护每个患者的事件窗口,替代Keyed Stream的分组能力
  2. 对每个新事件,先清理对应患者窗口中超过5秒的旧事件
  3. 将新事件加入窗口,再用Flink CEP匹配该窗口内的事件模式

完整代码示例

import org.apache.flink.cep.CEP
import org.apache.flink.cep.PatternSelectFunction
import org.apache.flink.cep.pattern.Pattern
import org.apache.flink.cep.pattern.conditions.SimpleCondition
import java.util.concurrent.ConcurrentHashMap

// 定义常量
val patientKey = "patient"
val hrKey = "hr"
val WINDOW_DURATION = 5000L // 5秒窗口,毫秒

// 维护患者的事件窗口:key=患者ID,value=该患者5秒内的事件列表
val patientEventWindows = ConcurrentHashMap<Int, MutableList<Event>>()

// 定义CEP模式:匹配同一患者的任意数量心率事件(5秒内的所有)
val hrPattern = Pattern.begin<Event>("hr-events")
    .where(SimpleCondition { it.type == "hr" })
    .oneOrMore() // 匹配一个或多个心率事件
    .within(WINDOW_DURATION) // 时间窗口5秒

fun processEvent(event: Event) {
    // 提取患者ID
    val patientId = event.properties[patientKey] as Int
    // 获取或创建该患者的事件窗口
    val eventList = patientEventWindows.computeIfAbsent(patientId) { mutableListOf() }

    // 第一步:清理窗口内超过5秒的旧事件
    val currentTime = event.timestamp
    eventList.removeIf { currentTime - it.timestamp > WINDOW_DURATION }

    // 第二步:添加当前事件到窗口
    eventList.add(event)

    // 第三步:用Flink CEP匹配当前窗口的事件
    val patternStream = CEP.pattern(eventList, hrPattern)
    // 提取匹配结果
    val matches = patternStream.select(PatternSelectFunction { map ->
        map["hr-events"] as List<Event>
    }).collect()

    // 输出匹配结果(如果有)
    matches.forEachIndexed { index, matchedEvents ->
        println("第${index+1}组匹配:${matchedEvents.joinToString { it.properties[hrKey].toString() }}")
    }
}

// 测试代码
fun main() {
    // 生成带时间戳的事件(假设事件按顺序发生,时间戳递增)
    val baseTimestamp = System.currentTimeMillis()
    val p1e1 = Event("hr", mapOf(patientKey to 1, hrKey to 1), baseTimestamp)
    val p1e2 = Event("hr", mapOf(patientKey to 1, hrKey to 2), baseTimestamp + 1000)
    val p2e1 = Event("hr", mapOf(patientKey to 2, hrKey to 1), baseTimestamp + 1500)
    val p1e3 = Event("hr", mapOf(patientKey to 1, hrKey to 3), baseTimestamp + 2000)
    val p2e2 = Event("hr", mapOf(patientKey to 2, hrKey to 2), baseTimestamp + 2500)
    val p3e1 = Event("hr", mapOf(patientKey to 3, hrKey to 1), baseTimestamp + 3000)
    val p2e3 = Event("hr", mapOf(patientKey to 2, hrKey to 3), baseTimestamp + 3500)
    val p3e2 = Event("hr", mapOf(patientKey to 3, hrKey to 2), baseTimestamp + 4000)
    val p3e3 = Event("hr", mapOf(patientKey to 3, hrKey to 3), baseTimestamp + 4500)

    // 逐个处理事件
    listOf(p1e1, p1e2, p2e1, p1e3, p2e2, p3e1, p2e3, p3e2, p3e3).forEach { processEvent(it) }
}

关键细节说明

  • 分组容器:用ConcurrentHashMap保证多线程场景下的安全,每个患者对应一个独立的事件窗口,完全替代Keyed Stream的分组逻辑。
  • 窗口清理:每次处理新事件时,先移除窗口内时间差超过5秒的旧事件,确保窗口内始终是符合时间要求的事件。
  • CEP匹配:对每个患者的窗口事件单独应用CEP模式,获取该患者5秒内的所有心率事件组,和你期望的分组结果完全一致。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.16 07:10:33