Spark Streaming集成Drools时滑动窗口规则全局KieSession为空问题
你的核心问题源于Spark的分布式执行模型与Driver/Executor的内存隔离,具体点:
@transient变量的序列化行为:你标记
globalKieSessionRef为@transient lazy val,这意味着该变量不会被序列化。当ForeachWriter实例被Driver序列化发送到Executor节点时,Executor会重新创建ForeachWriter的实例,此时globalKieSessionRef会重新初始化(回到null)——而你只在Driver端执行了initializeGlobalKieSession(),Executor端从未执行过这个初始化逻辑,所以process方法里拿到的永远是null。deviceList的同步问题:
deviceList是普通var变量,同样会被序列化到Executor,但后续Driver端的更新(比如timer任务里的修改)不会同步到Executor,Executor端的deviceList始终是初始的空列表,可能连deviceList.contains(Deviceid)的判断都无法触发。Timer任务的局限性:你在Driver端启动的Timer只会在Driver进程中运行,Executor端完全感知不到,定时更新
globalKieSessionRef的操作也只会影响Driver端的实例,无法同步到Executor。
针对Spark Streaming与Drools集成的场景,需要适配Spark的分布式模型,避免直接在Driver和Executor之间共享内存对象,以下是可行的方案:
1. 用广播变量传递规则,Executor端懒加载KieSession
KieSession本身不可序列化,因此不能直接广播,但可以广播滑动窗口规则的DRL内容,让每个Executor节点在需要时自行创建并维护KieSession:
调整步骤:
- 将滑动窗口规则和设备列表封装为Spark广播变量,在Driver端初始化并定时更新。
- 在
ForeachWriter的open方法中(每个Partition的Task初始化时执行),基于广播的规则创建KieSession,存储为@transient成员变量。 - 定时更新广播变量,触发Executor端重新创建KieSession。
代码示例:
// 用广播变量存储滑动窗口规则和设备列表 private var slidingWindowRulesBroadcast: Broadcast[String] = _ private var deviceListBroadcast: Broadcast[List[String]] = _ private lateinit var spark: SparkSession def initializeResources(): Unit = { val (allSlidingWindowRules, devices) = fetchSlidingWindowRules() slidingWindowRulesBroadcast = spark.sparkContext.broadcast(allSlidingWindowRules) deviceListBroadcast = spark.sparkContext.broadcast(devices) } // 定时更新广播变量 def recreateResources(): Unit = { val (allSlidingWindowRules, devices) = fetchSlidingWindowRules() // 先释放旧的广播变量资源 slidingWindowRulesBroadcast.unpersist() deviceListBroadcast.unpersist() // 创建新的广播变量 slidingWindowRulesBroadcast = spark.sparkContext.broadcast(allSlidingWindowRules) deviceListBroadcast = spark.sparkContext.broadcast(devices) } def Start_RuleEngine(): Unit = { spark = SparkSession.builder .appName("SparkSQL") .master("local[*]") .getOrCreate() initializeResources() // 定时更新资源(每30分钟) val timer = new java.util.Timer() val task = new java.util.TimerTask { def run(): Unit = recreateResources() } timer.schedule(task, 0, 30 * 60 * 1000) import spark.implicits._ val Delta_table_location = sys.env.getOrElse("DELTA_TABLE", "") val read_deltatable = spark.readStream.delta(Delta_table_location) val process_deltatable = read_deltatable.writeStream.foreach(new ForeachWriter[Row] { // 每个Partition专属的KieSession,@transient避免序列化问题 @transient private var slidingKieSession: KieSession = _ def open(partitionId: Long, epochId: Long): Boolean = { // 从广播变量获取规则,初始化KieSession val rules = slidingWindowRulesBroadcast.value slidingKieSession = createKieSessionWithRules(rules) true } def process(value: Row): Unit = { val Telemetery = value.getString(0) val jsonObj = new JSONObject(Telemetery) val Deviceid = jsonObj.getString("DeviceId") // 使用广播的设备列表判断,同时检查KieSession是否有效 if (deviceListBroadcast.value.contains(Deviceid) && slidingKieSession != null) { fireKieSession(slidingKieSession, jsonObj, Deviceid) } // 独立规则的KieSession逻辑保持不变 val rules = GetApprovedDroolRule(Deviceid) val kieSession = createKieSessionWithRules(rules) fireKieSession(kieSession, jsonObj, Deviceid) kieSession.dispose() } def close(errorOrNull: Throwable): Unit = { // 关闭当前Partition的KieSession,释放资源 if (slidingKieSession != null) { slidingKieSession.dispose() } } }) process_deltatable.start().awaitTermination() }
2. 结合Spark状态流管理窗口状态
如果滑动窗口规则依赖跨批次的状态,建议使用Spark本身的Stateful Streaming来管理状态,而不是依赖Drools KieSession的内存状态——因为Executor节点重启时,KieSession的内存状态会丢失,而Spark的状态流可以将状态持久化到外部存储(如HDFS、S3),保证容错性。
3. 注意KieSession的线程安全
Drools的KieSession是线程不安全的,每个Task(对应一个Partition)维护自己的KieSession实例即可,因为Spark的Task是单线程执行的,不会出现并发访问问题。
内容的提问来源于stack exchange,提问作者Basil Saju

