【发布时间】: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