Spark中Drools性能优化问询:规则激增致处理时长暴增
Great question—let’s break down why your processing time skyrocketed and how to fix it.
What’s Wrong With Your Current Code?
The biggest issue here is per-item KieSession initialization:
- Every time
seqOpruns, you create a newStatelessKieSessionfrom the broadcastedKieBase. Creating a KieSession is an expensive operation—Drools has to initialize the rule engine, set up working memory, and prepare all your rules for matching. With 10k rules, this initialization time is non-trivial, and multiplying it by every item in your RDD is what’s turning 20 minutes into 10 hours. - Processing one fact at a time negates Drools’ core optimizations. Drools is designed to handle batches of facts efficiently, as it can build indexes and run pattern matches across the entire set of facts in one go. Processing single facts forces the engine to repeat work for every item.
Fixes to Get Your Performance Back
1. Reuse KieSessions Per Partition
Instead of creating a session for each item, create one session per RDD partition. Spark partitions data across workers, so initializing a session once per partition amortizes the setup cost across all items in that partition.
Here’s a refactored approach using mapPartitions:
def processPartition[T: ClassTag](broadcastRules: Broadcast[KieBase])(partition: Iterator[T]): Iterator[MyAggregator] = { // Initialize a single session for the entire partition val kieBase = broadcastRules.value val session = kieBase.newStatelessKieSession() // Initialize your aggregator for the partition val partitionAggregator = MyAggregator() session.setGlobal("aggregator", partitionAggregator) // Batch-process all facts in the partition val facts = partition.toList session.execute(CommandFactory.newInsertElements(facts)) // Return the aggregated result for the partition Iterator(partitionAggregator) } // Usage in your main job: val finalAggregator = rdd .mapPartitions(processPartition(broadcastRules)) .reduce((agg1, agg2) => agg1.merge(agg2)) // Merge results from all partitions
2. Batch Insert Facts
Using newInsertElements instead of newInsert lets Drools process all facts in one pass. This leverages Drools’ internal optimizations like pattern matching indexes and reduces the overhead of repeated rule engine invocations.
3. Additional Optimizations
- Optimize Your KieBase: Ensure your KieBase is compiled with performance in mind. For example:
- Enable sequential mode (if rule execution order doesn’t affect your results) via
kiebase.setSequential(true). - Enable indexing for your facts to speed up pattern matching.
- Enable sequential mode (if rule execution order doesn’t affect your results) via
- Simplify Rules: Audit your 10k rules to remove redundant conditions or actions. Complex rules with nested patterns can slow down matching—look for opportunities to simplify or combine rules.
- Check Serialization: Verify that your
KieBaseis being broadcast efficiently. If serialization/deserialization is adding overhead, consider using a serializable wrapper or ensuring your KieBase is built with serialization in mind.
Why This Works
By reusing sessions per partition and batch-processing facts, you cut out the massive overhead of repeated KieSession initialization. With 10k rules, this alone will bring your processing time back down to something manageable. Drools’ batch processing capabilities mean the rule engine only does the heavy lifting of preparing rules once per partition, not per item.
内容的提问来源于stack exchange,提问作者Zied Koubaa

