基于独立flink-cep实现无Keyed Stream的患者心率事件分组
解决方案:手动模拟分组+时间窗口,结合Flink CEP独立库实现匹配
因为没法用Flink的Keyed Stream,我们可以手动实现按患者ID分组和5秒时间窗口过滤的逻辑,再结合Flink CEP的核心API完成匹配。
前提补充
假设你的Event类包含时间戳字段(事件时间,毫秒级),定义如下:
data class Event( val type: String, val properties: Map<String, Any>, val timestamp: Long // 事件发生的时间戳,必须有这个才能做时间窗口 )
核心实现逻辑
- 用线程安全的Map维护每个患者的事件窗口,替代Keyed Stream的分组能力
- 对每个新事件,先清理对应患者窗口中超过5秒的旧事件
- 将新事件加入窗口,再用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
相关产品推荐
相关产品推荐

