【问题标题】:Need help in filtering records according to set of rules with Apache Spark在使用 Apache Spark 根据规则集过滤记录时需要帮助
【发布时间】:2017-11-19 06:28:03
【问题描述】:

在我遇到的使用 Apache Spark 根据一组规则过滤记录的用例中,我需要帮助。

由于实际数据的字段太多,例如,可以考虑如下数据(为简单起见,以JSON格式给出数据),

records : [{
             "recordId": 1, 
             "messages": [{"name": "Tom","city": "Mumbai"}, 
                         {"name": "Jhon","address": "Chicago"}, .....]
           },....]

rules : [{
           ruleId: 1,
           ruleName: "rule1",
           criterias: {
                          name: "xyz",
                          address: "Chicago, Boston"
                      }
          }, ....]

我想根据所有规则匹配所有记录。这是伪代码:

 var matchedRecords = []
 for(record <- records)
    for(rule <- rules)
       for(message <- record.message)
           if(!isMatch(message, rule.criterias)) 
               break;
       if(allMessagesMatched) // If loop completed without break
           matchedRecords.put((record.id, ruleId))



 def isMatch(message, criteria) = 
           for(each field in crieteria)
                if(field.value contains comma)
                    if(! message.field containsAny field.value)
                        return false
                else if(!message.field equals field.value) // value doesnt contain comma
                     return false
           return true // if loop completed that means all criterias are matched

有数千条记录包含数千条消息,并且有数百条这样的规则。

解决此类问题的方法是什么?任何特定的模块都会有帮助,比如(SparkSQL、Spark Mlib、Spark GraphX)?我需要使用任何第三方库吗?

方法一:

  • 有列表[规则] & RDD[记录]

  • 广播列表[规则],因为它们的数量较少。

  • 将每条记录与所有规则匹配。

在这种情况下,仍然不会发生将每条消息与条件匹配的并行计算。

【问题讨论】:

    标签: apache-spark apache-spark-sql apache-spark-mllib spark-graphx


    【解决方案1】:

    我认为您建议的方法是好的方向。如果我必须解决这个任务,我将从使用负责匹配的方法实现通用特征开始:

    trait FilterRule extends Serializable {
       def match(record: Record): Boolean
    }
    

    然后我会实现特定的过滤器,例如:

    class EqualsRule extends FilterRule
    class RegexRule extends FilterRule
    

    然后我会实现复合过滤器,例如:

    class AndRule extends FilterRule
    class OrRule extends FilterRule
    ...
    

    然后你可以过滤你的rdd或DataSet:

    // constructing rule - in reality reading json from configuration, parsing json and creating FilterRule object
    val rule = AndRule(EqualsRule(...), EqualsRule(...), ...) 
    
    // applying rule
    rdd.filter(record => rule.match(r))
    

    第二个选项是尝试使用现有的 Spark SQL 函数和 DataFrame 进行过滤,您可以在其中使用 and 或 or 多个列构建非常复杂的表达式。这种方法的缺点是类型不安全,并且单元测试会更复杂。

    【讨论】:

      猜你喜欢
      • 2011-09-28
      • 1970-01-01
      • 2020-05-29
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2015-07-13
      • 1970-01-01
      • 2022-10-13
      相关资源
      最近更新 更多